飞熊出品 · Apache Flink 调研:26K stars 流处理事实标准,2.3 重新定义批流一体

41次阅读
飞熊出品 · Apache Flink 调研:26K stars 流处理事实标准,2.3 重新定义批流一体

副标题 :从 Exactly-Once 状态一致性 · 到 Flink Agents · 4 大类能力 + 6 大集成矩阵 + 数据栈闭环
赛道 :流处理 · 补足 338 全景图缺失的 ” 实时计算层 ”(Iceberg 379 之后的数据栈闭环)
作者 :yunying(增长运营官)
时间:2026-08-19


一、引子:Iceberg 之后为什么要写 Flink

8/19 发的《Apache Iceberg 调研》?p=379》讲清楚了 ” 数据怎么存 ”(湖仓格式),但 没讲 ” 数据怎么流进来 ”。数据从 Kafka → 数据湖(Iceberg)→ OLAP(Doris/StarRocks)→ BI 之间的 ” 实时计算 ” 层一直是空白。

Apache Flink 是这个空白的事实标准——12 年沉淀,26K stars,Apache 顶级,2026 还在大版本迭代(2.3.0 2026-06-25)。

本文补足:Flink = 流处理引擎 + 跟现有赛道的集成矩阵 + 实战跑通。

二、它解决什么问题(30 秒版)

传统实时计算有 3 个老问题:

  1. 只能 ” 至少一次 ” 或 ” 最多一次 ”——故障重跑可能重复或丢失
  2. 事件时间和处理时间不一致——凌晨 3 点的订单按业务时应该是昨天的
  3. 流批分离两套代码——同一业务写两次(Spark Streaming 写一份,Hive 写一份)

Flink = 解决这 3 个问题——Exactly-Once + Event Time + 流批统一 API。

实时数据(2026-08-19 拉取)

维度 数值
Stars 26,272
Forks 14,003
Open Issues 376(开源治理成熟,issue 量低)
License Apache-2.0
最新版本 2.3.0(2026-06-25)
最后 push 2026-08-19 14:20 UTC(今天!)
创建时间 2014-06-07(12 年
治理 Apache 顶级项目
Topics big-data / flink / java / python / scala / sql
官网 flink.apache.org
2026 新增 Flink Agents 0.3.1(2026-07-25)+ Native S3 FileSystem(2026-06-26)

关键判断 :Flink 是 流处理的事实标准——26K stars + 14K forks + 12 年 + Apache 顶级。Apache Beam / Spark Structured Streaming / Kafka Streams 都是 ” 参与者 ” 而非 ” 挑战者 ”。

三、4 大类核心能力

1. Correctness guarantees(正确性保证)

// Flink 标配的 Exactly-Once
env.enableCheckpointing(1000); // 每秒 checkpoint
env.getCheckpointConfig().setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE);
  • Exactly-Once 状态一致性——故障重跑不重复、不丢失
  • Event-Time 处理——按事件真实时间算窗口,不受处理延迟影响
  • 迟到数据处理——Watermark + Allowed Lateness,乱序数据也能进窗口

实战价值:金融交易 / 计费 / 审计场景必备 Exactly-Once。

2. Layered APIs(分层 API)

┌────────────────────────────────────────────────────────┐
│  L3: SQL on Stream & Batch Data(最高层)│
│      写 SQL 就能跑流处理(跟 Spark SQL 风格一致)│
├────────────────────────────────────────────────────────┤
│  L2: DataStream API(Java/Scala/Python)│
│      函数式 API:map / filter / keyBy / window / aggregate│
├────────────────────────────────────────────────────────┤
│  L1: ProcessFunction(最底层)│
│      精确控制时间 + 状态:自定义 Watermark / Trigger     │
└────────────────────────────────────────────────────────┘

实战价值:从 SQL 到 ProcessFunction,3 层 API 覆盖 95% 实时计算场景

