书接上文 《Kafka 从入门到实践》——Kafka 负责"消息去哪儿",Flink 负责"数据怎么变"。读完这篇,你将能够:用 Flink 消费 Kafka 数据,写自己的处理算子,把结果写回 Kafka。


一、Flink 是什么?为什么需要它?

1.1 一个场景

上篇博客我们搭好了 Kafka:上传 .dat 日志文件 → file-topic → 消费者保存到本地。

现在老板说:"别光存啊,给我分析每行有多少字符、多少单词,把结果存到另一个 topic。"

你当然可以在 Spring 消费者里加一段分析代码:

@KafkaListener(topics = "file-topic", ...)
public void onFile(ConsumerRecord<String, byte[]> record) {
    String content = new String(record.value());
    for (String line : content.split("\n")) {
        // 分析、统计、再发到 processed-topic...
    }
}

能跑。但问题是:

  • 文件大了怎么办?100MB 的文件在一个消费者线程里逐行处理,内存直接爆了

  • 想加个"每分钟统计日志级别数量"的需求?代码越来越臃肿

  • 三个消费者都要做同样的分析?代码复制三份

你需要一个专门干"数据处理"的框架——这就是 Flink。

1.2 Flink vs 自己写代码

自己在 Consumer 里写 用 Flink
大数据量 单线程扛不住,得自己写多线程 自动分布式并行处理
可靠性 崩溃后手动补数据 自动 checkpoint,重启后精确恢复
处理逻辑 跟消费逻辑耦合在一起 Source / Transform / Sink 清晰分离
扩展性 加机器要改代码 加个 TaskManager 就完事
复杂计算 自己实现窗口、聚合、join 内置 window、aggregate、connect

一句话:Kafka 负责可靠地传递数据,Flink 负责对数据做"实时计算"。


二、核心概念(用人话解释)

2.1 一张图看懂 Flink 的"三板斧"

 Source(数据入口)       Transform(转换)         Sink(数据出口)
      │                       │                       │
      ▼                       ▼                       ▼
┌──────────┐   byte[]   ┌──────────┐   String   ┌──────────┐
│  Kafka   │──────────▶│ flatMap  │──────────▶│  Kafka   │
│  Source  │            │ 逐行解析  │            │  Sink    │
└──────────┘            └──────────┘            └──────────┘

每个环节在 Flink 里叫一个算子(Operator)

算子类型 干什么的 类比
Source 从外部系统拉数据 水管工把水从水库引入管道
Transform 对数据做计算、转换 净水器,水进去→过滤→出来
Sink 把处理结果写回外部系统 水管工把水送到你家水龙头

2.2 几个关键概念

概念 一句话解释 类比
DataStream 无限的数据流 一条永不关闸的水管,数据就是水流
flatMap 一个输入,零到多个输出 拆包裹:一个大箱子进去,里面一个个小盒子出来
map 一个输入,恰好一个输出 翻译机:一句中文进,一句英文出
filter 通过条件才放行 筛子:大颗粒留下,小颗粒穿过
并行度 同时干活的人数 3 个人同时拧 3 个水龙头
Checkpoint 处理进度的快照 游戏存档,崩溃后读档继续
Watermark 事件时间的标记 迟到容忍度——超过这个时间的数据就不等了

2.3 算子的序列化要求(踩坑点)

Flink 会把你的算子对象序列化成字节,发送到不同机器上执行。所以:

  • ❌ 匿名内部类:.flatMap((s, out) -> { ... }) ← Flink 序列化不了 lambda

  • ❌ 非静态内部类:隐式持有外部类的 this 引用,外部类可能不可序列化

  • public static class:没有外部依赖,可以安全序列化

// ✅ 正确写法:必须 public static
public static class DatLineProcessor implements FlatMapFunction<byte[], String> {
    @Override
    public void flatMap(byte[] value, Collector<String> out) {
        // 你的处理逻辑
    }
}

三、环境准备

3.1 依赖

<properties>
    <flink.version>1.19.1</flink.version>  <!-- 用 1.19,别用 1.18,有坑 -->
</properties>
​
<dependencies>
    <!-- Flink 流处理核心 -->
    <dependency>
        <groupId>org.apache.flink</groupId>
        <artifactId>flink-streaming-java</artifactId>
        <version>${flink.version}</version>
        <!-- 排除 log4j-slf4j-impl,避免和 Spring Boot 的 Logback 抢 SLF4J 绑定 -->
        <exclusions>
            <exclusion>
                <groupId>org.apache.logging.log4j</groupId>
                <artifactId>log4j-slf4j-impl</artifactId>
            </exclusion>
        </exclusions>
    </dependency>
