LabHub

博客

流式与批处理的重新定义:Flink・RisingWave・Materialize、CDC、Streaming SQL,实时的实用主义(2025)

한국어English日本語中文

Season 5 Ep 2 — 如果说 Ep 1 讲的是“数据存在哪里”,那么 Ep 2 讲的就是“数据流动得有多快”。从 2020–2023 的实时狂热,回到 2024–2025 的实用主义。

Prologue — “实时不是基本配置,而是一个选项”

2019–2022 年数据大会上的每一场主题演讲都在说“所有数据都必须变成实时的”。2025 年的现实并非如此:

2025 年的正确答案:

“按每份数据、每个指标的 SLA 划分鲜度层级,再为每一层选用合适的工具。”

本文就把这些层级与工具选型具体化。


第1章 · 鲜度层级(Freshness Tiers)

1.1 五层框架

层级延迟示例工具
Real-timems–秒交易监控・异常检测・fraudFlink, Kafka Streams
Near-real-time1–5 分钟运营仪表盘・alertFlink, RisingWave, Materialize
Fresh5–60 分钟库存・广告优化Streaming append + rollup
Daily24 小时BI, 报表Spark, dbt, SQL warehouse
Historical周/月分析・ML 训练Batch 年/月

1.2 把每个指标与表映射到层级

1.3 决策原则


第2章 · Lambda 与 Kappa 架构的 2025 版本

2.1 Lambda (2014)

2.2 Kappa (2014, Jay Kreps)

2.3 2025 年:Unified on Lakehouse

2.4“Streaming + Materialized”模式


第3章 · 四大流处理引擎对比

3.2 Spark Structured Streaming

3.3 Kafka Streams / ksqlDB

3.4 RisingWave

3.5 Materialize

3.6 对比表

引擎延迟复杂度SQL韩国使用度特点
Flinkms–秒O业界标准,状态管理强
Spark SS秒–分O非常多与 Databricks 亲和
Kafka StreamsksqlDB一般Kafka 内置
RisingWavePostgres增加中运维简单,SaaS/OSS
MaterializeO少见Incremental view 是强项

3.7 选型指南


第4章 · CDC (Change Data Capture)

4.1 为什么 CDC 是核心

4.2 实现方式

4.3 工具

4.4 CDC → Iceberg 模式

PostgresDebeziumKafkaFlinkIceberg

4.5 实务陷阱


第5章 · Iceberg v3 与实时 Upsert

5.1 Iceberg 版本沿革

5.2 Row-level delete 的两种方式

5.3 实时 Upsert 的工作流

  1. Flink 读取 CDC 事件
  2. 基于 PK 生成 Equality delete + insert
  3. Iceberg 把它反映到快照中
  4. 通过周期性 compaction 清理 delete 文件

5.4 性能注意事项


第6章 · Streaming SQL 的崛起

6.1 为什么是 SQL

6.3 ksqlDB

6.4 RisingWave 的 Postgres 兼容

6.5 Materialize


第7章 · 成本与延迟的权衡

7.1 成本构成

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 降本手段

7.4 按延迟目标推荐的架构


第8章 · 可观测性与调试

8.1 核心指标

8.2 观测工具

8.3 调试

8.4 告警


第9章 · 故障・恢复・SLA

9.1 SLA 设计

9.2 恢复策略

9.3 重跑

9.4 Multi-region


第10章 · 流式 + Lakehouse 实战模式

10.1 Medallion 之上的流式

10.2 CDC → Silver

10.3 事件溯源

10.4 实时 Feature Store

10.5 Streaming ETL 流水线


第11章 · 三个实战案例

11.1 电商订单流水线

11.2 金融交易监控

11.3 游戏遥测


第12章 · 韩国企业的流处理

12.1 传统模式

12.2 最新动向

12.3 合规考量

12.4 难点


第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 项


第15章 · 下一篇预告 — Season 5 Ep 3:“OLAP 引擎 2025 对比”

既然流式与批处理已经共享同一份存储,下一个问题就是“在它之上谁查询得最快”。

承认“一个引擎做不了所有事”这个 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 年的实用主义。

评论

还没有评论。

登录后即可发表评论