我们用Flink最经典的入门案例—— 实时单词计数来演示完整流程:模拟从端口源源不断接收文本数据流,Flink实时统计每个单词的累计出现次数,全程低延迟、自动维护计数状态。

一、完整代码示例(Java版)

这是一个可直接运行的最小流处理程序:

import org.apache.flink.api.common.functions.FlatMapFunction;
import org.apache.flink.api.java.tuple.Tuple2;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.util.Collector;

public class StreamWordCount {
    public static void main(String[] args) throws Exception {
        // 1. 初始化流执行环境(Flink程序入口)
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

        // 2. 定义数据源 Source:从本地 9999 端口读取实时文本流
        DataStream<String> textStream = env.socketTextStream("localhost", 9999);

        // 3. 核心转换逻辑 Transformation
        DataStream<Tuple2<String, Integer>> resultStream = textStream
                // 拆分每行文本,每个单词转为 (单词, 1) 的二元组
                .flatMap((FlatMapFunction<String, Tuple2<String, Integer>>) (line, out) -> {
                    for (String word : line.split(" ")) {
                        out.collect(Tuple2.of(word, 1));
                    }
                })
                // 按单词分组,相同单词会进入同一个计算任务
                .keyBy(tuple -> tuple.f0)
                // 对计数字段累加,Flink自动维护每个单词的累计状态
                .sum(1);

        // 4. 定义输出 Sink:结果打印到控制台
        resultStream.print();

        // 5. 触发任务执行(Flink懒加载,必须调用才会真正启动任务)
        env.execute("实时单词计数任务");
    }
}

二、逐步骤流程解析

整个流程严格遵循 Source(数据源)→ Transformation(转换计算)→ Sink(结果输出) 的流处理模型:

  1. 环境初始化StreamExecutionEnvironment 是所有Flink流程序的入口,负责配置运行参数、管理作业生命周期。
  2. 接入数据源(Source)
    这里用最简单的Socket端口作为数据源,模拟实时数据流;生产环境通常替换为Kafka、MySQL Binlog、日志采集器等。
  3. 数据转换计算(Transformation)
  • flatMap:数据清洗与格式转换,把一行文本拆成单个单词,并统一转为(单词, 1)的格式
  • keyBy:按单词进行分区分组,相同单词一定会被分发到同一个并行任务中处理,保证计数准确
  • sum:状态化聚合计算,Flink会在本地状态中自动保存每个单词的历史累计值,每来一条新数据就更新一次计数
  1. 结果输出(Sink)
    这里直接打印到控制台;生产环境通常输出到ClickHouse、Doris、Redis、Kafka等存储或下游系统。
  2. 触发执行
    Flink是懒执行机制:前面的代码只是构建计算逻辑图,不会真正运行;调用execute()后,才会生成作业图并提交给Flink集群启动运行。

三、运行效果演示

  1. 先在终端启动端口发送数据:
nc -lk 9999
  1. 启动Flink程序
  2. 在nc终端逐行输入文本,观察Flink控制台输出:
    • 输入 hello flink → 输出 (hello,1)(flink,1)
    • 再输入 hello world → 输出 (hello,2)(world,1)
    • 每输入一行,结果实时更新,计数自动累计

四、背后的完整执行链路

  1. 客户端:编译代码、生成作业图,提交给JobManager
  2. JobManager:负责任务调度、资源分配,把计算任务拆分分发到多个TaskManager
  3. TaskManager:真正执行计算,维护本地状态存储,数据在内存中流转计算
  4. 数据流转:数据从Source流入,经过各个算子逐次处理,最终输出到Sink,全程无磁盘IO(状态可配置持久化),实现毫秒级延迟

这个例子里的 .sum(1) 不是简单的内存变量循环累加,它背后是 Flink 核心的**键控状态(Keyed State)**机制在支撑——每个单词对应一个独立计数器,保存在计算节点本地,来一条数据就更新一次状态、输出一次最新结果。

下面拆解完整的累加原理和执行流程。


4.1、前提:keyBy 是累加的基础

在执行 sum 之前,代码先做了 .keyBy(tuple -> tuple.f0),这一步是分布式计数准确的核心:

  • 数据流按「单词」做哈希分区,相同的单词一定会被路由到同一个 TaskManager 的同一个并行子任务中处理
  • 每个子任务内,不同单词的计数状态完全隔离(hello 的计数和 flink 的计数互不干扰)
  • 只有 keyed 流才能使用键控状态,没有 keyBy 就无法做按维度的状态化累加

4.2、sum 算子的单条数据累加逻辑

sum 是 Flink 封装好的有状态聚合算子,底层本质是维护了一个 ValueState<Integer> 类型的键控状态,专门保存当前单词的累计值。

对于每一条到来的 (单词, 1) 数据,算子内部固定执行 4 步:

  1. 读状态:以当前单词为 key,从本地状态中读取已有的累计值;如果是该单词第一条数据,状态为空,默认取初始值 0
  2. 计算:将状态里的旧值 + 当前数据的计数字段(也就是 Tuple2 的第 2 位,值为 1),得到新累计值
  3. 写状态:把新的累计值写回本地状态,覆盖旧值
  4. 发结果:将最新的累计结果向下游输出(例子里就是打印到控制台)

注:.sum(1) 里的数字 1 是字段索引,代表对 Tuple2 的第 2 个字段(下标从 0 开始)做累加聚合。