​
    <!-- Flink Kafka 连接器 -->
    <dependency>
        <groupId>org.apache.flink</groupId>
        <artifactId>flink-connector-kafka</artifactId>
        <version>3.2.0-1.19</version>  <!-- 版本号必须跟 Flink 版本匹配 -->
    </dependency>
​
    <!-- KafkaSource 需要这个基础库 -->
    <dependency>
        <groupId>org.apache.flink</groupId>
        <artifactId>flink-connector-base</artifactId>
        <version>${flink.version}</version>
    </dependency>
​
    <!-- 本地运行客户端 -->
    <dependency>
        <groupId>org.apache.flink</groupId>
        <artifactId>flink-clients</artifactId>
        <version>${flink.version}</version>
    </dependency>
</dependencies>

版本兼容性是个大坑,详见第七章。

3.2 前提条件

确保你已经按上篇博客搭好了 Kafka 集群(4 节点 KRaft),并且依赖里已经有 Spring Kafka。


四、实战:构建 Kafka → Flink → Kafka 管线

4.1 完整架构

┌──────────┐     ┌──────────┐     ┌──────────┐     ┌──────────┐
│ 用户上传  │     │  Kafka   │     │  Flink   │     │  Kafka   │
│ .dat 文件 │────▶│file-topic│────▶│ 流处理    │────▶│processed │
└──────────┘     │ (byte[]) │     │ 逐行分析   │     │ -topic   │
                 └──────────┘     └──────────┘     │(String)  │
                                                   └────┬─────┘
                                                        │
                                          ┌─────────────▼─────────────┐
                                          │ Spring Kafka 消费者        │
                                          │ 把分析结果存到内存,API 返回 │
                                          └───────────────────────────┘

4.2 第一步:写 Flink Job

public class DataFileProcessingJob implements Serializable {
​
    private final String bootstrapServers;
​
    public DataFileProcessingJob(String bootstrapServers) {
        this.bootstrapServers = bootstrapServers;
    }
​
    public void start() throws Exception {
        // 1. 创建本地执行环境(开发和测试用这一个就够了)
        Configuration flinkConfig = new Configuration();
        flinkConfig.set(RestOptions.BIND_PORT, "0");  // 随机端口,避免冲突
        final StreamExecutionEnvironment env =
                StreamExecutionEnvironment.createLocalEnvironment(1, flinkConfig);
​
        // 2. Kafka Source —— 从 file-topic 拉取 byte[] 消息
        KafkaSource<byte[]> source = KafkaSource.<byte[]>builder()
                .setBootstrapServers(bootstrapServers)
                .setTopics("file-topic")
                .setGroupId("flink-file-processor")
                .setStartingOffsets(OffsetsInitializer.latest())
                .setValueOnlyDeserializer(new RawBytesDeserializer())
                .build();
​
        DataStream<byte[]> rawStream = env.fromSource(
                source, WatermarkStrategy.noWatermarks(), "Kafka-Source");
​
        // 3. Transform —— flatMap 拆成逐行
        DataStream<String> processed = rawStream
                .flatMap(new DatLineProcessor())
                .name("Parse-Lines");
​
        // 4. Kafka Sink —— 分析结果写回 processed-topic
        KafkaSink<String> sink = KafkaSink.<String>builder()
                .setBootstrapServers(bootstrapServers)
                .setRecordSerializer(
                        KafkaRecordSerializationSchema.builder()
                                .setTopic("processed-topic")
                                .setValueSerializationSchema(new SimpleStringSchema())
                                .build())
                .build();
​
        processed.sinkTo(sink);
​
        // 5. 启动(这个方法会阻塞,直到任务被取消)
        env.execute("Flink-DAT-Processor");
    }
​
    // ==========================================
    // 算子类:必须是 public static
    // ==========================================
​
    /** 字节数组反序列化器 */
    public static class RawBytesDeserializer
            extends AbstractDeserializationSchema<byte[]> {
        @Override
        public byte[] deserialize(byte[] message) {
            return message;
        }
    }
​
    /** 逐行处理器 */
    public static class DatLineProcessor
            implements FlatMapFunction<byte[], String> {
        @Override
        public void flatMap(byte[] fileData, Collector<String> out) {
            String content = new String(fileData, StandardCharsets.UTF_8);
            int lineNumber = 0;
​
            for (String line : content.split("\\r?\\n")) {
                lineNumber++;
                if (line.trim().isEmpty()) continue;
​
                int chars = line.length();
                int words = line.trim().split("\\s+").length;
                String summary = line.length() > 50
                        ? line.substring(0, 50) + "..." : line;
​
                out.collect(String.format(
                        "L%d | 字:%d | 词:%d | %s",
                        lineNumber, chars, words, summary));
            }
        }
    }
}

