Skip to content

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

中文
0%
정확도 0%
💡 왼쪽 원문을 읽으면서 오른쪽에 따라 써보세요. Tab 키로 힌트를 받을 수 있습니다.

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

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

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

  • 实时流水线的成本是批处理的 3–10 倍
  • 大多数 BI 仪表盘即使延迟 15 分钟到 1 小时也没人抱怨
  • 真正需要实时的只有全部数据的 5–15%

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 把每个指标与表映射到层级

  • 库存指标: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章 · 四大流处理引擎对比

  • 真正基于事件时间(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韩国使用度特点
Flinkms–秒O业界标准,状态管理强
Spark SS秒–分O非常多与 Databricks 亲和
Kafka StreamsksqlDB一般Kafka 内置
RisingWavePostgres增加中运维简单,SaaS/OSS
MaterializeO少见Incremental view 是强项

3.7 选型指南

  • 事件时间・复杂状态:Flink
  • Databricks 生态:Spark SS
  • 以 Kafka 为中心・逻辑简单:Kafka Streams / ksqlDB
  • 只用 SQL 写流式应用:RisingWaveMaterialize

第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 模式

PostgresDebeziumKafkaFlinkIceberg
  • 事件保存在 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 的工作流

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

5.4 性能注意事项

  • Delete 文件堆积会拖慢读取性能
  • Compaction 的周期很关键(分钟–小时级)
  • 建议同时运行转换为 Copy-on-Write 的后台 job

第6章 · Streaming SQL 的崛起

6.1 为什么是 SQL

  • Flink Java API 学起来吃力
  • 数据工程师 + 分析师能够共同使用的通用语言
  • 容易与 dbt・Dagster 之类的工具集成
  • 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 年的现实并非如此:

작성 글자: 0원문 글자: 7,714작성 단락: 0/230