大数据技术硬核教程(中篇)
实时计算 · 湖仓一体 · 数据建模
承上启下:上篇解决了"数据放哪、怎么算",中篇解决"数据怎么实时流动、怎么在湖上做事务、怎么组织成可用的资产"。 难度提升:本篇包含分布式快照算法、MVCC 实现、Exactly-Once 语义证明等硬核内容。 实操环境:EMR(Hadoop 3.2.1 + Spark 3.5.3 + Hive 3.1.3)+ 本篇新增 Flink / Kafka / Iceberg。
目录
- 第 8 章:流处理的本质 —— 从批到流的范式跃迁
- 第 9 章:Flink 内核解剖 —— 时间、状态、容错
- 第 10 章:Exactly-Once 的真相
- 第 11 章:Kafka 深度剖析 —— 数据总线的工程美学
- 第 12 章:湖仓一体 —— Iceberg / Hudi / Delta 三国志
- 第 13 章:CDC 全链路 —— 让数据库变成流
- 第 14 章:数据建模 —— 被低估的核心竞争力
- 第 15 章:实时 OLAP 引擎对决
- 第 16 章:终极实战 —— 端到端实时数仓
第 8 章:流处理的本质 —— 从批到流的范式跃迁
8.1 一个颠覆性的认知:批是流的特例
传统认知:
批处理:处理一堆静态数据
流处理:处理源源不断的数据
→ 两种不同的东西Flink 的世界观(也是现代流处理的共识):
世界上只有一种数据:无界流(Unbounded Stream)
批处理 = 对无界流开了一个【有界窗口】
实时处理 = 对无界流开了一个【滑动的小窗口】
→ 批是流的特例,不是流是批的加速版这个认知为什么重要?
因为它决定了架构走向。如果你认为批和流是两种东西,你会建两套系统(Lambda 架构);如果你认为批是流的特例,你会建一套系统(Kappa / 流批一体)。
8.2 三代架构的血泪史
第一代:Lambda 架构(2011-2018)
┌─────────────────────┐
│ 批处理层(准确) │
┌────►│ Spark/MR → HDFS │────┐
│ │ T+1 全量重算 │ │
┌──────────┐ │ └─────────────────────┘ │ ┌──────────┐
│ 数据源 │──┤ ├──►│ 服务层 │
│ Kafka │ │ ┌─────────────────────┐ │ │ 合并结果 │
└──────────┘ │ │ 速度层(快但可能错) │ │ └──────────┘
└────►│ Storm/Flink → Redis│────┘
│ 秒级近似结果 │
└─────────────────────┘致命问题:
- 同一套业务逻辑要写两遍(Spark 一份 + Flink 一份)
- 两套代码的口径几乎必然会漂移,对不上数
- 运维两套集群,成本翻倍
- 改一次需求要改两处,还要保证一致
💬 真实案例:某大厂曾出现批流结果差异 3%,排查两周才发现是流层用了
<而批层用了<=。
第二代:Kappa 架构(2014 提出)
┌──────────┐ ┌─────────────────────────┐ ┌──────────┐
│ 数据源 │───►│ Kafka(长期保留) │───►│ Flink │───► 结果
└──────────┘ │ 可回溯到任意历史位点 │ └──────────┘
└─────────────────────────┘
│
需要重算时:从头重放核心思想:既然流能表达一切,那就只用流。需要"批处理"时,把 Kafka 从头重放一遍。
理想很美好,现实问题:
- Kafka 保存全量历史数据成本极高(Kafka 用本地磁盘,不是廉价存储)
- 重放 3 年数据可能要跑好几天
- 流式 SQL 表达能力当年不如批(复杂 JOIN、多维分析很吃力)
- 历史数据的 Schema 演进很难处理
第三代:湖仓流批一体(2020- 现在)
┌──────────────────────────────┐
│ 统一存储:数据湖表格式 │
┌──────────┐ │ Iceberg / Hudi / Delta │
│ 数据源 │ │ ✓ ACID ✓ 时间旅行 │
│ MySQL │────┐ │ ✓ 增量读取(关键!) │
│ 日志 │ │ └──────┬──────────────┬────────┘
│ Kafka │ │ ▲ │
└──────────┘ │ │ │
│ ┌──────┴──────┐ ┌────▼─────────┐
└───►│ Flink 写入 │ │ 读取层 │
│ (流式) │ │ Flink(流读) │
└─────────────┘ │ Spark(批读) │
│ Trino(即席) │
│ StarRocks(OLAP)│
└──────────────┘关键突破点:表格式层(Iceberg/Hudi)提供了 增量读取(Incremental Read) 能力——同一张表,既能全量批读,也能只读新增的 commit,实现"流读一张表"。
这一代的核心价值:
- 一份数据、一份存储、一套 SQL
- 流处理和批处理读的是同一张表,口径天然一致
- 存储用廉价对象存储,成本可控
- 重算时可以用 Spark 批处理高效跑(不用一条条重放)
📌 2026 年的现状:大厂基本都在往第三代走,但完全的"流批一体"仍有大量工程难题(延迟差异、状态一致性、资源隔离)。这是目前最活跃的技术方向之一。
8.3 流处理的三个核心难题
任何流处理引擎都必须回答这三个问题:
难题一:时间是什么?(Time Semantics)
一条订单数据:
用户在 10:00:00 点击下单 ← Event Time(事件时间)
数据 10:00:03 到达 Kafka ← Ingestion Time(摄入时间)
Flink 10:00:07 处理到它 ← Processing Time(处理时间)
问:"10:00-10:01 这一分钟有多少订单?" —— 按哪个时间算?答案必须是 Event Time,否则:
- 网络抖动会导致统计结果变化
- 任务重启后重算,结果和第一次不一样(不可重现!)
- 数据迟到就永远统计不到
但用 Event Time 带来新问题:数据乱序。
难题二:数据乱序了怎么办?(Out-of-Orderness)
到达顺序: [10:00:01] [10:00:05] [10:00:02] [10:00:03] [10:00:09]
↑ 迟到了
问:处理到 10:00:05 时,能不能关闭 "10:00:00-10:00:03" 这个窗口?
如果关了,后面来的 10:00:02 和 10:00:03 就丢了
如果不关,要等到什么时候?永远等下去吗?这就是 Watermark(水位线) 要解决的问题——它是一个关于"完整性"的启发式判断。
难题三:状态怎么办?(State Management)
批处理:任务失败 → 重跑一遍 → 结果一样(幂等)
流处理:任务跑了 3 天,内存里累积了 500GB 的聚合状态
突然挂了 → 重新从头算?3 天的数据要重跑!
→ 必须能【保存和恢复】状态这引出了 Checkpoint(检查点) 机制。
第 9 章:Flink 内核解剖 —— 时间、状态、容错
9.1 Watermark:流处理中最精妙的设计
本质定义
Watermark(t) = 一个断言:"我认为所有 EventTime <= t 的数据都已经到齐了"它是一条特殊的数据记录,和业务数据一起在流中流动。
工作机制图解
数据流(从右往左流动):
... [E:10:08] [W:10:05] [E:10:06] [E:10:03] [W:10:02] ...
↑ ↑
Watermark=10:05 Watermark=10:02
当 Watermark=10:05 流过某个窗口算子时:
→ 所有 endTime <= 10:05 的窗口被触发计算并关闭
→ 窗口 [10:00, 10:05) 触发输出
→ 之后再来 EventTime=10:04 的数据 = 迟到数据生成策略实操
// Flink 1.18+ 标准写法
DataStream<Order> stream = env
.fromSource(kafkaSource, WatermarkStrategy
// ① 允许 5 秒乱序(核心参数)
.<Order>forBoundedOutOfOrderness(Duration.ofSeconds(5))
// ② 指定从哪个字段取 Event Time
.withTimestampAssigner((order, ts) -> order.getCreateTime())
// ③ 空闲分区处理(重要!否则某个分区没数据会卡住全局 watermark)
.withIdleness(Duration.ofMinutes(1)),
"kafka-source");计算公式:
Watermark = 当前观察到的最大 EventTime - 允许乱序时间(maxOutOfOrderness) - 1ms⚠️ Watermark 的三大经典坑
坑 1:多并行度下取最小值
Source 并行度 = 3
Subtask 0: Watermark = 10:05
Subtask 1: Watermark = 10:03 ← 下游算子的 Watermark 取这个最小值
Subtask 2: Watermark = 10:07
原因:必须保证所有上游都到齐了才能推进,否则会丢数据
后果:★ 一个慢的分区会拖住整个作业的时间推进 ★坑 2:某个分区没数据 → Watermark 永远不推进
Kafka Topic 有 6 个分区,但业务低峰期只有 2 个分区有数据
→ 另外 4 个分区的 Watermark 停在很久之前
→ 全局 Watermark 取最小值 → 卡死,窗口永远不触发!
解法:withIdleness(Duration.ofMinutes(1))
标记空闲分区,计算全局 Watermark 时忽略它坑 3:乱序时间设置的两难
设置太小(如 1s)→ 大量数据迟到被丢弃 → 数据不准
设置太大(如 10min)→ 窗口迟迟不触发 → 延迟高
★ 这是【延迟】和【完整性】的根本权衡,没有银弹 ★工业界的三层兜底方案:
stream
.keyBy(Order::getUserId)
.window(TumblingEventTimeWindows.of(Time.minutes(1)))
// 第 1 层:Watermark 允许 5s 乱序
// 第 2 层:允许迟到 1 分钟,迟到数据触发窗口【增量更新】
.allowedLateness(Time.minutes(1))
// 第 3 层:超过 1 分钟的极端迟到数据,输出到侧流单独处理
.sideOutputLateData(lateTag)
.aggregate(new OrderAggregator());
// 侧流数据可以写到另一张表,T+1 时合并修正
DataStream<Order> lateStream = result.getSideOutput(lateTag);
lateStream.sinkTo(lateDataSink);💭 思考题 9.1:为什么
allowedLateness会增加状态大小?(提示:窗口关闭后状态才能清理)在什么场景下它会导致 OOM?
9.2 状态管理:流处理的记忆
状态的两大分类
┌─────────────────────────────────────────────────────┐
│ Keyed State(键控状态)—— 最常用 │
│ 只能在 keyBy() 之后使用,每个 key 独立一份状态 │
│ ├─ ValueState<T> 单值(如:用户最后登录时间) │
│ ├─ ListState<T> 列表(如:用户最近 10 次点击) │
│ ├─ MapState<K,V> 映射(如:用户各品类购买次数) │
│ ├─ ReducingState<T> 自动聚合(省内存) │
│ └─ AggregatingState<I,O> 自动聚合,输入输出类型可不同 │
├─────────────────────────────────────────────────────┤
│ Operator State(算子状态)—— 主要框架内部用 │
│ 每个算子并行实例一份 │
│ 典型用途:Kafka Source 保存各分区的 offset │
└─────────────────────────────────────────────────────┘实战:用状态实现"用户连续 3 次下单未支付"预警
public class PayAlertFunction extends KeyedProcessFunction<Long, Order, Alert> {
// 声明状态:记录连续未支付次数
private transient ValueState<Integer> unpaidCount;
private transient ValueState<Long> timerTs;
@Override
public void open(Configuration cfg) {
ValueStateDescriptor<Integer> desc =
new ValueStateDescriptor<>("unpaid-count", Integer.class);
// ★ 关键:配置 TTL,否则状态无限膨胀
StateTtlConfig ttl = StateTtlConfig
.newBuilder(Time.days(7))
.setUpdateType(StateTtlConfig.UpdateType.OnCreateAndWrite)
.setStateVisibility(StateTtlConfig.StateVisibility.NeverReturnExpired)
.cleanupInRocksdbCompactFilter(1000) // RocksDB 后台清理
.build();
desc.enableTimeToLive(ttl);
unpaidCount = getRuntimeContext().getState(desc);
timerTs = getRuntimeContext().getState(
new ValueStateDescriptor<>("timer", Long.class));
}
@Override
public void processElement(Order order, Context ctx, Collector<Alert> out)
throws Exception {
if ("pay".equals(order.getBehavior())) {
unpaidCount.clear(); // 支付了,清零
clearTimer(ctx);
return;
}
if ("order".equals(order.getBehavior())) {
Integer cnt = unpaidCount.value();
cnt = (cnt == null ? 0 : cnt) + 1;
unpaidCount.update(cnt);
if (cnt >= 3) {
out.collect(new Alert(order.getUserId(), "连续3次下单未支付", cnt));
unpaidCount.clear();
}
// 注册 30 分钟定时器,超时自动清理
long ts = ctx.timerService().currentProcessingTime() + 30 * 60 * 1000;
ctx.timerService().registerProcessingTimeTimer(ts);
timerTs.update(ts);
}
}
@Override
public void onTimer(long ts, OnTimerContext ctx, Collector<Alert> out)
throws Exception {
unpaidCount.clear(); // 超时清理状态,防止无限增长
timerTs.clear();
}
private void clearTimer(Context ctx) throws Exception {
Long ts = timerTs.value();
if (ts != null) {
ctx.timerService().deleteProcessingTimeTimer(ts);
timerTs.clear();
}
}
}🎯 状态管理三铁律:
- 必设 TTL——不设 TTL 的状态是定时炸弹,跑几个月必 OOM
- 必用定时器清理——TTL 是兜底,业务逻辑的主动清理更及时
- 状态越小越好——Checkpoint 时要序列化,状态大直接导致 Checkpoint 超时
状态后端选型(生产关键决策)
| 状态后端 | 存储位置 | 容量上限 | 访问延迟 | 适用场景 |
|---|---|---|---|---|
| HashMapStateBackend | JVM 堆内存 | 受堆大小限制(<10GB) | 极快(纳秒) | 小状态、低延迟要求 |
| EmbeddedRocksDBStateBackend | 本地磁盘(RocksDB) | TB 级 | 较慢(微秒,需序列化) | 大状态生产标配 |
RocksDB 调优(生产必备):
# flink-conf.yaml
state.backend: rocksdb
state.backend.incremental: true # ★ 增量 Checkpoint,必开
state.backend.rocksdb.memory.managed: true # 由 Flink 统一管理内存
state.backend.rocksdb.memory.fixed-per-slot: 512mb
# 使用预定义调优模板
state.backend.rocksdb.predefined-options: SPINNING_DISK_OPTIMIZED_HIGH_MEM
# 可选值:DEFAULT / SPINNING_DISK_OPTIMIZED / SPINNING_DISK_OPTIMIZED_HIGH_MEM / FLASH_SSD_OPTIMIZED
# Block Cache 与 Write Buffer
state.backend.rocksdb.block.cache-size: 256mb
state.backend.rocksdb.writebuffer.size: 64mb
state.backend.rocksdb.writebuffer.count: 4
# 使用多块磁盘分散 IO(重要!)
state.backend.rocksdb.localdir: /mnt/disk1/rocksdb,/mnt/disk2/rocksdb📐 为什么 RocksDB 用 LSM-Tree 而不是 B+Tree?
B+Tree: 写需要原地更新 → 随机写 → 磁盘上很慢 读快(O(log n) 次随机读) LSM-Tree:写只追加到内存 MemTable,满了顺序刷盘 → 【顺序写】→ 极快 读需要查多层 SSTable → 用 Bloom Filter 快速排除 流处理的特征:写远多于读(每条数据都要更新状态) → LSM-Tree 的写优化正好对症代价:读放大和空间放大。Compaction(合并)会消耗大量 CPU 和 IO,这是 RocksDB 调优的核心矛盾。
9.3 Checkpoint:分布式快照的艺术
问题的难度在哪?
一个 Flink 作业有 100 个算子实例,分布在 20 台机器上
每个实例都有自己的状态
问:怎么给这 100 个实例拍一张【一致的】快照?
难点:不能停止整个作业(stop-the-world 代价太大)
但如果不停,A 实例快照时处理到第 1000 条,
B 实例快照时处理到第 1050 条 → 不一致!Chandy-Lamport 算法:1985 年的智慧
Flink 的 Checkpoint 基于 Chandy-Lamport 分布式快照算法的变体(称为 ABS - Asynchronous Barrier Snapshotting)。
核心思想:用一个特殊标记(Barrier)在流中传播,把数据流"切"成快照前后两部分。
步骤演示(单输入):
时刻 T0:JobManager 向所有 Source 注入 Barrier-N
┌────────┐
│ Source │ ──[数据][数据]──▶
└────────┘
时刻 T1:Source 快照自己的状态(Kafka offset),发出 Barrier
┌────────┐
│ Source │ ──[数据][Barrier-N][数据]──▶
└────────┘ ↑ 已快照 offset=1000
时刻 T2:算子收到 Barrier
┌────────┐
│ Map │ 收到 Barrier-N
└────────┘ → ① 暂停处理
→ ② 异步快照自己的状态到 HDFS/OSS
→ ③ 把 Barrier 转发给下游
→ ④ 继续处理
时刻 T3:所有 Sink 都确认收到 Barrier-N
→ JobManager 标记 Checkpoint-N 完成Barrier 对齐(Barrier Alignment)—— 多输入的关键
当一个算子有多个输入流时(如 Join、Union):
输入流 A: ──[a3][a2][Barrier-N][a1]──▶ ┐
├─▶ [Join 算子]
输入流 B: ──[b5][b4][b3][b2][Barrier-N]▶ ┘
Barrier 从 A 先到达:
→ 算子【阻塞】输入流 A(把 a2、a3 缓存起来)
→ 继续处理输入流 B,直到 B 的 Barrier 也到
→ 两个 Barrier 都到齐 → 开始快照 → 解除阻塞
★ 这就是"对齐":保证快照的是同一个逻辑时刻的状态 ★对齐的代价:如果 A 和 B 的 Barrier 到达时间差很大(数据倾斜/反压),阻塞时间长 → 延迟飙升。
Unaligned Checkpoint(Flink 1.11+):反压场景的救星
Aligned(对齐):
等所有 Barrier 到齐 → 反压时可能等几分钟 → Checkpoint 超时失败
Unaligned(非对齐):
第一个 Barrier 到达就立即快照
把【正在传输中的数据】(in-flight data)也存进 Checkpoint
→ Checkpoint 时间不受反压影响
→ 代价:Checkpoint 体积变大# 生产配置建议
execution.checkpointing.unaligned: true
execution.checkpointing.aligned-checkpoint-timeout: 30s
# 含义:先尝试对齐,30 秒内没对齐完就自动切换到非对齐模式
# ★ 这是两种模式的最佳折中 ★Checkpoint 完整配置模板
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
CheckpointConfig cc = env.getCheckpointConfig();
// 间隔:1 分钟(根据可容忍的恢复时间设定)
env.enableCheckpointing(60_000, CheckpointingMode.EXACTLY_ONCE);
// 两次 Checkpoint 之间最小间隔(防止 Checkpoint 挤占正常处理)
cc.setMinPauseBetweenCheckpoints(30_000);
// 超时时间(超过则失败)
cc.setCheckpointTimeout(10 * 60_000);
// 最大并发 Checkpoint 数(通常设 1)
cc.setMaxConcurrentCheckpoints(1);
// ★ 容忍连续失败次数(不设的话一次失败作业就挂)
cc.setTolerableCheckpointFailureNumber(3);
// ★ 作业取消后保留 Checkpoint(生产必配,否则无法从 Checkpoint 恢复)
cc.setExternalizedCheckpointCleanup(
ExternalizedCheckpointCleanup.RETAIN_ON_CANCELLATION);
// 非对齐 Checkpoint
cc.enableUnalignedCheckpoints(true);
cc.setAlignedCheckpointTimeout(Duration.ofSeconds(30));
// 存储位置
cc.setCheckpointStorage("hdfs:///flink/checkpoints");Checkpoint vs Savepoint
| 维度 | Checkpoint | Savepoint |
|---|---|---|
| 触发方 | Flink 自动定期 | 人工触发 |
| 目的 | 故障自动恢复 | 版本升级、集群迁移、A/B 测试 |
| 格式 | 引擎内部格式,可能不跨版本兼容 | 标准格式,保证跨版本兼容 |
| 存储 | 增量(RocksDB) | 全量 |
| 生命周期 | 自动清理旧的 | 手动管理 |
# 触发 Savepoint
flink savepoint <jobId> hdfs:///flink/savepoints
# 优雅停止并生成 Savepoint(推荐的停机方式)
flink stop --savepointPath hdfs:///flink/savepoints <jobId>
# 从 Savepoint 恢复(改了代码也能恢复,只要算子 UID 没变)
flink run -s hdfs:///flink/savepoints/savepoint-xxx -c com.Main job.jar⚠️ 生产铁律:所有有状态算子必须显式设置 UID
javastream.keyBy(...).process(new MyFunction()).uid("my-process-v1").name("订单风控");不设 UID,Flink 会根据算子拓扑自动生成 hash。你改一行代码、加一个算子,hash 就变了 → Savepoint 恢复失败,状态全丢。 这是新手最容易踩的致命坑。
第 10 章:Exactly-Once 的真相
10.1 先破除幻想:Exactly-Once 不是"只处理一次"
很多人的误解:
"Exactly-Once 意味着每条数据只被处理一次"
真相:
Exactly-Once 指的是 Exactly-Once State Consistency(状态一致性), 即:从结果上看,就像每条数据只被处理了一次。
实际上数据可能被处理多次(失败重放),但通过状态回滚保证最终效果等价。
10.2 三种语义的严格定义
At-Most-Once(至多一次):
数据可能丢,不会重
实现:不做任何容错,直接处理
场景:监控指标采样,丢一点无所谓
At-Least-Once(至少一次):
数据不会丢,可能重
实现:失败后从上次 Checkpoint 重放
场景:★ 下游能幂等去重时的最佳选择(性能好)★
Exactly-Once(精确一次):
不丢不重(效果上)
实现:Checkpoint + 事务性 Sink(两阶段提交)
代价:延迟增加(数据必须等事务提交才可见)10.3 端到端 Exactly-Once 的三个必要条件
条件 1:Source 可重放(Replayable Source)
✅ Kafka(可按 offset 重读)
✅ Pulsar、Kinesis
❌ Socket、随机数生成器(无法重放)
条件 2:引擎内部状态一致(Checkpoint 机制)
✅ Flink Checkpoint + Barrier 对齐
条件 3:Sink 支持事务或幂等(Transactional / Idempotent Sink)
方案 A:事务性 Sink(两阶段提交)→ Kafka 事务、JDBC XA
方案 B:幂等 Sink(写入天然去重)→ 按主键 upsert 的 MySQL/HBase/ES三者缺一不可。少一个,端到端 Exactly-Once 就不成立。
10.4 两阶段提交(2PC)在 Flink 中的实现
Flink 用 TwoPhaseCommitSinkFunction 实现事务性写出:
┌──────────────────────────────────────────────────────────┐
│ Phase 1: Pre-commit(预提交) │
│ │
│ Checkpoint Barrier 到达 Sink │
│ → Sink 把当前事务的数据写入外部系统,但【不提交】 │
│ (Kafka: 写入但事务未 commit,消费者读不到) │
│ → 记录事务 ID 到 Checkpoint 状态 │
│ → 开启一个新事务,接收后续数据 │
│ → 向 JobManager 汇报"我预提交成功了" │
└──────────────────────────────────────────────────────────┘
↓
JobManager 收到【所有】算子的成功确认
↓
┌──────────────────────────────────────────────────────────┐
│ Phase 2: Commit(正式提交) │
│ │
│ JobManager 广播 notifyCheckpointComplete │
│ → Sink 提交事务(Kafka commitTransaction) │
│ → 数据对下游消费者可见 │
└──────────────────────────────────────────────────────────┘故障恢复的四种情况:
| 故障时机 | 恢复行为 | 结果 |
|---|---|---|
| Pre-commit 之前挂 | 从上个 Checkpoint 恢复,重新处理 | 无影响 |
| Pre-commit 过程中挂 | 事务未提交,外部系统自动回滚/超时清理 | 无脏数据 |
| Pre-commit 完成、Commit 前挂 | 恢复后根据 Checkpoint 中的事务 ID 重新提交 | 数据正确 |
| Commit 之后挂 | 事务已提交,从下个 Checkpoint 继续 | 数据正确 |
关键洞察:第三种情况是 2PC 的精髓——Checkpoint 里保存了事务 ID,恢复时能把"悬挂的事务"提交掉。这要求外部系统的事务在一段时间内不能被清理。
// Kafka 事务超时必须 > Checkpoint 间隔 + 恢复时间
// 否则恢复时事务已被 Kafka 清理 → 数据丢失
properties.setProperty("transaction.timeout.ms", "900000"); // 15 分钟
KafkaSink<String> sink = KafkaSink.<String>builder()
.setBootstrapServers("kafka:9092")
.setDeliveryGuarantee(DeliveryGuarantee.EXACTLY_ONCE)
.setTransactionalIdPrefix("my-app-") // ★ 必须唯一,多个作业不能重复
.setKafkaProducerConfig(properties)
.setRecordSerializer(...)
.build();⚠️ 著名的生产事故模式:
transaction.timeout.ms默认 1 小时,但 Kafka Broker 端transaction.max.timeout.ms默认 15 分钟。设置超过 Broker 上限会直接报错。而设置太小,作业恢复慢时事务已超时 → 数据丢失。
10.5 幂等写入:更简单的 Exactly-Once
很多场景下,与其用复杂的 2PC,不如让 Sink 幂等:
-- MySQL 幂等写入
INSERT INTO result (dt, category, pv)
VALUES (?, ?, ?)
ON DUPLICATE KEY UPDATE pv = VALUES(pv);
-- ClickHouse ReplacingMergeTree(按主键去重)
CREATE TABLE result (
dt Date, category String, pv UInt64, version UInt64
) ENGINE = ReplacingMergeTree(version)
ORDER BY (dt, category);
-- Doris/StarRocks Unique Key 模型(天然 upsert)
CREATE TABLE result (
dt DATE, category VARCHAR(64), pv BIGINT
) UNIQUE KEY(dt, category)
DISTRIBUTED BY HASH(dt);决策树:
需要 Exactly-Once?
├─ 否 → At-Least-Once(性能最好)
└─ 是 → 下游能按主键幂等吗?
├─ 能 → ★ 用幂等写入(简单、低延迟)★
└─ 不能 → 用 2PC 事务 Sink(复杂、有延迟)💭 思考题 10.1:一个 Flink 作业写 Kafka 用 Exactly-Once,下游消费者的
isolation.level设置为默认的read_uncommitted会发生什么?
第 11 章:Kafka 深度剖析 —— 数据总线的工程美学
11.1 Kafka 为什么能用磁盘做到百万 TPS?
这是一个反直觉的事实。Kafka 的性能秘诀有四个:
秘诀一:顺序写(Sequential I/O)
机械磁盘随机写: ~100 IOPS → 约 0.5 MB/s
机械磁盘顺序写: ~600 MB/s → 快 1000 倍!
Kafka 的 Partition 就是一个【只追加的日志文件】
→ 永远顺序写,永不随机写秘诀二:页缓存(Page Cache)
传统做法:应用自己维护缓存(JVM 堆内)
问题:① GC 压力大 ② 进程重启缓存全丢 ③ 数据在堆内和 PageCache 存两份
Kafka 做法:不维护应用层缓存,全靠操作系统 Page Cache
优势:① 无 GC 压力 ② 进程重启缓存还在 ③ 内存零拷贝这就是为什么 Kafka Broker 的 JVM 堆内存推荐设置得很小(6GB 左右), 把绝大部分物理内存留给 Page Cache。给 Kafka 配 64GB 堆是典型的错误配置。
秘诀三:零拷贝(Zero-Copy / sendfile)
传统读取发送流程(4 次拷贝 + 4 次上下文切换):
磁盘 → 内核 Page Cache → 用户态应用 Buffer → 内核 Socket Buffer → 网卡
① ② ③ ④
零拷贝(sendfile 系统调用,2 次拷贝 + 2 次切换):
磁盘 → 内核 Page Cache ─────────────────────────────▶ 网卡
① (DMA 直接传输,数据不进用户态)②性能提升:CPU 消耗降低约 50%,吞吐提升数倍。
⚠️ 注意:开启 SSL/TLS 加密后零拷贝会失效(因为需要在用户态做加密),吞吐会明显下降。这是安全和性能的真实权衡。
秘诀四:批量 + 压缩
// Producer 端攒批
props.put("batch.size", 65536); // 16KB → 64KB
props.put("linger.ms", 10); // 攒 10ms 再发(关键的延迟/吞吐旋钮)
props.put("compression.type", "zstd"); // lz4/snappy/zstd
// ★ 精妙之处:压缩在 Producer 端做,Broker 端【不解压】直接存
// Consumer 端解压 → Broker 的 CPU 完全不消耗在压缩上11.2 存储结构:Segment 与索引
topic-order-0/ ← Partition 目录
├── 00000000000000000000.log ← Segment 数据文件(默认 1GB 滚动)
├── 00000000000000000000.index ← 偏移量索引(稀疏索引)
├── 00000000000000000000.timeindex ← 时间戳索引
├── 00000000000001234567.log ← 下一个 Segment,文件名 = 起始 offset
├── 00000000000001234567.index
└── leader-epoch-checkpoint查找 offset=1234600 的消息:
① 二分查找文件名 → 定位到 00000000000001234567.log
② 在对应 .index 中二分查找 → 找到最接近的索引项(如 offset=1234580 → position=4096)
③ 从 position=4096 开始顺序扫描 → 找到 1234600
★ 稀疏索引:默认每 4KB 数据建一个索引项
→ 索引文件小,可全部放内存
→ 用少量顺序扫描换极小的索引空间📐 这个设计和上篇讲的 Parquet Row Group 统计信息是同一个思想: 用粗粒度索引快速定位到一个小范围,再顺序扫描。 在大数据领域反复出现。
11.3 副本机制:ISR 与高水位
Partition 的副本集合:
AR (Assigned Replicas) = 所有副本
ISR (In-Sync Replicas) = 与 Leader 保持同步的副本 ★核心概念★
OSR (Out-of-Sync) = 落后太多的副本
判定标准:replica.lag.time.max.ms(默认 30s)
→ 超过 30 秒没追上 Leader,就被踢出 ISR两个水位线:
LEO (Log End Offset) = 每个副本最后一条消息的 offset + 1
HW (High Watermark) = ISR 中【最小】的 LEO ★消费者只能读到 HW 之前的数据★
Leader: [0][1][2][3][4][5] LEO=6
Follower1: [0][1][2][3][4] LEO=5
Follower2: [0][1][2][3] LEO=4 ← 最慢
HW = min(6, 5, 4) = 4
→ 消费者只能读到 offset 0~3
为什么?因为 offset 4、5 还没被所有 ISR 副本同步,
如果此时 Leader 挂了,新 Leader 可能没有这些数据 → 会出现"读到的数据消失"11.4 生产配置:可靠性与性能的权衡矩阵
// ===== 场景 A:金融级不丢数据(吞吐降低约 30%)=====
props.put("acks", "all"); // 等所有 ISR 确认
props.put("retries", Integer.MAX_VALUE);
props.put("enable.idempotence", true); // 幂等 Producer,防重复
props.put("max.in.flight.requests.per.connection", 5);
// Broker/Topic 端配套:
// min.insync.replicas = 2 (ISR 少于 2 个时拒绝写入)
// unclean.leader.election.enable = false ★ 禁止非 ISR 副本当 Leader ★
// ===== 场景 B:日志采集,追求吞吐(可容忍少量丢失)=====
props.put("acks", "1");
props.put("linger.ms", 100);
props.put("batch.size", 131072);
props.put("compression.type", "lz4");
// ===== 场景 C:低延迟(如实时风控)=====
props.put("acks", "1");
props.put("linger.ms", 0); // 不攒批,立即发送
props.put("compression.type", "none");⚠️
unclean.leader.election.enable=true是数据丢失的头号元凶。 它允许落后的副本成为 Leader(保住可用性),代价是丢失 Leader 独有的数据。 这是 CAP 中 A 和 C 的赤裸裸选择——金融场景必须设为 false。
11.5 前沿:KRaft 模式取代 ZooKeeper
Kafka 2.8 之前:依赖 ZooKeeper 存元数据
问题:① 多一套系统运维 ② 元数据规模受限(分区数上限约 20 万)
③ Controller 故障切换慢(需要从 ZK 全量拉取元数据)
Kafka 3.3+ KRaft 模式(Kafka Raft):
✓ 用 Kafka 自己的 Raft 实现管理元数据
✓ 元数据本身就是一个 Kafka topic(__cluster_metadata)
✓ 分区数上限提升到 ★ 百万级 ★
✓ Controller 故障切换从分钟级降到秒级
✓ 部署简化,不用再装 ZK
Kafka 4.0(2025):★ 彻底移除 ZooKeeper 支持 ★📌 给学习者:2026 年新搭 Kafka 集群请直接用 KRaft 模式,不要再学 ZK 部署方式了。
第 12 章:湖仓一体 —— Iceberg / Hudi / Delta 三国志
12.1 为什么需要"表格式"这一层?
Hive 表的原罪
Hive 表的本质 = 一个目录 + Metastore 里的分区列表
/warehouse/orders/dt=2026-09-01/part-00000.parquet
/part-00001.parquet
致命缺陷:
① 无 ACID —— 写到一半任务挂了,读者会看到部分文件(脏读)
② 无法行级更新/删除 —— GDPR 要删某个用户数据?只能重写整个分区
③ 分区必须显式指定 —— WHERE ts > '...' 无法裁剪 dt 分区
④ Schema 演进危险 —— 改列名/类型可能导致历史数据读不出来
⑤ 元数据瓶颈 —— 百万分区时 Metastore 查询慢如蜗牛
⑥ 无快照/回滚 —— 误删数据无法恢复表格式层的解法
核心思路:★ 不再用"目录"定义表,而是用【元数据文件】显式记录表包含哪些数据文件 ★
传统:表 = "这个目录下的所有文件"(隐式,依赖文件系统 list)
现代:表 = "元数据文件里列出的这批文件"(显式,原子切换指针)
→ 提交 = 原子地切换指向新元数据文件的指针
→ 天然获得:ACID、快照、时间旅行、无需 list 的高效规划12.2 Iceberg 元数据结构逐层拆解
catalog (Hive Metastore / Glue / Nessie / JDBC)
│ 保存一个指针:current_metadata_location
▼
┌──────────────────────────────────────────────────┐
│ v3.metadata.json ← 表的当前状态 │
│ ├ schema(含 schema 演进历史) │
│ ├ partition-spec(含分区演进历史) │
│ ├ current-snapshot-id: 8712 │
│ └ snapshots: [ │
│ {id:8710, manifest-list: snap-8710.avro}, │
│ {id:8711, manifest-list: snap-8711.avro}, │
│ {id:8712, manifest-list: snap-8712.avro} │ ← 当前
│ ] │
└──────────────┬───────────────────────────────────┘
▼
┌──────────────────────────────────────────────────┐
│ snap-8712.avro (Manifest List) │
│ 列出本快照包含哪些 manifest 文件,每条记录含: │
│ - manifest_path │
│ - partition 范围统计(用于分区级剪枝) │
│ - added/existing/deleted 文件数 │
└──────────────┬───────────────────────────────────┘
▼
┌──────────────────────────────────────────────────┐
│ manifest-a.avro (Manifest File) │
│ 列出具体的数据文件,每条记录含: │
│ - file_path: s3://.../data-001.parquet │
│ - partition: {dt: 2026-09-01} │
│ - record_count: 1000000 │
│ - ★ 每列的 min/max/null_count ★ ← 文件级剪枝 │
└──────────────┬───────────────────────────────────┘
▼
data-001.parquet data-002.parquet ...一次查询的执行过程
SELECT * FROM orders WHERE dt = '2026-09-01' AND amount > 5000;① 读 catalog → 拿到 v3.metadata.json 路径
② 读 metadata.json → 拿到 current snapshot 的 manifest-list
③ 读 manifest-list → 用 partition 统计【剪掉】不含 dt=2026-09-01 的 manifest
④ 读剩余 manifest → 用 amount 的 min/max【剪掉】 max(amount)<=5000 的文件
⑤ 只读剩下的那几个 Parquet 文件
★ 全程【零次 list 操作】★
→ 对象存储上性能提升巨大(list 慢且收费)
→ 百万文件的表,规划耗时依然是秒级💡 这就是 Iceberg 相比 Hive 最大的性能优势来源。Hive 需要 list 目录找文件,S3 上 list 一个大分区可能要几十秒。
12.3 三者的核心差异(选型决策)
| 维度 | Iceberg | Hudi | Delta Lake |
|---|---|---|---|
| 出身 | Netflix → Apache | Uber → Apache | Databricks |
| 设计重心 | 通用表格式、多引擎中立 | 流式 upsert、增量处理 | Spark 生态深度集成 |
| 元数据 | 三层(metadata→manifest list→manifest) | Timeline + 文件组 | _delta_log JSON + checkpoint parquet |
| 更新模式 | CoW + MoR(v2) | CoW + MoR(最成熟) | CoW(Deletion Vector 后接近 MoR) |
| 主键/索引 | 无内置主键 | 有主键 + 多种索引(Bloom/HBase/Bucket) | 无主键(靠 MERGE) |
| 隐藏分区 | ✅ 独有杀手锏 | ❌ | ❌ |
| 分区演进 | ✅ 独有 | ❌ | ❌(需重写) |
| 引擎兼容 | ✅ 最好(Spark/Flink/Trino/Doris/StarRocks/Snowflake/BigQuery) | 中(Spark/Flink 好) | Spark 最好,其他靠 UniForm |
| 2026 势头 | ✅ 事实标准(AWS S3 Tables / Snowflake / Databricks 全面拥抱) | 实时场景仍强 | Databricks 生态内强 |
Iceberg 的杀手锏 1:隐藏分区(Hidden Partitioning)
-- ❌ Hive 的痛点:必须显式写分区列,否则全表扫描
CREATE TABLE orders (ts TIMESTAMP, amount DOUBLE)
PARTITIONED BY (dt STRING); -- 需要额外维护一个 dt 列
-- 查询必须这样写才能裁剪:
SELECT * FROM orders WHERE dt = '2026-09-01'; -- ✅ 裁剪
SELECT * FROM orders WHERE ts >= '2026-09-01'; -- ❌ 全表扫描!
-- ✅ Iceberg:分区是"从某列计算出来的函数"
CREATE TABLE orders (ts TIMESTAMP, amount DOUBLE)
PARTITIONED BY (days(ts)); -- 不需要额外的 dt 列
-- 查询直接用原始列,Iceberg 自动推导分区裁剪
SELECT * FROM orders WHERE ts >= '2026-09-01'; -- ✅ 自动裁剪!支持的分区变换函数:years() months() days() hours() bucket(N, col) truncate(N, col)
🎯 这一个特性就避免了无数生产事故——再也不会有人写错分区条件导致扫全表了。
Iceberg 的杀手锏 2:分区演进(Partition Evolution)
-- 业务初期数据少,按月分区
ALTER TABLE orders SET PARTITION SPEC (months(ts));
-- 一年后数据暴涨,改成按天分区
ALTER TABLE orders SET PARTITION SPEC (days(ts));
-- ★ 关键:历史数据【不需要重写】★
-- 旧数据继续用月分区规则,新数据用天分区规则
-- 查询时 Iceberg 会分别应用各自的规则Hive 做同样的事需要:重写全表 → 可能是几天的作业 + 几十 TB 的 IO。
CoW vs MoR:更新策略的根本权衡
Copy-on-Write(写时复制):
更新 1 行 → 读取整个文件 → 修改 → 写出新文件 → 旧文件标记删除
写放大:★★★★★(改 1 行可能重写 128MB)
读性能:★★★★★(直接读,无合并开销)
适用:★ 读多写少、批量更新(如 T+1 数仓)★
Merge-on-Read(读时合并):
更新 1 行 → 只写一条"删除记录"或"更新记录"到 Delta 文件
读取时 → 基础文件 + Delta 文件实时合并
写放大:★(极低)
读性能:★★(需要合并,Delta 越多越慢)
适用:★ 高频更新、实时入湖(如 CDC 同步)★
必须:定期 Compaction 合并 Delta 文件决策口诀:
CDC 实时入湖 → MoR + 定期 Compaction
T+1 批量更新 → CoW
读多写极少 → CoW12.4 Iceberg 实操(在你的 EMR 集群上)
环境准备
# 下载 Iceberg runtime jar(匹配 Spark 3.5 + Scala 2.12)
cd $SPARK_HOME/jars
wget https://repo1.maven.org/maven2/org/apache/iceberg/\
iceberg-spark-runtime-3.5_2.12/1.6.1/iceberg-spark-runtime-3.5_2.12-1.6.1.jar# 启动配置了 Iceberg 的 spark-sql
spark-sql \
--conf spark.sql.extensions=org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions \
--conf spark.sql.catalog.ice=org.apache.iceberg.spark.SparkCatalog \
--conf spark.sql.catalog.ice.type=hive \
--conf spark.sql.catalog.ice.uri=thrift://172.16.146.220:9083 \
--conf spark.sql.catalog.ice.warehouse=hdfs:///lab/iceberg完整 CRUD 实操
-- 建库建表
CREATE DATABASE IF NOT EXISTS ice.lakehouse;
CREATE TABLE ice.lakehouse.orders (
order_id BIGINT,
user_id BIGINT,
amount DECIMAL(10,2),
status STRING,
create_time TIMESTAMP
)
USING iceberg
PARTITIONED BY (days(create_time)) -- ★ 隐藏分区
TBLPROPERTIES (
'write.format.default' = 'parquet',
'write.parquet.compression-codec'= 'zstd',
'write.target-file-size-bytes' = '134217728', -- 128MB
'format-version' = '2' -- ★ v2 才支持行级删除
);
-- 插入
INSERT INTO ice.lakehouse.orders VALUES
(1, 1001, 299.00, 'paid', TIMESTAMP '2026-09-01 10:00:00'),
(2, 1002, 1599.00,'pending', TIMESTAMP '2026-09-01 11:30:00'),
(3, 1001, 89.50, 'paid', TIMESTAMP '2026-09-02 09:15:00');
-- ★ 行级 UPDATE(Hive 做不到)
UPDATE ice.lakehouse.orders SET status = 'paid' WHERE order_id = 2;
-- ★ 行级 DELETE(GDPR 合规必备)
DELETE FROM ice.lakehouse.orders WHERE user_id = 1002;
-- ★ MERGE INTO(数仓 SCD / 幂等入库的核心语法)
MERGE INTO ice.lakehouse.orders t
USING (SELECT 1 AS order_id, 1001 AS user_id, 399.00 AS amount,
'refund' AS status, TIMESTAMP '2026-09-03 12:00:00' AS create_time) s
ON t.order_id = s.order_id
WHEN MATCHED THEN UPDATE SET t.status = s.status, t.amount = s.amount
WHEN NOT MATCHED THEN INSERT *;时间旅行与回滚(运维救命功能)
-- 查看所有快照
SELECT snapshot_id, committed_at, operation, summary
FROM ice.lakehouse.orders.snapshots;
-- 查看某个历史时刻的数据
SELECT * FROM ice.lakehouse.orders
TIMESTAMP AS OF '2026-09-01 12:00:00';
-- 查看某个快照版本
SELECT * FROM ice.lakehouse.orders VERSION AS OF 8712345678901234;
-- ★ 误操作回滚(生产救命)★
CALL ice.system.rollback_to_snapshot('lakehouse.orders', 8712345678901234);
-- 增量读取(流批一体的基础)
SELECT * FROM ice.lakehouse.orders
VERSION AS OF 8712345678901234;
-- Spark API 方式的增量读:
-- .option("start-snapshot-id", "...")
-- .option("end-snapshot-id", "...")🔬 硬核实验 12.1:模拟一次误删并恢复
sql-- 1. 记录当前快照 ID SELECT snapshot_id FROM ice.lakehouse.orders.snapshots ORDER BY committed_at DESC LIMIT 1; -- 2. 模拟事故 DELETE FROM ice.lakehouse.orders WHERE 1=1; SELECT COUNT(*) FROM ice.lakehouse.orders; -- 0 行,凉了 -- 3. 回滚 CALL ice.system.rollback_to_snapshot('lakehouse.orders', <刚才记录的ID>); SELECT COUNT(*) FROM ice.lakehouse.orders; -- 数据回来了!这个实验建议让学生做一遍——它会让人对"数据湖也能有事务"产生直观震撼。
表维护(生产必备的定时任务)
-- ① 小文件合并(最重要的维护操作)
CALL ice.system.rewrite_data_files(
table => 'lakehouse.orders',
strategy=> 'binpack',
options => map('target-file-size-bytes','134217728',
'min-input-files','5')
);
-- ② 按列排序重写(提升查询剪枝率)
CALL ice.system.rewrite_data_files(
table => 'lakehouse.orders',
strategy => 'sort',
sort_order => 'user_id ASC NULLS LAST'
);
-- ③ 清理过期快照(释放存储,默认保留 5 天)
CALL ice.system.expire_snapshots(
table => 'lakehouse.orders',
older_than => TIMESTAMP '2026-09-01 00:00:00',
retain_last => 10
);
-- ④ 清理孤儿文件(失败任务残留)
CALL ice.system.remove_orphan_files(
table => 'lakehouse.orders',
older_than => TIMESTAMP '2026-09-01 00:00:00'
);
-- ⑤ 重写 manifest(元数据太多时)
CALL ice.system.rewrite_manifests('lakehouse.orders');⚠️ 不做维护的 Iceberg 表会慢慢变慢。流式写入的表可能每分钟产生一批小文件, 一周后有几十万个小文件 + 几万个快照 → 查询规划都要几分钟。 把这四个 CALL 写成每日调度任务,是数据平台的基本功。
第 13 章:CDC 全链路 —— 让数据库变成流
13.1 CDC 方案对比
| 方案 | 原理 | 延迟 | 对源库影响 | 能否捕获 DELETE | 评价 |
|---|---|---|---|---|---|
| 定时全量拉取 | SELECT * | 小时级 | 大 | ❌ | 数据量大时不可行 |
| 基于时间戳增量 | WHERE update_time > ? | 分钟级 | 中 | ❌ | 无法捕获删除,且依赖业务字段规范 |
| 触发器 | DB Trigger 写影子表 | 秒级 | 大 | ✅ | 侵入性强,影响源库性能 |
| 日志解析(Binlog) | 解析 MySQL binlog | 秒级/亚秒级 | 极小 | ✅ | ★ 唯一正确答案 ★ |
为什么 binlog 是最优解?
Binlog 本身就是数据库为了主从复制而产生的【完整变更流】
→ CDC 工具伪装成一个 MySQL Slave
→ 对主库来说就是多了一个从库,几乎零额外开销
→ 且天然包含 INSERT/UPDATE/DELETE 的完整前后镜像13.2 MySQL 准备工作
-- 检查 binlog 配置
SHOW VARIABLES LIKE 'log_bin'; -- 必须 ON
SHOW VARIABLES LIKE 'binlog_format'; -- 必须 ROW(不能是 STATEMENT/MIXED)
SHOW VARIABLES LIKE 'binlog_row_image'; -- 建议 FULL
-- my.cnf 配置
-- [mysqld]
-- server-id = 1
-- log_bin = mysql-bin
-- binlog_format = ROW
-- binlog_row_image = FULL
-- expire_logs_days = 7
-- 创建 CDC 专用账号(最小权限原则)
CREATE USER 'cdc_user'@'%' IDENTIFIED BY 'StrongPass123!';
GRANT SELECT, RELOAD, SHOW DATABASES, REPLICATION SLAVE, REPLICATION CLIENT
ON *.* TO 'cdc_user'@'%';
FLUSH PRIVILEGES;⚠️ 为什么必须是 ROW 格式?
STATEMENT格式记录的是 SQL 语句(如UPDATE t SET a=a+1),CDC 无法知道具体哪些行变了、变成了什么值。ROW格式记录每一行的前后镜像,才能还原出完整的变更事件。
13.3 Flink CDC 实战:MySQL → Iceberg
方式一:纯 SQL(推荐,开发效率高)
-- 在 Flink SQL Client 中执行
SET 'execution.checkpointing.interval' = '30s';
SET 'execution.checkpointing.mode' = 'EXACTLY_ONCE';
-- ① 定义 CDC 源表
CREATE TABLE mysql_orders (
order_id BIGINT,
user_id BIGINT,
amount DECIMAL(10,2),
status STRING,
update_time TIMESTAMP(3),
PRIMARY KEY (order_id) NOT ENFORCED -- ★ CDC 必须有主键
) WITH (
'connector' = 'mysql-cdc',
'hostname' = '192.168.1.50',
'port' = '3306',
'username' = 'cdc_user',
'password' = 'StrongPass123!',
'database-name' = 'shop',
'table-name' = 'orders',
'server-time-zone' = 'Asia/Shanghai',
-- ★ 增量快照算法(2.0+),支持并行读全量 + 无锁
'scan.incremental.snapshot.enabled' = 'true',
'scan.incremental.snapshot.chunk.size' = '8096',
'scan.startup.mode' = 'initial' -- initial:全量+增量 / latest-offset:仅增量
);
-- ② 定义 Iceberg 目标表
CREATE CATALOG ice WITH (
'type' = 'iceberg',
'catalog-type' = 'hive',
'uri' = 'thrift://172.16.146.220:9083',
'warehouse' = 'hdfs:///lab/iceberg'
);
CREATE TABLE IF NOT EXISTS ice.lakehouse.ods_orders (
order_id BIGINT,
user_id BIGINT,
amount DECIMAL(10,2),
status STRING,
update_time TIMESTAMP(3),
PRIMARY KEY (order_id) NOT ENFORCED
) WITH (
'format-version' = '2',
'write.upsert.enabled' = 'true', -- ★ 开启 upsert
'write.distribution-mode' = 'hash'
);
-- ③ 一行 SQL 启动实时同步
INSERT INTO ice.lakehouse.ods_orders
SELECT order_id, user_id, amount, status, update_time FROM mysql_orders;就这样,MySQL 的每一次增删改都会在 30 秒内同步到数据湖。
Flink CDC 2.0 增量快照算法(核心创新)
CDC 1.x 的痛点:
读全量时需要【全局锁】(FLUSH TABLES WITH READ LOCK)
→ 锁住整个库,业务受影响
→ 且全量阶段只能单并行度,大表要读几小时
→ 全量期间失败 → 从头再来
CDC 2.0 的解法(借鉴 DBLog 论文):
① 把表按主键切成多个 chunk(如每 8096 行一个)
② 每个 chunk 独立读取,可【多并行度】
③ 读 chunk 时记录 binlog 位点,读完后回补这期间的 binlog 变更
④ 用【水位线 + 去重】保证一致性,全程【无锁】
⑤ 每个 chunk 完成即 Checkpoint → 失败只重跑该 chunk这是 Flink CDC 最重要的工程突破,也是它能在生产大规模使用的关键。
13.4 CDC 的四种数据变更类型
Flink 内部用 Changelog 表示变更流:
+I (INSERT) 新增一行
-U (UPDATE_BEFORE) 更新前的旧值(撤回)
+U (UPDATE_AFTER) 更新后的新值
-D (DELETE) 删除一行
例:MySQL 执行 UPDATE orders SET amount=399 WHERE order_id=1
Flink 收到两条记录:
-U (1, 1001, 299.00, 'paid')
+U (1, 1001, 399.00, 'paid')
★ 为什么需要 -U?★
因为下游可能有聚合:SUM(amount)
收到 -U 时:sum -= 299
收到 +U 时:sum += 399
→ 结果正确
如果只发 +U,sum 会变成 299+399=698,错误!💡 这就是"回撤流(Retract Stream)"机制,是 Flink SQL 能正确处理流式聚合的根本保障。
13.5 整库同步(生产实用技巧)
一张表一个作业,100 张表就要 100 个 Flink 作业 + 100 个 binlog 连接 → 源库扛不住。
-- ★ 用正则一次同步整库(一个作业、一个 binlog 连接)
CREATE TABLE all_tables (
...
) WITH (
'connector' = 'mysql-cdc',
'database-name' = 'shop',
'table-name' = '.*', -- 正则匹配所有表
...
);或使用 Flink CDC 3.x 的 YAML Pipeline(2024 年新增,极大简化整库同步):
source:
type: mysql
hostname: 192.168.1.50
port: 3306
username: cdc_user
password: StrongPass123!
tables: shop.\.* # 整库
server-id: 5400-5404
sink:
type: iceberg
catalog.type: hive
uri: thrift://172.16.146.220:9083
warehouse: hdfs:///lab/iceberg
transform:
- source-table: shop.orders
projection: \*, 'CN' AS region # 加字段
route: # 分库分表合并
- source-table: shop.orders_\.*
sink-table: lakehouse.ods_orders
pipeline:
name: mysql-to-iceberg-sync
parallelism: 4bin/flink-cdc.sh mysql-to-iceberg.yamlFlink CDC 3.x 的核心能力:整库同步、自动建表、Schema 变更自动同步(源库 ALTER TABLE 会自动传导到湖表)、分库分表路由合并。
🔥 这是 2024-2026 数据集成领域最大的效率提升。以前接入 500 张表要写 500 个任务,现在一个 YAML 搞定。
第 14 章:数据建模 —— 被低估的核心竞争力
14.1 为什么建模比调优更重要?
技术调优的收益:让一个 SQL 从 10 分钟跑到 3 分钟(3 倍)
好的建模的收益:让这个 SQL 根本不需要跑(∞ 倍)
★ 一个建模糟糕的数仓,再牛的调优也救不回来 ★
★ 一个建模优秀的数仓,普通配置也能跑得飞快 ★更现实的问题:技术组件三年一换,但数据模型是企业的长期资产。建模能力是数据工程师最保值的技能。
14.2 维度建模:Kimball 方法论
四步建模法
第 1 步:选择业务过程(Business Process)
↓ 「下单」「支付」「退款」「物流揽收」
↓ 原则:一个业务过程 = 一张事实表
第 2 步:声明粒度(Grain)★ 最关键的一步 ★
↓ 「一行 = 一个订单」 还是 「一行 = 订单中的一个商品」?
↓ 铁律:粒度越细越好,能聚合上去,不能拆下来
第 3 步:确定维度(Dimensions)
↓ 用户、商品、时间、地区、渠道、活动...
↓ 问自己:业务方会怎么切分数据看?
第 4 步:确定事实(Facts)
↓ 金额、数量、折扣、运费...
↓ 原则:必须是【可加】的数值⚠️ 粒度错误是最昂贵的建模错误。如果你把事实表建成"一行=一个订单",后来业务要看"商品维度的销量",你得推倒重来。
星型 vs 雪花型
★ 星型模型(Star Schema)—— 大数据首选
dim_user dim_time
\ /
\ /
fact_order ← 维度表直接连事实表,不再拆分
/ \
/ \
dim_item dim_region
优点:JOIN 少(最多 1 跳)、查询快、业务方易懂
缺点:维度表有冗余(如 dim_item 里冗余了品类名称)
❄ 雪花模型(Snowflake Schema)—— 传统数据库首选
dim_item ──▶ dim_category ──▶ dim_category_l1
│
fact_order
优点:无冗余、存储省
缺点:JOIN 层级多、查询慢
★ 大数据时代的选择:坚定选星型 ★
理由:存储很便宜,但 Shuffle Join 很贵渐变维(SCD, Slowly Changing Dimension)—— 高频面试题
问题:用户的城市从"北京"变成了"上海",历史订单应该算北京还是上海?| 类型 | 做法 | 结果 | 适用 |
|---|---|---|---|
| SCD Type 1 | 直接覆盖 | 历史全变成"上海",丢失历史 | 修正错误数据 |
| SCD Type 2 | 新增一行,用生效时间区间标记 | 历史订单关联到"北京"那一行 | ★ 需要历史准确性时的标准做法 ★ |
| SCD Type 3 | 加一列 previous_city | 只能保留上一个值 | 只关心"变更前后"对比 |
SCD Type 2 的表结构:
CREATE TABLE dim_user (
sk_user_id BIGINT, -- ★ 代理键(surrogate key),每个版本唯一
user_id BIGINT, -- 业务键(自然键),同一用户多行相同
user_name STRING,
city STRING,
start_date DATE, -- 此版本生效开始
end_date DATE, -- 此版本失效(当前版本填 '9999-12-31')
is_current BOOLEAN
) USING iceberg;
-- 数据示例:
-- sk=1001, user_id=1, city='北京', start='2024-01-01', end='2026-06-30', is_current=false
-- sk=1002, user_id=1, city='上海', start='2026-07-01', end='9999-12-31', is_current=true
-- ★ 事实表存 sk_user_id(而非 user_id),自动关联到正确的历史版本用 Iceberg MERGE 实现 SCD2:
-- 步骤 1:关闭发生变化的旧版本
MERGE INTO dim_user t
USING stg_user s
ON t.user_id = s.user_id AND t.is_current = true
WHEN MATCHED AND (t.city <> s.city OR t.user_name <> s.user_name)
THEN UPDATE SET t.end_date = CURRENT_DATE() - INTERVAL 1 DAY,
t.is_current = false;
-- 步骤 2:插入新版本
INSERT INTO dim_user
SELECT
hash(s.user_id, CURRENT_DATE()) AS sk_user_id,
s.user_id, s.user_name, s.city,
CURRENT_DATE() AS start_date,
DATE '9999-12-31' AS end_date,
true AS is_current
FROM stg_user s
LEFT JOIN dim_user t ON s.user_id = t.user_id AND t.is_current = true
WHERE t.user_id IS NULL
OR t.city <> s.city OR t.user_name <> s.user_name;14.3 数仓分层:工程化的秩序
┌────────────────────────────────────────────────────────────┐
│ ADS 应用数据层 (Application Data Service) │
│ 面向具体报表/接口的最终结果,高度聚合 │
│ 例:ads_gmv_daily_report、ads_user_profile_tag │
├────────────────────────────────────────────────────────────┤
│ DWS 数据汇总层 (Data Warehouse Summary) │
│ 按主题的轻度聚合宽表,复用性强 │
│ 例:dws_user_action_1d(用户日汇总) │
│ dws_item_sale_1d(商品日销售汇总) │
├────────────────────────────────────────────────────────────┤
│ DWD 数据明细层 (Data Warehouse Detail) ★ 最重要的一层 ★ │
│ ① 清洗(去脏数据、格式统一) │
│ ② 脱敏(手机号/身份证加密) │
│ ③ 维度退化(把常用维度字段冗余进来,减少 JOIN) │
│ ④ 明细粒度,不聚合 │
│ 例:dwd_order_detail_di │
├────────────────────────────────────────────────────────────┤
│ DIM 维度层 │
│ dim_user、dim_item、dim_region、dim_date │
├────────────────────────────────────────────────────────────┤
│ ODS 原始数据层 (Operational Data Store) │
│ ★ 与源系统【一模一样】,不做任何加工 ★ │
│ 作用:出问题时可回溯、可重跑 │
│ 例:ods_mysql_orders_di │
└────────────────────────────────────────────────────────────┘分层的价值(说服业务方的话术):
| 价值 | 说明 |
|---|---|
| 复用 | DWS 算一次,10 个 ADS 报表都能用,算力节省 90% |
| 解耦 | 源库表结构变了,只改 ODS→DWD,上游 100 个报表不用动 |
| 可溯源 | 数据对不上时能一层层往下查,快速定位问题层 |
| 可重跑 | ODS 保留原始数据,任何一层出错都能重算 |
| 权限治理 | 敏感数据在 DWD 层脱敏,上层无需再管 |
命名规范(团队协作的基础)
{层次}_{业务域}_{描述}_{刷新周期}{类型}
刷新周期:
d = 日 h = 小时 m = 月 rt = 实时
类型:
i = 增量(incremental) f = 全量(full)
示例:
ods_shop_orders_di → ODS 层、电商域、订单、日增量
dwd_trade_order_detail_di → DWD 层、交易域、订单明细、日增量
dws_trade_user_order_1d → DWS 层、交易域、用户订单、1日汇总
ads_trade_gmv_report_df → ADS 层、交易域、GMV 报表、日全量
dim_user_df → 维度表、用户、日全量14.4 前沿:dbt / SQLMesh 与 Analytics Engineering
传统数仓开发的痛点:
几千个 SQL 文件散落各处 → 没人知道依赖关系
改一张表 → 不知道会影响哪些下游
没有测试 → 数据错了下游才发现
没有文档 → 新人接手要看三个月dbt(data build tool)的解法:把软件工程实践引入数据开发。
-- models/dws/dws_user_order_1d.sql
{{ config(
materialized='incremental',
incremental_strategy='merge',
unique_key=['dt','user_id'],
file_format='iceberg',
partition_by=['dt']
) }}
SELECT
dt,
user_id,
COUNT(*) AS order_cnt,
SUM(amount) AS total_amount,
COUNT(DISTINCT item_id) AS item_cnt
FROM {{ ref('dwd_order_detail_di') }} -- ★ ref() 自动构建依赖图
WHERE 1=1
{% if is_incremental() %}
AND dt > (SELECT MAX(dt) FROM {{ this }}) -- 增量逻辑
{% endif %}
GROUP BY dt, user_id# models/dws/schema.yml —— 数据测试与文档
version: 2
models:
- name: dws_user_order_1d
description: "用户日粒度订单汇总,供用户画像和 GMV 报表使用"
columns:
- name: user_id
description: "用户ID"
tests:
- not_null
- relationships:
to: ref('dim_user')
field: user_id
- name: total_amount
tests:
- not_null
- dbt_utils.accepted_range:
min_value: 0dbt run # 按依赖顺序执行所有模型
dbt test # 运行所有数据质量测试
dbt docs generate && dbt docs serve # ★ 自动生成带血缘图的文档网站 ★dbt 带来的四个革命性改变:
- 依赖自动推导 ——
ref()自动构建 DAG,不用手写调度依赖 - 数据测试即代码 —— 唯一性、非空、外键、值域全部自动化检查
- 文档与血缘自动生成 —— 点击任何一张表能看到上下游全链路
- 版本控制 + CI/CD —— SQL 进 Git,PR 评审,自动化部署
🔮 2026 新势力:SQLMesh 相比 dbt 的改进:
- 列级血缘(dbt 只有表级)
- 虚拟数据环境:不用复制数据就能创建开发环境
- 自动判断变更是否破坏性:改了个注释 vs 改了个 JOIN 条件,前者不用重算
- 原生支持多引擎、原生支持时间旅行
第 15 章:实时 OLAP 引擎对决
15.1 为什么需要专门的 OLAP 引擎?
数据湖(Iceberg + Trino):
✓ 存储便宜、容量无限、格式开放
✗ 查询延迟秒级到分钟级(对象存储 IO + 无索引)
✗ 高并发差(几十 QPS 就吃力)
实时 OLAP(ClickHouse/Doris/StarRocks):
✓ 亚秒级响应、高并发(几千 QPS)
✓ 本地 SSD + 精细索引 + 向量化执行
✗ 存储贵、容量有限
★ 典型架构:湖做全量存储和离线加工,OLAP 做最近 N 个月的高频查询服务 ★15.2 三大引擎深度对比
| 维度 | ClickHouse | Apache Doris | StarRocks |
|---|---|---|---|
| 架构 | 无中心节点、Shared-Nothing | FE(元数据) + BE(存储计算) | FE + BE/CN,支持存算分离 |
| 单表查询 | ★★★★★ 最强 | ★★★★ | ★★★★★ |
| 多表 JOIN | ★★ 弱(最大短板) | ★★★★ | ★★★★★ 最强(CBO 优秀) |
| 实时写入 | ★★★(不宜高频小批) | ★★★★★ | ★★★★★ |
| 更新删除 | ★★(ReplacingMergeTree 异步) | ★★★★(Unique Key) | ★★★★★(主键模型+部分列更新) |
| 并发能力 | ★★(几十 QPS) | ★★★★ | ★★★★ |
| 运维复杂度 | 高(分布式表、副本手动管) | 低(自动分片副本) | 低 |
| 数据湖分析 | ★★★ | ★★★★ | ★★★★★(外表性能接近内表) |
| 物化视图 | 有限 | 较好 | ★★★★★ 异步+同步,自动改写查询 |
| 国内生态 | 强 | ★★★★★ 最强(百度开源) | 强(鼎石) |
选型决策树
你的核心场景是什么?
│
├─ 单表海量日志分析、写多读少(如 APM、日志检索)
│ → ★ ClickHouse ★(单表性能无敌)
│
├─ 复杂多表 JOIN 的实时数仓、需要星型模型
│ → ★ StarRocks ★(JOIN 能力和 CBO 最强)
│
├─ 需要湖仓联邦查询(同时查 Iceberg 外表和内表)
│ → ★ StarRocks ★(外表性能最好)
│
├─ 高频 upsert(CDC 实时同步到 OLAP)
│ → ★ StarRocks 主键模型 / Doris Unique Key ★
│
├─ 团队运维能力弱、需要开箱即用、要中文社区支持
│ → ★ Doris ★(运维最省心,国内生态最好)
│
└─ 已有 Hadoop 集群、只做离线即席查询
→ ★ Trino/Presto ★(无需额外存储,直接查湖)15.3 ClickHouse 核心:MergeTree 家族
CREATE TABLE events (
event_date Date,
event_time DateTime,
user_id UInt64,
event_type LowCardinality(String), -- ★ 低基数列专用类型,自动字典编码
url String,
duration UInt32
) ENGINE = MergeTree()
PARTITION BY toYYYYMM(event_date) -- 分区(不宜太多,建议按月)
ORDER BY (event_date, user_id, event_time) -- ★ 排序键 = 主索引,最重要的设计 ★
SETTINGS index_granularity = 8192; -- 每 8192 行一个索引标记ORDER BY 的选择是 ClickHouse 性能的命门:
原则 1:最常用的过滤列放最前面
原则 2:基数低的列放前面(利于压缩和跳过)
原则 3:一般按 (时间, 高频过滤维度, 细粒度ID) 排列
★ ClickHouse 的主键索引是【稀疏索引】★
每 8192 行记录一个标记 → 索引极小,可全内存
查询时通过索引定位到候选数据块(granule),再扫描跳数索引(Data Skipping Index)—— 二级过滤:
ALTER TABLE events ADD INDEX idx_url url TYPE bloom_filter(0.01) GRANULARITY 4;
ALTER TABLE events ADD INDEX idx_dur duration TYPE minmax GRANULARITY 4;MergeTree 家族选型:
| 引擎 | 用途 |
|---|---|
| MergeTree | 基础,只追加 |
| ReplacingMergeTree | 按 ORDER BY 去重(后台异步,不保证即时) |
| SummingMergeTree | 自动按维度求和,预聚合 |
| AggregatingMergeTree | 保存聚合中间态,配合物化视图 |
| CollapsingMergeTree | 用 sign 列实现"折叠"更新 |
| Replicated*MergeTree | 上述各种的副本版本 |
15.4 StarRocks 核心:主键模型与物化视图
-- ★ 主键模型:支持高频 upsert 且查询性能不衰减
CREATE TABLE orders_rt (
order_id BIGINT NOT NULL,
user_id BIGINT,
amount DECIMAL(10,2),
status VARCHAR(32),
update_time DATETIME
)
PRIMARY KEY (order_id) -- ★ 主键模型
DISTRIBUTED BY HASH(order_id) BUCKETS 16
PROPERTIES (
"replication_num" = "3",
"enable_persistent_index" = "true" -- 主键索引持久化,降低内存占用
);四种数据模型对比:
| 模型 | 特点 | 适用 |
|---|---|---|
| Duplicate Key | 允许重复,只追加 | 日志明细 |
| Aggregate Key | 导入时自动聚合 | 固定维度的报表 |
| Unique Key | 按主键 upsert(Merge-on-Read) | 有更新需求的场景 |
| Primary Key | Delete+Insert 实现,查询无合并开销 | ★ CDC 实时同步首选 ★ |
异步物化视图(自动查询改写,杀手级特性):
CREATE MATERIALIZED VIEW mv_daily_gmv
DISTRIBUTED BY HASH(dt)
REFRESH ASYNC EVERY (INTERVAL 5 MINUTE)
AS
SELECT DATE(update_time) AS dt,
status,
COUNT(*) AS order_cnt,
SUM(amount) AS gmv
FROM orders_rt
GROUP BY DATE(update_time), status;
-- ★ 神奇之处:用户查原表,优化器【自动改写】去查物化视图 ★
SELECT DATE(update_time), SUM(amount) FROM orders_rt
WHERE update_time >= '2026-09-01' GROUP BY DATE(update_time);
-- 用户不用改 SQL,性能自动提升几十倍湖仓联邦(直接查 Iceberg):
CREATE EXTERNAL CATALOG iceberg_catalog PROPERTIES (
"type" = "iceberg",
"iceberg.catalog.type" = "hive",
"hive.metastore.uris" = "thrift://172.16.146.220:9083"
);
-- ★ 内表(热数据)JOIN 外表(冷数据),一条 SQL 搞定
SELECT o.user_id, SUM(o.amount), h.first_order_date
FROM orders_rt o -- StarRocks 内表
JOIN iceberg_catalog.lakehouse.dim_user_history h -- Iceberg 外表
ON o.user_id = h.user_id
GROUP BY o.user_id, h.first_order_date;🎯 这是 2026 年最重要的架构趋势之一:OLAP 引擎直接查数据湖, 热数据放本地 SSD、冷数据留在湖上,一套 SQL 打通,不用再做数据搬运。
第 16 章:终极实战 —— 端到端实时数仓
16.1 架构设计
┌──────────┐ binlog ┌────────────┐ ┌───────────────┐
│ MySQL │───────────►│ Flink CDC │─────────►│ Kafka │
│ 业务库 │ │ 整库同步 │ │ ods_* topics │
└──────────┘ └────────────┘ └───────┬───────┘
│
┌──────────────────────────────┤
▼ ▼
┌─────────────────┐ ┌─────────────────┐
│ Flink SQL │ │ Flink SQL │
│ DWD 清洗打宽 │ │ DWS 实时聚合 │
└────────┬────────┘ └────────┬────────┘
│ │
▼ ▼
┌───────────────────────┐ ┌──────────────────┐
│ Iceberg 数据湖 │ │ StarRocks │
│ ODS/DWD/DWS 全量历史 │ │ 实时服务层 │
│ (Spark 批处理加工) │ │ (亚秒级查询) │
└───────────┬───────────┘ └────────┬─────────┘
│ │
└──────────┬───────────────────┘
▼
┌─────────────────┐
│ BI / 实时大屏 │
└─────────────────┘16.2 分层实现
Layer 1:ODS —— CDC 入湖入 Kafka
-- Flink SQL
SET 'execution.checkpointing.interval' = '60s';
SET 'table.exec.sink.upsert-materialize' = 'NONE';
CREATE TABLE ods_orders_cdc (
order_id BIGINT, user_id BIGINT, item_id BIGINT,
amount DECIMAL(10,2), status STRING, create_time TIMESTAMP(3),
update_time TIMESTAMP(3),
PRIMARY KEY (order_id) NOT ENFORCED
) WITH (
'connector'='mysql-cdc', 'hostname'='192.168.1.50', 'port'='3306',
'username'='cdc_user', 'password'='***',
'database-name'='shop', 'table-name'='orders',
'scan.incremental.snapshot.enabled'='true'
);
-- 双写:既进 Kafka(给实时链路),也进 Iceberg(给离线链路)
CREATE TABLE kafka_ods_orders (
order_id BIGINT, user_id BIGINT, item_id BIGINT,
amount DECIMAL(10,2), status STRING, create_time TIMESTAMP(3),
update_time TIMESTAMP(3),
PRIMARY KEY (order_id) NOT ENFORCED
) WITH (
'connector'='upsert-kafka', -- ★ upsert-kafka 支持 changelog
'topic'='ods_orders',
'properties.bootstrap.servers'='kafka:9092',
'key.format'='json', 'value.format'='json'
);
EXECUTE STATEMENT SET
BEGIN
INSERT INTO kafka_ods_orders SELECT * FROM ods_orders_cdc;
INSERT INTO ice.lakehouse.ods_orders SELECT * FROM ods_orders_cdc;
END;💡
STATEMENT SET让多个 INSERT 共用一个 Source,只读一次 binlog,这是生产必备技巧。
Layer 2:DWD —— 实时打宽
-- 维度表(Lookup Join 用)
CREATE TABLE dim_user_lookup (
user_id BIGINT, user_name STRING, city STRING, level INT,
PRIMARY KEY (user_id) NOT ENFORCED
) WITH (
'connector'='jdbc',
'url'='jdbc:mysql://192.168.1.50:3306/shop',
'table-name'='users',
'lookup.cache.max-rows'='100000', -- ★ 缓存,避免每条数据查库
'lookup.cache.ttl'='10min',
'lookup.max-retries'='3'
);
-- 打宽:订单 + 用户维度
CREATE TABLE dwd_order_wide (
order_id BIGINT, user_id BIGINT, user_name STRING, city STRING,
user_level INT, item_id BIGINT, amount DECIMAL(10,2),
status STRING, create_time TIMESTAMP(3),
PRIMARY KEY (order_id) NOT ENFORCED
) WITH ('connector'='upsert-kafka', 'topic'='dwd_order_wide', ...);
INSERT INTO dwd_order_wide
SELECT
o.order_id, o.user_id, u.user_name, u.city, u.level,
o.item_id, o.amount, o.status, o.create_time
FROM kafka_ods_orders AS o
LEFT JOIN dim_user_lookup FOR SYSTEM_TIME AS OF o.proctime AS u -- ★ 时态 Join
ON o.user_id = u.user_id;📐
FOR SYSTEM_TIME AS OF是 Lookup Join 的语法,含义是"用处理这条数据的时刻去查维表"。 相比双流 Join,它不需要保存维表状态,内存占用极小,是维度关联的首选。
Layer 3:DWS —— 实时聚合
CREATE TABLE dws_city_gmv (
window_start TIMESTAMP(3), window_end TIMESTAMP(3),
city STRING, order_cnt BIGINT, gmv DECIMAL(20,2), uv BIGINT,
PRIMARY KEY (window_start, city) NOT ENFORCED
) WITH (
'connector'='starrocks',
'jdbc-url'='jdbc:mysql://sr-fe:9030',
'load-url'='sr-fe:8030',
'database-name'='rt_dw', 'table-name'='dws_city_gmv',
'username'='root', 'password'='',
'sink.buffer-flush.interval-ms'='5000'
);
-- ★ Flink 1.13+ 的 Window TVF 语法(比老的 GROUP BY TUMBLE 更强大)
INSERT INTO dws_city_gmv
SELECT
window_start, window_end, city,
COUNT(DISTINCT order_id) AS order_cnt,
SUM(amount) AS gmv,
COUNT(DISTINCT user_id) AS uv
FROM TABLE(
TUMBLE(TABLE dwd_order_wide, DESCRIPTOR(create_time), INTERVAL '5' MINUTE)
)
WHERE status = 'paid'
GROUP BY window_start, window_end, city;Layer 4:离线补数与校验(Lambda 的正确用法)
即使有了实时链路,仍然需要离线链路做数据校验和修正:
-- Spark SQL:T+1 用 Iceberg 全量数据重算,作为"真值"
CREATE OR REPLACE TABLE ice.lakehouse.dws_city_gmv_daily AS
SELECT
DATE(create_time) AS dt, city,
COUNT(DISTINCT order_id) AS order_cnt,
SUM(amount) AS gmv,
COUNT(DISTINCT user_id) AS uv
FROM ice.lakehouse.dwd_order_wide
WHERE status = 'paid' AND DATE(create_time) = '${dt}'
GROUP BY DATE(create_time), city;
-- ★ 对账 SQL:实时 vs 离线的差异监控 ★
SELECT
r.dt, r.city,
r.gmv AS rt_gmv,
b.gmv AS batch_gmv,
ROUND(ABS(r.gmv - b.gmv) / NULLIF(b.gmv,0) * 100, 4) AS diff_pct
FROM (SELECT dt, city, SUM(gmv) gmv FROM starrocks.rt_dw.dws_city_gmv
WHERE dt='${dt}' GROUP BY dt, city) r
FULL OUTER JOIN ice.lakehouse.dws_city_gmv_daily b
ON r.dt = b.dt AND r.city = b.city
WHERE ABS(COALESCE(r.gmv,0) - COALESCE(b.gmv,0)) / NULLIF(b.gmv,1) > 0.001; -- 差异 >0.1% 告警🎯 这是实时数仓的黄金法则:永远要有离线链路做真值校验。 实时链路会因为乱序、迟到、重启等原因产生偏差,没有对账机制的实时数仓是不可信的。
16.3 生产运维 Checklist
【Flink 作业】
□ 所有有状态算子设置了 UID
□ Checkpoint 间隔与超时合理(间隔 1min、超时 10min)
□ 开启 Unaligned Checkpoint + 对齐超时 30s
□ setTolerableCheckpointFailureNumber(3)
□ ExternalizedCheckpoint 设为 RETAIN_ON_CANCELLATION
□ 所有 State 配置了 TTL
□ RocksDB 用多块磁盘、开启增量 Checkpoint
□ 配置了反压告警和 Checkpoint 失败告警
【Kafka】
□ acks / min.insync.replicas 按场景配置
□ unclean.leader.election.enable = false
□ 事务超时 > Checkpoint 间隔 + 恢复时间
□ 分区数 ≥ Flink 并行度(否则有空闲 subtask)
□ 磁盘留有足够 buffer(至少 30% 空闲)
【Iceberg】
□ 每日定时 rewrite_data_files(小文件合并)
□ 每日定时 expire_snapshots(清理旧快照)
□ 每周 remove_orphan_files
□ format-version = 2(支持行级删除)
□ 监控表的文件数和元数据大小
【数据质量】
□ 实时 vs 离线对账任务(每小时/每日)
□ 主键唯一性校验
□ 关键指标的环比/同比波动告警
□ 数据延迟监控(业务时间 vs 处理时间的差值)
□ 数据量突增突降告警📋 中篇能力清单
理论层
- [ ] 能解释为什么说"批是流的特例"
- [ ] 能讲清 Lambda / Kappa / 湖仓流批一体各自的问题
- [ ] 能画出 Watermark 的推进过程,说出三个经典坑
- [ ] 能描述 Chandy-Lamport 算法和 Barrier 对齐
- [ ] 能说明端到端 Exactly-Once 的三个必要条件
- [ ] 能解释两阶段提交在故障各阶段的恢复行为
- [ ] 能说出 Kafka 高性能的四个原因
- [ ] 能画出 Iceberg 的三层元数据结构
- [ ] 能解释 CoW 和 MoR 的权衡
- [ ] 能说清 SCD Type 1/2/3 的差异和适用场景
实操层
- [ ] 能用 Flink SQL 完成 MySQL → Iceberg 的 CDC 同步
- [ ] 能配置生产级的 Checkpoint 参数
- [ ] 能在 Iceberg 上做 UPDATE/DELETE/MERGE 和时间旅行回滚
- [ ] 能写出 Iceberg 表的四个维护 CALL 语句
- [ ] 能用 Window TVF 写实时聚合
- [ ] 能用 Lookup Join 做维度打宽
- [ ] 能设计实时 vs 离线的对账 SQL
判断力层
- [ ] 能在 Iceberg / Hudi / Delta 之间做出有依据的选型
- [ ] 能在 ClickHouse / Doris / StarRocks 之间做出选型
- [ ] 能判断某场景该用 Exactly-Once 还是 At-Least-Once
- [ ] 能识别一个数仓分层设计的问题
🔮 下篇预告
《大数据技术硬核教程(下篇):治理 · 性能 · AI 融合 · 架构演进》
- 数据治理体系:元数据管理、血缘追踪、数据质量框架、数据合约(Data Contract)
- 成本优化实战:FinOps 方法论、存储分层、计算成本归因、Spot 实例策略
- 性能调优大全:数据倾斜的 7 种解法、Join 优化决策树、内存模型深度剖析
- 安全与合规:Ranger 权限体系、数据脱敏、GDPR/个保法 落地
- AI 与大数据融合:Feature Store、向量数据库、RAG 数据管道、Text-to-SQL 实战
- 架构演进:Data Mesh vs Data Fabric、Zero-ETL 趋势、Lakebase、Agentic Data Engineering
- 面试与职业:高频考点、系统设计题拆解、大数据工程师的能力模型与成长路径
📮 写给中篇读者
如果说上篇讲的是"怎么存、怎么算",中篇讲的就是"怎么让数据流动起来并且可信"。
实时计算最反直觉的一点是:它的难点不在"快",而在"准"。 Watermark、Checkpoint、Exactly-Once、对账机制——这些机制存在的唯一目的, 就是在一个数据会乱序、机器会宕机、网络会抖动的世界里,给出一个可以被信任的数字。
当你下次看到大屏上跳动的 GMV 时,想想它背后经过了多少道保障。