4.3 第二步:在 Spring Boot 中管理 Flink 生命周期

env.execute()阻塞调用——它会一直占据当前线程直到任务被取消。如果在 Spring Boot 主线程里直接调用,整个应用就卡死了。

所以必须放在独立线程里运行:

@Service
public class FlinkJobService {
​
    @Value("${spring.kafka.bootstrap-servers}")
    private String bootstrapServers;
​
    private final AtomicBoolean running = new AtomicBoolean(false);
    private ExecutorService executor;
​
    /** 在独立线程中启动 Flink Job */
    public synchronized void startJob() {
        if (running.get()) return;
​
        running.set(true);
        executor = Executors.newSingleThreadExecutor(r -> {
            Thread t = new Thread(r, "flink-job");
            t.setDaemon(true);  // 守护线程,主线程退出时自动结束
            return t;
        });
​
        executor.submit(() -> {
            try {
                new DataFileProcessingJob(bootstrapServers).start();
            } catch (Exception e) {
                e.printStackTrace();
            } finally {
                running.set(false);
            }
        });
    }
​
    /** 停止 Flink Job */
    public synchronized void stopJob() {
        if (!running.get()) return;
        running.set(false);
        if (executor != null) executor.shutdownNow();
    }
​
    public boolean isRunning() { return running.get(); }
​
    @PreDestroy
    public void onShutdown() { stopJob(); }
}

关键点:

  • t.setDaemon(true) 保证 Spring Boot 关闭时 Flink 线程不会阻止 JVM 退出

  • shutdownNow() 会中断 env.execute() 的阻塞,让任务停止

4.4 第三步:消费处理结果

Flink 把分析结果写回了 processed-topic,再用 Spring Kafka 监听它:

@Service
public class FlinkResultConsumer {
​
    private final List<String> results = Collections.synchronizedList(new ArrayList<>());
​
    @KafkaListener(topics = "processed-topic", groupId = "flink-result-group")
    public void onResult(String result) {
        System.out.println("[Flink结果] " + result);
        results.add(result);
    }
​
    public List<String> getResults() {
        return new ArrayList<>(results);
    }
}

4.5 第四步:加个 Controller,用 REST 控制

@RestController
@RequestMapping("/api/flink")
public class FlinkController {
​
    private final FlinkJobService flinkJobService;
    private final FlinkResultConsumer resultConsumer;
​
    @PostMapping("/start")
    public Map<String, Object> start() {
        flinkJobService.startJob();
        return Map.of("status", "ok", "running", flinkJobService.isRunning());
    }
​
    @PostMapping("/stop")
    public Map<String, Object> stop() {
        flinkJobService.stopJob();
        return Map.of("status", "ok", "running", flinkJobService.isRunning());
    }
​
    @GetMapping("/status")
    public Map<String, Object> status() {
        return Map.of("running", flinkJobService.isRunning());
    }
​
    @GetMapping("/results")
    public Map<String, Object> results() {
        List<String> list = resultConsumer.getResults();
        return Map.of("count", list.size(), "results", list);
    }
​
    @DeleteMapping("/results")
    public Map<String, Object> clear() {
        resultConsumer.clearResults();
        return Map.of("status", "ok");
    }
}

4.6 跑起来

# 1. 启动应用
mvn spring-boot:run
​
# 2. 启动 Flink Job
curl -X POST http://localhost:8080/api/flink/start
​
# 3. 上传一个 .dat 文件
curl -X POST http://localhost:8080/api/file/upload -F "file=@sample.dat"
​
# 4. 查看 Flink 分析结果
curl http://localhost:8080/api/flink/results

结果长这样:

{
  "count": 7,
  "results": [
    "L1 | 字:57 | 词:6 | 2024-01-15 08:30:22 INFO 用户登录 userId=1001 ip=192.1...",
    "L2 | 字:76 | 词:7 | 2024-01-15 08:31:05 INFO 查询订单 orderId=ORD-2024-001...",
    "L3 | 字:76 | 词:7 | 2024-01-15 08:32:18 INFO 创建订单 orderId=ORD-2024-002...",
    "L4 | 字:91 | 词:8 | 2024-01-15 08:33:47 WARN 库存不足 productId=PROD-887 s...",
    "L5 | 字:81 | 词:7 | 2024-01-15 08:34:12 INFO 支付成功 orderId=ORD-2024-001...",
    "L6 | 字:75 | 词:7 | 2024-01-15 08:35:03 ERROR 支付失败 orderId=ORD-2024-00...",
    "L7 | 字:58 | 词:6 | 2024-01-15 08:36:21 INFO 用户登出 userId=1001 sessionT..."
  ]
}

