Skip to content

大数据技术硬核教程(中篇) ​

实时计算 · 湖仓一体 · 数据建模 ​

承上启下:上篇解决了"数据放哪、怎么算",中篇解决"数据怎么实时流动、怎么在湖上做事务、怎么组织成可用的资产"。 难度提升:本篇包含分布式快照算法、MVCC 实现、Exactly-Once 语义证明等硬核内容。 实操环境:EMR(Hadoop 3.2.1 + Spark 3.5.3 + Hive 3.1.3)+ 本篇新增 Flink / Kafka / Iceberg。


目录 ​


第 8 章:流处理的本质 —— 从批到流的范式跃迁 ​

8.1 一个颠覆性的认知:批是流的特例 ​

传统认知:

批处理:处理一堆静态数据
流处理:处理源源不断的数据
→ 两种不同的东西

Flink 的世界观(也是现代流处理的共识):

世界上只有一种数据:无界流(Unbounded Stream)

批处理  = 对无界流开了一个【有界窗口】
实时处理 = 对无界流开了一个【滑动的小窗口】

→ 批是流的特例,不是流是批的加速版

这个认知为什么重要?

因为它决定了架构走向。如果你认为批和流是两种东西,你会建两套系统(Lambda 架构);如果你认为批是流的特例,你会建一套系统(Kappa / 流批一体)。

8.2 三代架构的血泪史 ​

第一代:Lambda 架构(2011-2018) ​

                    ┌─────────────────────┐
                    │   批处理层(准确)     │
              ┌────►│  Spark/MR → HDFS    │────┐
              │     │  T+1 全量重算        │    │
┌──────────┐  │     └─────────────────────┘    │   ┌──────────┐
│ 数据源    │──┤                                ├──►│ 服务层    │
│ Kafka    │  │     ┌─────────────────────┐    │   │ 合并结果  │
└──────────┘  │     │   速度层(快但可能错)  │    │   └──────────┘
              └────►│  Storm/Flink → Redis│────┘
                    │  秒级近似结果         │
                    └─────────────────────┘

致命问题:

  1. 同一套业务逻辑要写两遍(Spark 一份 + Flink 一份)
  2. 两套代码的口径几乎必然会漂移,对不上数
  3. 运维两套集群,成本翻倍
  4. 改一次需求要改两处,还要保证一致

💬 真实案例:某大厂曾出现批流结果差异 3%,排查两周才发现是流层用了 < 而批层用了 <=。

第二代:Kappa 架构(2014 提出) ​

┌──────────┐    ┌─────────────────────────┐    ┌──────────┐
│ 数据源    │───►│  Kafka(长期保留)        │───►│ Flink    │───► 结果
└──────────┘    │  可回溯到任意历史位点      │    └──────────┘
                └─────────────────────────┘
                         │
                需要重算时:从头重放

核心思想:既然流能表达一切,那就只用流。需要"批处理"时,把 Kafka 从头重放一遍。

理想很美好,现实问题:

  1. Kafka 保存全量历史数据成本极高(Kafka 用本地磁盘,不是廉价存储)
  2. 重放 3 年数据可能要跑好几天
  3. 流式 SQL 表达能力当年不如批(复杂 JOIN、多维分析很吃力)
  4. 历史数据的 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 的数据 = 迟到数据

生成策略实操 ​

java
// 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)→ 窗口迟迟不触发 → 延迟高

★ 这是【延迟】和【完整性】的根本权衡,没有银弹 ★

工业界的三层兜底方案:

java
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 次下单未支付"预警 ​

java
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();
        }
    }
}

🎯 状态管理三铁律:

  1. 必设 TTL——不设 TTL 的状态是定时炸弹,跑几个月必 OOM
  2. 必用定时器清理——TTL 是兜底,业务逻辑的主动清理更及时
  3. 状态越小越好——Checkpoint 时要序列化,状态大直接导致 Checkpoint 超时

状态后端选型(生产关键决策) ​

状态后端存储位置容量上限访问延迟适用场景
HashMapStateBackendJVM 堆内存受堆大小限制(<10GB)极快(纳秒)小状态、低延迟要求
EmbeddedRocksDBStateBackend本地磁盘(RocksDB)TB 级较慢(微秒,需序列化)大状态生产标配

RocksDB 调优(生产必备):

yaml
# 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 到达时间差很大(数据倾斜/反压),阻塞时间长 → 延迟飙升。

Aligned(对齐):
  等所有 Barrier 到齐 → 反压时可能等几分钟 → Checkpoint 超时失败

