温馨提示×

温馨提示×

您好,登录后才能下订单哦!

密码登录×
登录注册×
其他方式登录
点击 登录注册 即表示同意《亿速云用户服务条款》

如何使用Flink进行流式数据计算

发布时间:2026-07-01 18:08:13 来源:亿速云 阅读:98 作者:小樊 栏目:软件技术

使用 Apache Flink 进行流式数据计算,通常包括 环境搭建 → 数据源接入 → 流式计算逻辑 → 数据输出 → 运行与监控 几个步骤。下面我从 概念 + 示例 + 常见场景 三个层面,系统性地讲一下。


一、Flink 流式计算核心概念

1️⃣ 什么是流式计算

流式计算是指 数据持续不断到达,系统实时处理并产生结果,而不是等数据全部就绪再计算。

Flink 的特点:

  • 真正的流式计算(不是微批)
  • 事件时间 / 处理时间 / 摄入时间
  • 有状态计算
  • 精确一次(Exactly-once)语义
  • 高吞吐、低延迟

2️⃣ Flink 流式计算核心模型

DataStream API(主流)
  ↓
Source → Transformation → Sink

核心抽象:

  • DataStream:流式数据集
  • Transformation:map / filter / keyBy / window 等
  • State:状态(ValueState / ListState)
  • Time:时间语义 + 窗口
  • Checkpoint:容错机制

二、Flink 流式计算基本流程

1️⃣ 准备环境

StreamExecutionEnvironment env =
        StreamExecutionEnvironment.getExecutionEnvironment();

常用配置:

env.enableCheckpointing(5000); // 每5秒一次 checkpoint
env.setParallelism(2);

2️⃣ 接入数据源(Source)

示例:Socket 数据源(上手用)

DataStream<String> stream = env
        .socketTextStream("localhost", 9999);

常见生产级 Source:

  • Kafka(最常用)
  • Pulsar
  • Kinesis
  • MySQL CDC(CDC 流)

Kafka Source 示例:

KafkaSource<String> source = KafkaSource
        .<String>builder()
        .setBootstrapServers("localhost:9092")
        .setTopics("input-topic")
        .setGroupId("flink-group")
        .setValueOnlyDeserializer(new SimpleStringSchema())
        .build();

DataStream<String> stream = env.fromSource(
        source,
        WatermarkStrategy.noWatermarks(),
        "Kafka Source"
);

3️⃣ 定义流式计算逻辑

1. 简单转换

DataStream<Tuple2<String, Integer>> result =
        stream
            .flatMap((String line, Collector<Tuple2<String, Integer>> out) -> {
                for (String word : line.split(" ")) {
                    out.collect(new Tuple2<>(word, 1));
                }
            })
            .returns(Types.TUPLE(Types.STRING, Types.INT))
            .keyBy(value -> value.f0)
            .sum(1);

2. 窗口计算(核心)

滚动窗口(Tumbling Window)

DataStream<WordCount> result =
        stream
            .map(line -> new WordCount(line, 1))
            .keyBy(WordCount::getWord)
            .window(TumblingProcessingTimeWindows.of(Time.seconds(10)))
            .sum("count");

滑动窗口(Sliding Window)

.window(SlidingProcessingTimeWindows.of(Time.seconds(30), Time.seconds(10)))

3. 事件时间 + Watermark(推荐)

WatermarkStrategy<String> strategy =
        WatermarkStrategy
            .<String>forBoundedOutOfOrderness(Duration.ofSeconds(5))
            .withTimestampAssigner((event, timestamp) -> extractTimestamp(event));

DataStream<String> stream = env
        .fromSource(kafkaSource, strategy, "kafka");

4️⃣ 输出结果(Sink)

Kafka Sink

stream.sinkTo(kafkaSink);

打印(调试)

result.print();

MySQL / HBase / Redis(业务输出)

result.addSink(new MySQLSink());

5️⃣ 执行任务

env.execute("Flink Streaming Job");

⚠️ Flink 是懒执行,只有 execute 才会运行


三、一个完整示例(WordCount 流式版)

public class StreamWordCount {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env =
                StreamExecutionEnvironment.getExecutionEnvironment();

        DataStream<String> text =
                env.socketTextStream("localhost", 9999);

        DataStream<Tuple2<String, Integer>> counts =
                text.flatMap((String line, Collector<Tuple2<String, Integer>> out) -> {
                    for (String word : line.split(" ")) {
                        out.collect(new Tuple2<>(word, 1));
                    }
                })
                .returns(Types.TUPLE(Types.STRING, Types.INT))
                .keyBy(t -> t.f0)
                .sum(1);

        counts.print();

        env.execute("Stream WordCount");
    }
}

四、常见流式计算场景

场景 核心用法
实时 ETL map / filter
实时统计 keyBy + window
实时去重 state / keyBy
实时 TopN KeyedProcessFunction
实时对账 CoProcessFunction
CDC 同步 Flink CDC
复杂事件处理 CEP

五、Flink 流式计算的关键点(很重要)

✅ 1. 状态管理

  • Flink 能做 有状态流式计算
  • 状态支持:
    • ValueState
    • ListState
    • MapState

✅ 2. 精确一次语义

env.enableCheckpointing(5000);

✅ 3. 时间语义选择

  • 处理时间:简单但不可靠
  • 事件时间:推荐(真实业务时间)

六、学习建议

推荐学习路径

  1. 先跑 Socket / WordCount
  2. 再学 Kafka + 窗口
  3. 再学 状态 + 事件时间
  4. 最后学 CDC / CEP / SQL

官方文档(必看)

  • https://flink.apache.org
  • https://nightlies.apache.org/flink/flink-docs-stable/

如果你愿意,我可以: ✅ 给你一个 Kafka + Flink CDC 的完整示例
✅ 手把手教你 实时 PV / UV 统计
✅ 解释 窗口触发、Watermark、迟到数据
✅ 对比 Flink vs Spark Streaming

你可以直接告诉我你现在的 使用场景技术栈

向AI问一下细节

免责声明:本站发布的内容(图片、视频和文字)以原创、转载和分享为主,文章观点不代表本网站立场,如果涉及侵权请联系站长邮箱:is@yisu.com进行举报,并提供相关证据,一经查实,将立刻删除涉嫌侵权内容。

AI