KEEL · 龙骨 · A CURRICULUM FOR THE AI ERA
流式数据管道课 · 课程导读 — keel 龙骨
流式数据管道课:CDC、事件时间与数仓分层 的参考信息:流式数据管道课 · 课程导读
你现在的起点
- 会写 Kafka 生产者和消费者,也大概知道 PostgreSQL 的 WAL 是什么;
- 但把「把数据从 A 搬到 B」当成一个连接器的配置问题——填个
bootstrap.servers、选个 connector class,跑起来了就以为成了; - 出过或见过这几类事:改了表结构以后同步断了、下游重跑一遍数据就对不上、延迟告警和实际业务口径对不上、一个没人管的复制槽把主库磁盘写满。
如果你符合上面这几条,这门课补的是链路视角:一条变更从产生到在分析表里可查,中间跨了几个进程、几块存储,每一段各自的失败模式是什么。
这门课和相邻课程的分工
- 《Kafka 工程课》 讲 Kafka 自己:分区、段、ISR、
acks、消费组、事务。这门课里 Kafka 只是管道的一段,不重复讲它的内部机制。 - 《PostgreSQL 工程课》06 章 讲逻辑复制的原理和
pg_recvlogical的基本用法。这门课从复制槽往下走,把解码结果真正接进一条可运行的管道。 - 这门课要做的事:把 CDC、消息队列、流处理语义、数仓分层连成一条你能自己跑起来、也能自己排查的链路。
一条能走通的学习路径
一条变更的旅程(先建立"跨了几块存储"的整链路图)
→ CDC 的原理与实现(读 WAL 而不是轮询;复制槽是负债也是保险)
→ 管道拓扑与 Kafka Connect(source/sink 边界、offset 存哪、SMT 在哪一段)
→ 事件时间与水位线(乱序、迟到、用延迟换正确性)
→ 状态与容错(快照、重启从哪继续、至少一次与恰好一次的真实差别)
→ 数仓分层与建模(ODS/DWD/DWS、宽表与预聚合的取舍)
→ OLAP 选型与管道运维(列存/预聚合/湖表,以及上线要盯的指标)
前两章把「数据怎么被取出来」讲透,中间两章讲「取出来之后怎么算对」,最后两章讲「算完以后放哪、怎么运维」。
章节地图
| 章 | 关键问题 | 你会亲手验证什么 |
|---|---|---|
| 00 · 一条变更的旅程 | 业务库一行更新到分析表出现这行,中间经过哪些进程与存储;什么时候该上管道 | 端到端跑通一次 orders → 分析表;量出 test_decoding 输出里一个事务的边界;用 WAL 字节数判断管道值不值 |
| 01 · CDC 的原理与实现 | 为什么读 WAL 而不是查表轮询;复制槽不消费会怎样;update 的 old key 从哪来 | 建 cdc_slot 看解码原文;灌 10 万行量出 16.25 MB WAL 堆积;REPLICA IDENTITY FULL 前后 UPDATE 输出对比 |
| 02 · 管道拓扑与 Kafka Connect | source/sink 的责任边界;offset 存在哪、断了怎么续;SMT 在链路的哪一段生效 | 跑 standalone Connect 的 FileStreamSource/FileStreamSink;读 offset 文件;重启后从 offset=3 续读不重不漏;InsertField 用错类型看 task 报错 |
| 03 · 事件时间与水位线 | 处理时间与事件时间差在哪;水位线解决什么、又引入什么;迟到数据怎么定 | 用可运行切片跑同一批乱序数据:事件时间窗口金额 [40,138,255] vs 处理时间 [63,151,251];allowed_lateness 从 0 调到 5 让总金额从 433 回到 465 |
| 04 · 状态与容错 | 状态存哪、快照多久一次、重启后从哪继续;至少一次与恰好一次的真实差别 | kill -9 后从快照恢复,结果与不中断一致;确定性崩溃点让非幂等 sink 多写 1 条(13:10 写两次),幂等 sink 挡掉 |
| 05 · 数仓分层与建模 | ODS/DWD/DWS/ADS 各自解决什么;维度建模与宽表的取舍;回填怎么重跑 | 把 CDC 数据落进三层表(42500 / 28000 / 3 行);同一聚合查询 ODS 20.9 ms vs DWD 4.8 ms vs DWS 0.021 ms |
| 06 · OLAP 选型与管道运维 | 列存/预聚合/湖表各适合什么;该盯哪些指标;数据质量与血缘怎么管 | 100 万行 GROUP BY 312 ms vs 物化视图 0.018 ms;建视图 502 ms、刷新 488 ms;灌 20 万行后视图过期(100 万 vs 120 万) |
学完这门课你能做什么
- 拿到一条 CDC + Kafka + 数仓的链路,能说清每一跳跨了哪个进程、状态记在哪块存储上、断了从哪续;
- 能对 AI 或同事给出的管道方案做评审:复制槽有没有人消费、offset 存在哪、坏掉的是 task 还是 worker、迟到数据口径有没有写进契约;
- 能解释「至少一次」在真实系统里会以什么形式多算一次,以及幂等键应该建在哪一列上;
- 决定预聚合和宽表时,能拿出自己量到的查询耗时和刷新耗时,而不是只引用一张选型图。
前置要求
- 会基本的 SQL、能跑起 PostgreSQL 和 Kafka(或愿意照着章节里的命令跑);
- 了解日志追加与消费者位点这类通用概念(不懂也行,第 00、01 章会从现场讲起);
- 建议先读 《Kafka 工程课》00 章 · 提交日志不是队列,把「消息被消费后不会被删除」这个前提建立起来。