7 行日志,每行都统计了字符数、单词数,超过 50 字符自动截断。这就是一条完整的 Kafka → Flink → Kafka 流处理管线。


五、Flink 常用算子速查

5.1 Transform 算子

算子 输入→输出 用途 示例
map 1→1 逐条转换 dataStream.map(s -> s.toUpperCase())
flatMap 1→N 一条拆成多条 dataStream.flatMap(new LineSplitter())
filter 1→0或1 过滤 dataStream.filter(s -> s.contains("ERROR"))
keyBy 分组 按 key 分流 dataStream.keyBy(Log::getLevel)
reduce 聚合 合并计算 stream.keyBy(...).reduce((a, b) -> a + b)
window 时间窗口 按时间段聚合 stream.keyBy(...).window(TumblingEventTimeWindows.of(Time.seconds(5)))

5.2 Source / Sink

类型 常见的
Source Kafka、文件、Socket、JDBC、自定义
Sink Kafka、文件、JDBC、Redis、Elasticsearch、自定义

5.3 如何选择

简单转换(一行变一行)     → map
一行拆成多行              → flatMap
只要某些数据              → filter
按某个字段分组统计        → keyBy + reduce/window
多流合并                  → connect / union

六、Flink 在 Spring Boot 中的架构细节

6.1 为什么要在独立线程运行?

Spring Boot 主线程                         Flink 守护线程
      │                                        │
      │ startJob() ───────────────────────────▶│ env.execute() ← 阻塞在这里
      │ (立即返回)                              │   │
      │                                        │   │ 轮询 Kafka,处理数据
      │ 处理 HTTP 请求                          │   │
      │ 响应 /api/flink/status                 │   │
      │                                        │   │
      │ stopJob() ─── shutdownNow() ───────────▶│ InterruptedException
      │                                        │
      ▼                                        ▼
   正常运行                                  安全退出

6.2 关键配置项

// 并行度 = 1:开发环境一个 TaskManager 就够了
StreamExecutionEnvironment.createLocalEnvironment(1, flinkConfig);
​
// 随机 REST 端口:避免多实例冲突
flinkConfig.set(RestOptions.BIND_PORT, "0");
​
// 起始消费位置:
OffsetsInitializer.earliest()  // 从最早的消息开始(适合开发测试,会重复消费)
OffsetsInitializer.latest()    // 只消费新消息(适合生产环境)

七、踩坑记录

坑1:flink-connector-kafka 版本不兼容

现象:Flink 启动后立刻停止,没有任何报错,也没有消费到任何数据。

原因

flink-connector-kafka:3.0.2-1.18  ← 针对 kafka-clients 3.4.0 编译
Spring Boot 3.3.0 引入的 kafka-clients  ← 实际 3.7.0

二进制不兼容导致 Flink 内部 Kafka 消费者静默失败。

解决:升级到兼容版本组合。

组件 ❌ 旧版本 ✅ 新版本
Flink 1.18.1 1.19.1
flink-connector-kafka 3.0.2-1.18 3.2.0-1.19
flink-connector-base 1.19.1(新增)

版本选择黄金法则

  • Flink 版本 = flink-connector-kafka 后缀版本(如 3.2.0-1.19 → Flink 1.19)

  • flink-connector-kafka 的 kafka-clients 版本 ≥ Spring Kafka 的 kafka-clients 版本

坑2:Logback 和 Log4j 打架

Flink 自带 Log4j,Spring Boot 用 Logback。两个日志框架抢 SLF4J 绑定。

解决:只排除 log4j-slf4j-impl,保留 log4j-apilog4j-core

<exclusion>
    <groupId>org.apache.logging.log4j</groupId>
    <artifactId>log4j-slf4j-impl</artifactId>
</exclusion>
<!-- 不要排除 log4j-api 和 log4j-core!Flink 需要它们 -->

坑3:算子类序列化失败

现象NotSerializableException

原因:用了非静态内部类或 lambda。

