Season 5 Ep 2 — 如果说 Ep 1 讲的是“数据存在哪里”,那么 Ep 2 讲的就是“数据流动得有多快”。从 2020–2023 的实时狂热,回到 2024–2025 的实用主义。
- Prologue — “实时不是基本配置,而是一个选项”
- 第1章 · 鲜度层级(Freshness Tiers)
- 第2章 · Lambda 与 Kappa 架构的 2025 版本
- 第3章 · 四大流处理引擎对比
- 第4章 · CDC (Change Data Capture)
- 第5章 · Iceberg v3 与实时 Upsert
- 第6章 · Streaming SQL 的崛起
- 第7章 · 成本与延迟的权衡
- 第8章 · 可观测性与调试
- 第9章 · 故障・恢复・SLA
- 第10章 · 流式 + Lakehouse 实战模式
- 第11章 · 三个实战案例
- 第12章 · 韩国企业的流处理
- 第13章 · 十个反模式
- 第14章 · 检查清单 — 流式上线前 12 项
- 第15章 · 下一篇预告 — Season 5 Ep 3:“OLAP 引擎 2025 对比”
Prologue — “实时不是基本配置,而是一个选项”
2019–2022 年数据大会上的每一场主题演讲都在说“所有数据都必须变成实时的”。2025 年的现实并非如此:
- 实时流水线的成本是批处理的 3–10 倍
- 大多数 BI 仪表盘即使延迟 15 分钟到 1 小时也没人抱怨
- 真正需要实时的只有全部数据的 5–15%
2025 年的正确答案:
“按每份数据、每个指标的 SLA 划分鲜度层级,再为每一层选用合适的工具。”
本文就把这些层级与工具选型具体化。
第1章 · 鲜度层级(Freshness Tiers)
1.1 五层框架
| 层级 | 延迟 | 示例 | 工具 |
|---|---|---|---|
| Real-time | ms–秒 | 交易监控・异常检测・fraud | Flink, Kafka Streams |
| Near-real-time | 1–5 分钟 | 运营仪表盘・alert | Flink, RisingWave, Materialize |
| Fresh | 5–60 分钟 | 库存・广告优化 | Streaming append + rollup |
| Daily | 24 小时 | BI, 报表 | Spark, dbt, SQL warehouse |
| Historical | 周/月 | 分析・ML 训练 | Batch 年/月 |
1.2 把每个指标与表映射到层级
- 库存指标:Fresh
- 登录异常检测:Real-time
- 月度营收:Daily
- 模型训练用特征:Fresh + Daily 混合
1.3 决策原则
- 业务问题:“指标晚 1 小时会发生什么?”
- 答案是“什么都不会发生”→ Daily 或 Fresh
- 答案是“损失 10 万韩元”→ Near-real-time
- 答案是“有人会受伤/会有法律问题”→ Real-time
第2章 · Lambda 与 Kappa 架构的 2025 版本
2.1 Lambda (2014)
- Batch layer + Speed layer + Serving layer
- 同一套逻辑实现两遍 → 维护负担
- 2015–2020 年占据主导
2.2 Kappa (2014, Jay Kreps)
- 统一到一条流式日志(Kafka)上
- 重跑也通过流式完成
- 理论上简单,实务上困难
2.3 2025 年:Unified on Lakehouse
- 在 Iceberg/Delta 表之上,批处理与流式同时写入
- Flink + Iceberg, Spark Structured Streaming + Delta
- Serving 交给 OLAP 引擎或 materialized view
2.4“Streaming + Materialized”模式
- 源 → 流式流水线 → Lakehouse
- 高频查询用 Materialized view 预计算
- ad-hoc 由 OLAP 引擎直接查询
第3章 · 四大流处理引擎对比
3.1 Apache Flink
- 真正基于事件时间(Event time),有 watermark、exactly-once
- 状态(Keyed state)丰富,经过大规模生产验证
- 学习曲线陡峭,运维复杂
- 2024–2025 年 Flink CDC 2.x、Flink SQL 走向成熟
3.2 Spark Structured Streaming
- Spark 生态中的 micro-batch 流处理
- Databricks 的标准
- 延迟目标在 1–60 秒以上,并非真正的 ms 级
- SQL/Python 上手容易
3.3 Kafka Streams / ksqlDB
- Kafka 内置库,基于 JVM
- 轻量,专注处理 Kafka 中的数据时最合适
- 分布式运维需要自己实现
3.4 RisingWave
- 2022 年开源,专注 Streaming SQL
- 兼容 Postgres 的 SQL + Materialized view
- 状态放在独立存储(S3)— 运维简单
- 作为 Flink 的替代方案迅速崛起
3.5 Materialize
- Streaming + Incremental view maintenance
- 复杂 JOIN 与聚合自动保持最新
- 以企业级 SaaS 为主
3.6 对比表
| 引擎 | 延迟 | 复杂度 | SQL | 韩国使用度 | 特点 |
|---|---|---|---|---|---|
| Flink | ms–秒 | 高 | O | 多 | 业界标准,状态管理强 |
| Spark SS | 秒–分 | 中 | O | 非常多 | 与 Databricks 亲和 |
| Kafka Streams | 秒 | 中 | ksqlDB | 一般 | Kafka 内置 |
| RisingWave | 秒 | 低 | Postgres | 增加中 | 运维简单,SaaS/OSS |
| Materialize | 秒 | 中 | O | 少见 | Incremental view 是强项 |
3.7 选型指南
- 事件时间・复杂状态:Flink
- Databricks 生态:Spark SS
- 以 Kafka 为中心・逻辑简单:Kafka Streams / ksqlDB
- 只用 SQL 写流式应用:RisingWave 或 Materialize
第4章 · CDC (Change Data Capture)
4.1 为什么 CDC 是核心
- 把业务库(Postgres/MySQL)的变更实时复制到分析系统
- 相比批量导出,延迟与负载大幅下降
- 自动维护 Lakehouse 的 Silver 层
4.2 实现方式
- Log-based:Postgres WAL, MySQL binlog — 效率最高
- Trigger-based:用 DB 触发器生成变更事件 — 对数据库压力大
- Query-based:基于 updated_at 列的轮询 — 简单,但抓不到删除
4.3 工具
- Debezium (OSS):Postgres/MySQL/SQL Server/Oracle/Mongo → Kafka
- Fivetran, Airbyte:托管式 ELT,数百个源连接器
- Striim, HVR:企业级 CDC
- Flink CDC:在 Debezium 之上集成 Flink
4.4 CDC → Iceberg 模式
Postgres → Debezium → Kafka → Flink → Iceberg
- 事件保存在 Kafka 中 → 可以重跑
- Flink 用 MERGE INTO 向 Iceberg upsert
- 利用 Iceberg v2/v3 的 row-level delete
4.5 实务陷阱
- Schema change 的处理(新增列・删除列・类型变更)
- 初始快照 + 增量(Incremental snapshot)之间的协调
- Exactly-once 保证(同时防止重复与遗漏)
- 回填与重跑策略
第5章 · Iceberg v3 与实时 Upsert
5.1 Iceberg 版本沿革
- v1:Append-only
- v2:Row-level delete(position/equality)、MERGE INTO
- v3 (2024–2025):Deletion vectors, row lineage, V3 partition transforms
5.2 Row-level delete 的两种方式
- Position delete:指向某行所在文件与位置的删除文件
- Equality delete:基于条件的删除(实现 UPDATE 时很有用)
5.3 实时 Upsert 的工作流
- Flink 读取 CDC 事件
- 基于 PK 生成 Equality delete + insert
- Iceberg 把它反映到快照中
- 通过周期性 compaction 清理 delete 文件
5.4 性能注意事项
- Delete 文件堆积会拖慢读取性能
- Compaction 的周期很关键(分钟–小时级)
- 建议同时运行转换为 Copy-on-Write 的后台 job
第6章 · Streaming SQL 的崛起
6.1 为什么是 SQL
- Flink Java API 学起来吃力
- 数据工程师 + 分析师能够共同使用的通用语言
- 容易与 dbt・Dagster 之类的工具集成
6.2 Flink SQL
- 2023–2024 年趋于稳定
- 事件时间・窗口・状态全都可以用 SQL 表达
- 支持 UDF・UDAF
6.3 ksqlDB
- Kafka 之上的 Streaming SQL
- Table 与 Stream 的概念清晰
- 适合简单的 ETL 与聚合
6.4 RisingWave 的 Postgres 兼容
CREATE MATERIALIZED VIEW→ 实时保持最新- Postgres 工具(Grafana, Superset 等)可以直接使用
- 运维复杂度低 → 2025 年快速增长
6.5 Materialize
- 建立在 Incremental View Maintenance (IVM) 的研究之上
- 复杂 JOIN 也能自动保持最新
- 提供 SaaS + 自托管两种选择
第7章 · 成本与延迟的权衡
7.1 成本构成
- 计算(流式集群 24/7 运行)
- 状态存储(RocksDB/S3)
- Kafka/MSK 的运维
- 网络(跨区域)
7.2 典型成本对比(每月,中等规模)
| 选项 | 每月成本 | 延迟 |
|---|---|---|
| 批处理(Airflow + Spark,日级) | 低($1–5k) | 24 小时 |
| Micro-batch(5 分钟) | 中($3–10k) | 5 分钟 |
| Structured Streaming | 中–高($5–20k) | 秒–分 |
| Flink 集群 | 高($10–30k+) | ms–秒 |
| Managed(RisingWave/Confluent) | 中–高($7–25k) | 秒 |
7.3 降本手段
- 把实时层做薄(只保留核心指标)
- 评估能否用 Near-real-time 层替代
- 淡季缩容
- 状态放 S3(RisingWave 式)vs 本地 RocksDB(Flink)之间的权衡
7.4 按延迟目标推荐的架构
- 100ms 以内:In-memory stream processor + Redis/RocksDB
- 秒级:Flink/RisingWave + Kafka
- 分钟级:Spark Structured Streaming + Delta/Iceberg
- 小时级:Micro-batch Airflow + dbt
- 天级:Batch Spark
第8章 · 可观测性与调试
8.1 核心指标
- End-to-end latency(从事件发生到结果生效)
- Lag(Kafka consumer lag, Flink checkpoint lag)
- Throughput(eps, rps)
- Back-pressure(按 Flink task)
- 状态大小(RocksDB size, checkpoint size)
8.2 观测工具
- Flink UI + Prometheus + Grafana
- Kafka Lag Exporter, Burrow, Conduktor
- 用 OpenLineage 记录流式流水线的血缘
- Datadog/NewRelic APM + 流式 extensions
8.3 调试
- 基于 checkpoint 与 savepoint 的重跑
- CEP(Complex Event Processing)规则的调试
- 用采样日志追踪事件
- 测试使用 MiniCluster + embedded Kafka
8.4 告警
- Lag > 阈值 → Slack/PagerDuty
- End-to-end 延迟违反 SLO
- Checkpoint 连续失败
第9章 · 故障・恢复・SLA
9.1 SLA 设计
- Availability:99.9% = 每月允许 43 分钟的宕机
- Freshness:事件发生后 X 秒内可查询
- Correctness:最终一致的正确性(eventual)vs exactly-once
9.2 恢复策略
- Flink:checkpoint(定期)、savepoint(手动)
- Kafka:Replication factor 3, min.insync.replicas 2
- Iceberg:用快照 + 分支回滚
9.3 重跑
- Kafka retention 要覆盖重跑周期(7–30 天很常见)
- 或者从 Iceberg 原始数据 → 重新运行流式 job
- 重跑期间防止重复(幂等性设计)
9.4 Multi-region
- Kafka MirrorMaker 2 / Confluent Cluster Linking
- Iceberg 依靠存储复制(S3 Cross-region)
- 流式作业通常采用 Active-Passive
第10章 · 流式 + Lakehouse 实战模式
10.1 Medallion 之上的流式
- Bronze:append 原始事件
- Silver:清洗・去重・JOIN upsert
- Gold:聚合・指标
10.2 CDC → Silver
- 数据库变更事件 → Flink → Iceberg Silver upsert
- Silver 承担业务库的“副本 + 历史”角色
- Gold 由 dbt/Spark 批处理生成
10.3 事件溯源
- 把领域事件永久保存在 Kafka 中
- 状态通过重放事件来重建
- 在 Iceberg Bronze 中长期保存
10.4 实时 Feature Store
- Feast/Tecton + Kafka + Iceberg
- Online feature(Redis)+ Offline feature(Iceberg)
- 监控 Online-offline skew
10.5 Streaming ETL 流水线
- 源 → Bronze(append)→ Silver(清洗)→ Gold(聚合)
- 每个阶段用 Flink/Spark SS 作业实现
- 重跑只需重置上游 offset
第11章 · 三个实战案例
11.1 电商订单流水线
- 订单事件 Kafka → Flink → 库存更新 + 异常检测 + 分析
- 库存:秒级实时
- 营收指标:5 分钟的 near-real-time
- 月度报表:Daily batch
11.2 金融交易监控
- 交易事件 Kafka → Flink CEP → Fraud score
- 要求延迟低于 100ms
- 状态:每位用户的交易历史(Flink keyed state)
- 结果:保存到 Redis + Iceberg
11.3 游戏遥测
- 客户端事件 Kinesis/Kafka → Flink/Spark SS
- 实时分析(同时在线人数・DAU)属于 Near-real-time
- 详细日志堆积到 Iceberg Bronze,之后再分析
- A/B 测试:Gold 聚合
第12章 · 韩国企业的流处理
12.1 传统模式
- 金融:Tibco EMS, IBM MQ, Kafka 混用
- 通信:Charging/Billing 使用实时流式
- 游戏:Kafka + Flink 或 Kafka Streams
- 电商:Kafka + Spark SS + ELK
12.2 最新动向
- Confluent Cloud / MSK 的采用在增加
- Flink 的采用面扩大(尤其是 Toss・Coupang・Naver・Kakao)
- RisingWave・Materialize 在 2024–2025 年仍处于导入初期
12.3 合规考量
- 金融:在网络隔离环境中自建 Kafka 集群
- 个人信息:CDC 事件的 PII 脱敏是必须的
- 审计日志:长期保存・不可篡改
12.4 难点
- 数据工程师人手不足 → 越来越倾向托管服务
- 与遗留数仓共存
- 24/7 on-call 文化仍在形成中
第13章 · 十个反模式
13.1“把一切都做成实时”
连不需要的表也做成流式 → 成本与复杂度激增。
13.2 迷信 Exactly-once
要在源端与汇端同时保证端到端 exactly-once 并不容易。幂等设计是必须的。
13.3 省略 CDC 的初始快照
出现遗漏,准确性下降。
13.4 没有模式变更的自动传播
下游流水线被打断。
13.5 Kafka retention 太短
无法重跑。
13.6 Checkpoint 周期太长
故障时的恢复成本与重跑量激增。
13.7 状态无限保留
Flink keyed state 无限增长 → OOM。
13.8 不做 Delete 文件的 compaction
Iceberg 的读取性能下降。
13.9 内存与 CPU 分配不足
Back-pressure 连锁反应。
13.10 缺少可观测性与告警
事故由客户先发现。
第14章 · 检查清单 — 流式上线前 12 项
- 鲜度层级映射(按表、按指标的 SLA)
- 引擎选型依据(Flink/Spark SS/RisingWave/Materialize)
- CDC 源的选定 + 初始快照策略
- Kafka retention + 重跑计划
- Iceberg/Delta 表设计 + compaction 自动化
- Checkpoint 与 savepoint 策略
- 幂等性与重复处理的设计
- 可观测性(lag・latency・throughput・back-pressure)
- SLA/SLO 的定义 + 告警
- 成本仪表盘(计算・存储・网络)
- 灾难恢复与 Multi-region 计划
- On-call 与故障响应手册
第15章 · 下一篇预告 — Season 5 Ep 3:“OLAP 引擎 2025 对比”
既然流式与批处理已经共享同一份存储,下一个问题就是“在它之上谁查询得最快”。
- DuckDB:单节点 OLAP 的革命
- ClickHouse:实时 OLAP 的标准
- Snowflake / BigQuery / Redshift:托管式巨头
- Databricks SQL / StarRocks / Doris / Pinot / Druid
- Trino / Presto 的联邦查询
- 真实基准测试中的陷阱
- 成本 vs 延迟 vs 运维负担
- MPP 与单节点之间的边界
- 面向韩国企业的选型指南
- “因地制宜”的引擎布局模式
承认“一个引擎做不了所有事”这个 2025 年的现实之后,才真正有意思。
下一篇文章再见。
总结:2025 年的流处理已经从“全部实时”重新定义为“基于 SLA 的鲜度层级”。按照 Real-time / Near-real-time / Fresh / Daily / Historical 五个层级去布置 Flink・Spark SS・RisingWave・Materialize・ksqlDB,用 CDC 把业务库的变化流向 Lakehouse,再用 Iceberg v3 的 row-level delete 处理实时 upsert。不是 Lambda/Kappa,而是“Unified on Lakehouse”成为主导模式,成本、延迟与复杂度的权衡需要被有意识地设计。“实时不是基本配置,而是一个选项” — 这就是 2025 年的实用主义。
현재 단락 (1/230)
2019–2022 年数据大会上的每一场主题演讲都在说“所有数据都必须变成实时的”。2025 年的现实并非如此: