
副标题 :从 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 个老问题:
- 只能 ” 至少一次 ” 或 ” 最多一次 ”——故障重跑可能重复或丢失
- 事件时间和处理时间不一致——凌晨 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 个最值得用的理由
- 流处理事实标准——26K stars + 12 年 + Apache 顶级,国内 90% 大厂实时计算首选
- Exactly-Once + Event Time——金融交易 / 计费 / 审计场景必备
- 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(自然语言查询)
参考
- Apache Flink 官网:flink.apache.org
- GitHub 仓库:apache/flink
- Apache Flink 2.3.0 Release Notes:flink.apache.org/2026/06/25/
- Apache Iceberg 调研(379):《Apache Iceberg 调研》
- Apache Doris 调研(296):《Apache Doris 调研》
- Apache StarRocks 调研(332):《StarRocks 调研》
- Apache Druid + Pinot 横评(336):《Druid + Pinot 横评》
- BI 栈全景图 v2 (338):《2026 BI 栈全景图 v2》
📎 WordPress 链接
- 官方链接:《飞熊出品 · Apache Flink 调研:26K stars 流处理事实标准,2.3 重新定义批流一体》
- 短链:
https://east196.cn/?p={WP_POST_ID} - 状态:published · 2026-08-19