// ❌ 这样写 Flink 会报 NotSerializableException
DataStream<String> processed = rawStream.flatMap(
    (byte[] data, Collector<String> out) -> {
        // lambda 不能被 Flink 序列化
    }
);
​
// ✅ 必须单独定义一个 public static class
public static class DatLineProcessor
        implements FlatMapFunction<byte[], String> {
    @Override
    public void flatMap(byte[] value, Collector<String> out) { ... }
}

坑4:env.execute() 阻塞导致 Spring Boot 假死

现象:调用任何 API 都无响应。

原因:在主线程里直接调了 env.execute()

解决:放独立线程,设 setDaemon(true)


八、常见问题 FAQ

Q1: Flink 和 Spark Streaming 有什么区别?

Flink Spark Streaming
处理模型 真正的逐条流处理 微批次(攒一小批再处理)
延迟 毫秒级 秒级
状态管理 内置 State Backend,非常成熟 相对弱一些
学习曲线 中等 如果你已经会 Spark 就比较简单
适用场景 低延迟实时计算 Lambda 架构、与 Spark 生态集成

简单记:真正的实时用 Flink,准实时用 Spark Streaming。

Q2: Flink 能处理多大的数据量?

阿里用 Flink 处理双十一的每秒几十亿条交易消息。你的 .dat 文件完全不用担心。

Q3: 处理挂了怎么办?会丢数据吗?

Flink 有 Checkpoint 机制——定期把处理进度(当前读到哪个 offset、中间计算结果)存到持久化存储。挂了之后从最近的 checkpoint 恢复,不会丢数据。

本文用的是本地开发模式,没开 checkpoint。生产环境加几行配置就开启了:

env.enableCheckpointing(60_000);  // 每 60 秒做一次快照
env.getCheckpointConfig().setCheckpointStorage("file:///checkpoints");

Q4: Flink Job 一定要在 Spring Boot 里跑吗?

不一定。Flink 的标准部署方式是提交到 Flink Cluster(JobManager + 多个 TaskManager):

开发环境:本地跑(createLocalEnvironment)
生产环境:提交到 Flink 集群(flink run -c MainClass xxx.jar)

本文为了演示方便,嵌在了 Spring Boot 里。生产环境建议独立部署 Flink 集群。

Q5: 怎么知道 Flink 有没有在处理数据?

看控制台输出。正常运行的日志长这样:

[Flink-Source] 收到消息, 大小: 585 bytes
[Flink-Transform] 处理完成, 共 7 行
[Flink结果] L1 | 字:57 | 词:6 | 2024-01-15 08:30:22 INFO...

也可以通过 Flink 的 Web UI(默认 localhost:8081)查看详细的吞吐量、延迟指标。


九、学习路线图

第1步 ✅ 理解核心概念
    ├── Source → Transform → Sink
    ├── DataStream / flatMap / map / filter
    └── 并行度 / Checkpoint
​
第2步 ✅ 本地跑通第一条管线
    ├── KafkaSource → flatMap → KafkaSink
    ├── 理解序列化要求(public static class)
    └── Spring Boot 中管理 Flink 生命周期
​
第3步 ✅ 踩坑 & 调试
    ├── 版本兼容性(Flink + Kafka connector + kafka-clients)
    ├── 日志冲突(Log4j vs Logback)
    └── 算子序列化问题
​
第4步(进阶探索)
    ├── Window 窗口计算(每分钟聚合、滑动窗口)
    ├── KeyBy + Reduce 分组聚合
    ├── Checkpoint & Savepoint(故障恢复)
    ├── Event Time vs Processing Time(事件时间处理)
    ├── CEP(复杂事件匹配,如"30秒内连续3次登录失败")
    └── 提交到 Flink 集群(JobManager / TaskManager)

十、总结

回顾一下你在这篇博客里掌握了什么:

  • Flink 的定位:不是替代 Kafka,而是和 Kafka 配合——Kafka 负责传数据,Flink 负责算数据

  • 三板斧模式:Source 拉数据 → Transform 处理 → Sink 写结果,所有 Flink Job 都是这个套路

  • 算子必须 public static:因为 Flink 要序列化它们发到各台机器执行

  • Spring Boot 集成env.execute() 会阻塞,必须在独立守护线程里运行

  • 版本兼容:Flink 1.19.1 + flink-connector-kafka 3.2.0-1.19 + kafka-clients 3.7.0 是一组可用的版本组合

从 Kafka 到 Flink,你已经掌握了一条完整的实时数据处理链路:消息怎么传(Kafka)→ 数据怎么算(Flink)→ 结果怎么用(Spring Boot API)。剩下的就是在这个骨架上不断丰富业务逻辑。

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