3. Operational focus(运维能力)

  • Flexible Deployment——Standalone / YARN / Kubernetes / Mesos 全支持
  • High Availability——JobManager 主备自动切换
  • Savepoints——手动保存点(生产环境迁移 / 升级必备

4. Scalability + Performance(扩展 + 性能)

  • Scale-out architecture——加机器线性扩展
  • Very large state——TB 级别状态(用 RocksDB State Backend)
  • Incremental Checkpoints——增量 checkpoint,不全量重写
  • Low latency + High throughput——亚秒级延迟 + 百万 TPS

四、6 大集成矩阵(跟已发赛道怎么连)

集成目标 状态 相关文章
Apache Kafka ✅ 一等公民(Source + Sink) (待补)
Apache Iceberg ✅ 流式入湖(2026 重点) Iceberg 379
Apache Doris ✅ 流式写入(Streaming Sink) Doris 296
Apache StarRocks ✅ 流式写入 StarRocks 332
Apache Druid ✅ 实时 OLAP 输入 Druid+Pinot 336
Elasticsearch ✅ 全连接器 (待补)

实战数据栈闭环(飞熊组合拳)

┌─────────────┐   ┌─────────────┐   ┌─────────────┐
│   Kafka     │──▶│   Flink     │──▶│   Iceberg   │
│ 事件流入口  │   │ 流处理引擎  │   │ 数据湖存储  │
└─────────────┘   └─────────────┘   └─────────────┘
                                          │
                                          ▼
                              ┌───────────────────────┐
                              │  Doris / StarRocks   │
                              │  OLAP 实时查询       │
                              └───────────────────────┘
                                          │
                                          ▼
                              ┌───────────────────────┐
                              │  Superset / Metabase │
                              │  BI 仪表盘            │
                              └───────────────────────┘
                                          │
                                          ▼
                              ┌───────────────────────┐
                              │  AI Agent / Cube.js   │
                              │  自然语言查数据        │
                              └───────────────────────┘

五、Flink Agents(2026 新方向)

2026-07-25 Flink 团队发布了 Flink Agents 0.3.1——把 Agent 编排拉到 Flink 运行时。

核心能力

Stateful Agents——Agent 状态持久化(多轮对话不掉链)

Event-driven Actions——Agent 触发实时动作(IoT 报警 / 自动修复)

Flink 2.3 分布支持——无缝集成 Flink 生态

实战场景

客服 Agent——多轮对话 + 实时查订单(Iceberg)+ 触发发货动作

运维 Agent——监听 Kafka 日志 → 自动定位 + 调用 K8s API 修复

六、对比 Spark Structured Streaming / Kafka Streams

维度 Apache Flink Spark Structured Streaming Kafka Streams
Stars 26,272 41,000 (Spark) 7,800 (Kafka)
License Apache-2.0 ✅ Apache-2.0 ✅ Apache-2.0 ✅
延迟 亚秒级 微批(≥ 100ms) 亚秒级
Exactly-Once ✅ 标配 ⚠️ 需手动配 ✅ 标配
流批统一 ✅ 同一 API ✅ 同 Spark ❌ 仅流
SQL 能力 ✅ 完整 ✅ 完整 ⚠️ 弱(KSQL 替代)
生态 独立 + 集成 Kafka/Iceberg/Doris/StarRocks Spark 全家桶(SparkSQL/MLlib/Structured Streaming) 仅 Kafka 生态
学习曲线 中(流处理思维) (Spark 用户无缝切换) 低(Java 开发者)

选谁?

场景 推荐
亚秒级延迟 + Exactly-Once Flink
Spark 已有项目 Spark Structured Streaming
Kafka 全栈 + 轻量 Kafka Streams
AI Agent 流处理 Flink + Flink Agents 0.3
国内大厂(阿里 / 字节 / 美团) Flink ⭐(国内 90% 大厂实时计算首选)

七、实战 3 步 · 35 分钟跑通 ”Kafka → Flink → Iceberg”

Step 1 · Docker Compose 一键起 3 组件(10 分钟)

version: '3.8'
services:
  # Kafka 3.x
  kafka:
    image: apache/kafka:3.8.0
    environment:
      KAFKA_PROCESS_ROLES: broker,controller
      KAFKA_CONTROLLER_QUORUM_VOTERS: 1@localhost:9093
      KAFKA_NODE_ID: 1
    ports:
      - "9092:9092"

  # Flink JobManager + TaskManager
  jobmanager:
    image: apache/flink:2.3.0-java17
    command: jobmanager
    environment:
      FLINK_PROPERTIES: |
        jobmanager.rpc.address: jobmanager
    ports:
      - "8081:8081"  # Flink Web UI

  taskmanager:
    image: apache/flink:2.3.0-java17
    command: taskmanager
    depends_on:
      - jobmanager
    environment:
      FLINK_PROPERTIES: |
        jobmanager.rpc.address: jobmanager
        taskmanager.numberOfTaskSlots: 2

Step 2 · 写 Kafka → Iceberg 流处理作业(15 分钟)

FlinkKafka2Iceberg.java

StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.enableCheckpointing(60_000); // 60s checkpoint
env.getCheckpointConfig().setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE);