4.3、结合例子走一遍完整流程

假设我们依次输入两行文本:

hello flink
hello world

整个累加过程按数据到来顺序执行:

序号 到来数据 读取状态值 计算新值 更新后状态 输出结果
1 (hello, 1) 0(首次无状态) 0+1=1 hello → 1 (hello, 1)
2 (flink, 1) 0(首次无状态) 0+1=1 flink → 1 (flink, 1)
3 (hello, 1) 1 1+1=2 hello → 2 (hello, 2)
4 (world, 1) 0(首次无状态) 0+1=1 world → 1 (world, 1)

你在控制台看到的每一行输出,都对应一次「状态更新 + 结果下发」,是纯流处理模式:来一条算一条,不攒批、不等窗口。


4.4、补充关键细节

  1. 状态存储位置
    默认存在 TaskManager 的 JVM 堆内存中,读写速度极快,支撑毫秒级延迟;数据量大时可配置 RocksDB 状态后端,将状态落盘到本地磁盘。
  2. 故障与容错
    默认重启任务后计数会清零(内存状态丢失)。只要开启 Flink 的 Checkpoint 机制,状态会定期做快照持久化到分布式存储(HDFS/S3等),任务故障重启后可从快照恢复,保证计数不重复、不丢失(Exactly-Once 语义)。
  3. 底层等价实现sum 只是语法糖,你完全可以手动用 RichFlatMapFunction + ValueState 实现完全一样的累加逻辑,效果完全等价。

五、关于状态存储位置

Flink 通过**状态后端(State Backend)**控制运行时计数的存储位置,同时通过 Checkpoint/Savepoint 机制将状态快照持久化到外部存储用于故障恢复。除了默认的 JVM 堆内存,还有多类存储方案,下面按「运行时读写存储」和「持久化快照存储」分开说明。


一、运行时状态存储(任务运行中实时读写的载体)

这是计数数据实时读写的位置,直接决定读写性能和单任务可承载的状态规模。

1. JVM 堆内内存(HashMapStateBackend)

  • 存储位置:TaskManager 进程的 JVM 堆内存,也就是你例子里的默认方式。
  • 特点:读写速度最快,无需序列化;但容量受堆内存限制,大状态会引发频繁 GC 甚至 OOM,进程重启后状态直接丢失。
  • 适用场景:小状态(MB 级)、本地测试、短周期计算。

2. 本地磁盘 + 堆外缓存(EmbeddedRocksDBStateBackend)

这是目前 Flink 生产环境最主流的运行时存储方案。

  • 存储位置:TaskManager 节点的本地磁盘目录,由嵌入式 KV 数据库 RocksDB 管理;热数据会缓存到堆外内存加速访问,冷数据自动落盘。
  • 核心特点
    • 容量上限取决于本地磁盘大小,支持 TB 级超大状态,远超出堆内存的限制
    • 不占用 JVM 堆内存,大幅降低 GC 压力,适合长周期不间断运行的任务
    • 读写需要序列化/反序列化,性能略低于堆内存,但完全满足生产级吞吐
  • 适用场景:生产环境长周期累计计数、大状态任务(如用户全量行为标签、天级以上的指标统计)。

注意:RocksDB 是本地磁盘存储,仅提升单任务状态容量,不提供跨节点故障容错;节点宕机后本地状态依然会丢失,容错需要依赖下面的持久化快照。


二、持久化快照存储(用于故障恢复、任务重启)

无论运行时状态存在内存还是本地磁盘,节点故障、任务停止都会丢失数据。Flink 通过 Checkpoint 机制定期将全量状态快照持久化到外部存储,保证故障后计数不重复、不丢失(Exactly-Once 语义)。

1. 分布式文件系统 / 对象存储(生产主流)

  • 常见存储介质:HDFS、阿里云 OSS、AWS S3、MinIO、腾讯云 COS 等
  • 工作方式:任务按固定间隔(如 1 分钟)异步将全量状态快照写入远端存储;任务故障重启后,自动从最新快照加载所有计数状态,无缝恢复。
  • 适用场景:所有生产环境的有状态任务。

2. 本地文件系统(仅测试)

快照保存在服务器本地磁盘,仅单节点可用,节点故障会丢失快照,仅适合本地开发调试。


三、补充:手动持久化 Savepoint

除了自动触发的 Checkpoint,你还可以手动触发 Savepoint,将当前计数状态完整备份到外部存储(路径与 Checkpoint 一致),主要用于任务代码升级、版本迭代后无缝恢复计数,相当于状态的「手动全量备份」。


结合你的单词计数例子改造

只需要在环境初始化时加两行配置,就能把默认内存计数改成「本地磁盘运行 + HDFS 持久化」的生产级方案:

import org.apache.flink.contrib.streaming.state.EmbeddedRocksDBStateBackend;

// 1. 配置运行时状态后端为 RocksDB(本地磁盘存储)
env.setStateBackend(new EmbeddedRocksDBStateBackend());
// 2. 配置 Checkpoint 快照持久化到分布式存储
env.getCheckpointConfig().setCheckpointStorage("hdfs:///flink/checkpoints/wordcount");
// 3. 开启每60秒一次的自动快照
env.enableCheckpointing(60000);

配置后,即使任务重启、节点宕机,每个单词的累计计数也能完整恢复。

更多推荐