Unaligned(非对齐):
  第一个 Barrier 到达就立即快照
  把【正在传输中的数据】(in-flight data)也存进 Checkpoint
  → Checkpoint 时间不受反压影响
  → 代价:Checkpoint 体积变大
yaml
# 生产配置建议
execution.checkpointing.unaligned: true
execution.checkpointing.aligned-checkpoint-timeout: 30s
# 含义:先尝试对齐,30 秒内没对齐完就自动切换到非对齐模式
# ★ 这是两种模式的最佳折中 ★

Checkpoint 完整配置模板 ​

java
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 ​

维度CheckpointSavepoint
触发方Flink 自动定期人工触发
目的故障自动恢复版本升级、集群迁移、A/B 测试
格式引擎内部格式,可能不跨版本兼容标准格式,保证跨版本兼容
存储增量(RocksDB)全量
生命周期自动清理旧的手动管理
bash
# 触发 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

java
stream.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 就不成立。

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,恢复时能把"悬挂的事务"提交掉。这要求外部系统的事务在一段时间内不能被清理。

java
// 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 幂等:

sql
-- 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 加密后零拷贝会失效(因为需要在用户态做加密),吞吐会明显下降。这是安全和性能的真实权衡。

秘诀四:批量 + 压缩 ​

java
// 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 生产配置:可靠性与性能的权衡矩阵 ​

java
// ===== 场景 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  ...

一次查询的执行过程 ​

sql
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 三者的核心差异(选型决策) ​

维度IcebergHudiDelta Lake
出身Netflix → ApacheUber → ApacheDatabricks
设计重心通用表格式、多引擎中立流式 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) ​

sql
-- ❌ 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) ​

sql
-- 业务初期数据少,按月分区
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
读多写极少     → CoW

12.4 Iceberg 实操(在你的 EMR 集群上) ​

环境准备 ​

bash
# 下载 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
bash
# 启动配置了 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 实操 ​

sql
-- 建库建表
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 *;

时间旅行与回滚(运维救命功能) ​

sql
-- 查看所有快照
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;   -- 数据回来了!

这个实验建议让学生做一遍——它会让人对"数据湖也能有事务"产生直观震撼。

表维护(生产必备的定时任务) ​

sql
-- ① 小文件合并(最重要的维护操作)
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 准备工作 ​

sql
-- 检查 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 格式记录每一行的前后镜像,才能还原出完整的变更事件。

方式一:纯 SQL(推荐,开发效率高) ​

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 秒内同步到数据湖。

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 连接 → 源库扛不住。

sql
-- ★ 用正则一次同步整库(一个作业、一个 binlog 连接)
CREATE TABLE all_tables (
    ...
) WITH (
    'connector' = 'mysql-cdc',
    'database-name' = 'shop',
    'table-name' = '.*',          -- 正则匹配所有表
    ...
);

或使用 Flink CDC 3.x 的 YAML Pipeline(2024 年新增,极大简化整库同步):

yaml
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: 4
bash
bin/flink-cdc.sh mysql-to-iceberg.yaml

Flink 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 的表结构:

sql
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:

sql
-- 步骤 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)的解法:把软件工程实践引入数据开发。

sql
-- 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
yaml
# 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: 0
bash
dbt run          # 按依赖顺序执行所有模型
dbt test         # 运行所有数据质量测试
dbt docs generate && dbt docs serve   # ★ 自动生成带血缘图的文档网站 ★

dbt 带来的四个革命性改变:

  1. 依赖自动推导 —— ref() 自动构建 DAG,不用手写调度依赖
  2. 数据测试即代码 —— 唯一性、非空、外键、值域全部自动化检查
  3. 文档与血缘自动生成 —— 点击任何一张表能看到上下游全链路
  4. 版本控制 + 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 三大引擎深度对比 ​

维度ClickHouseApache DorisStarRocks
架构无中心节点、Shared-NothingFE(元数据) + 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 家族 ​

sql
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)—— 二级过滤:

sql
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 核心:主键模型与物化视图 ​

sql
-- ★ 主键模型:支持高频 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 KeyDelete+Insert 实现,查询无合并开销★ CDC 实时同步首选 ★

异步物化视图(自动查询改写,杀手级特性):

sql
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):

sql
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 ​

sql
-- 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 —— 实时打宽 ​

sql
-- 维度表(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 —— 实时聚合 ​

sql
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 的正确用法) ​

即使有了实时链路,仍然需要离线链路做数据校验和修正:

sql
-- 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 时,想想它背后经过了多少道保障。

基于 Vite 强力驱动 | 纯静态轻量托管