// Kafka Source
KafkaSource<String> source = KafkaSource.<String>builder()
    .setBootstrapServers("kafka:9092")
    .setGroupId("flink-consumer")
    .setTopics("orders")
    .setStartingOffsets(OffsetsInitializer.earliest())
    .setValueOnlyDeserializer(new SimpleStringSchema())
    .build();

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

// Iceberg Sink(关键:Flink 2.3+ 原生支持流式入湖)tableSink = TableUtils.createSink("iceberg://catalog/orders");

orders.map(order -> parseOrder(order))       // 解析 JSON
      .keyBy(Order::getCustomerId)          // 按客户分区
      .window(TumblingEventTimeWindows.of(Time.minutes(5)))  // 5 分钟窗口
      .aggregate(new OrderAggregator())     // 聚合 GMV
      .sinkTo(tableSink);                   // 写入 Iceberg

env.execute("Kafka to Iceberg");

Step 3 · 启动 + 验证(10 分钟)

# 1. 启动 Kafka + Flink
docker-compose up -d
sleep 30

# 2. 启动 Flink 作业
docker exec -it jobmanager ./bin/flink run \
  -c com.example.FlinkKafka2Iceberg \
  /opt/flink/usrlib/job.jar

# 3. 验证 Iceberg 表
docker exec -it jobmanager ./bin/sql-client.sh
# SQL> SELECT count(*) FROM iceberg.orders;

35 分钟跑通的成果:Kafka 实时事件 → Flink 流处理(Exactly-Once + 5 分钟窗口聚合)→ Iceberg 数据湖。

八、风险与坑

# 风险 应对
1 State Backend 选错(默认 HashMapStateBackend 不适合大状态) 生产用 RocksDBStateBackend(TB 级状态)
2 Checkpoint 频率 vs 性能 60s 是常见值,金融场景降到 10s
3 Savepoint 兼容性 Flink 大版本升级时 Savepoint 不一定能跨,需测试
4 Kafka Source 起始偏移 线上慎用 earliest,会全量重放
5 反压(Backpressure) 调 buffer 池 + 监控 TaskManager 指标
6 学习曲线(流处理思维) 团队 1-2 周上手,重点 Event Time + Window + Watermark

九、总结

3 个最值得用的理由

  1. 流处理事实标准——26K stars + 12 年 + Apache 顶级,国内 90% 大厂实时计算首选
  2. Exactly-Once + Event Time——金融交易 / 计费 / 审计场景必备
  3. 2026 新增 Flink Agents——把 AI Agent 流处理原生集成

1 句选型口诀

亚秒级 + Exactly-Once → Flink;Spark 用户 → Structured Streaming;纯 Kafka 轻量 → Kafka Streams。

国内大厂首选 Flink;AI Agent 时代 Flink + Flink Agents 是新方向。

飞熊数据栈闭环(今日成果)

Kafka 事件源
   ↓
Apache Flink(流处理 · 380)← 本文 P4
   ↓
Apache Iceberg(数据湖 · 379)← P3
   ↓
Doris / StarRocks / ClickHouse(OLAP)↓
Superset / Metabase(BI 仪表盘)↓
Cube.js / AI Agent(自然语言查询)

参考


📎 WordPress 链接

正文完