CalcGuide · 技术博客主页 / 一页纸学习计划
🟠

Schema 演进与批流一体:从表格式版本化到统一存储

分类:数据与人工智能 · 路径:docs/topics/schema-evolution-and-batch-stream-unified/README.md

#schema-evolution#lakehouse#batch-stream-unified#iceberg#hudi

用表格式 Schema Evolution + Lakehouse 统一存储,把批处理与流处理跑在同一份数据上

父主题

数据工程基础:从 ETL 到数仓与流处理

子主题(0)

Schema 演进与批流一体:从表格式版本化到统一存储

0. 元信息

1. 学习路线

Schema 演进兼容性矩阵 → Iceberg Schema Evolution → Hudi Schema Evolution → Delta Schema Evolution → 批流一体存储 → 流批 SQL 引擎 → 端到端一致性验证

2. 阶段周数分配

阶段1.5 周方案2 周方案备注
1. 兼容性矩阵0.5 天1 天向后 / 向前 / 双向兼容
2. Iceberg Schema Evolution1.5 天2 天加列 / 改类型 / 重命名
3. Hudi Schema Evolution1 天1.5 天Schema Provider + Inline
4. Delta Schema Evolution1 天1.5 天columnMapping 模式
5. 批流一体存储1.5 天2 天Iceberg + Spark 批读 + Flink 流写
6. 端到端一致性0.5 天1 天批流同查询对比

每天 1.5~2 小时。1.5 周方案专注前 5 阶段;2 周方案多一天做端到端一致性。

3. 核心知识 / 产出 / 标准表

阶段核心知识实践产出可观察学会标准
1. Schema 演进兼容性向后兼容 / 向后不兼容、列默认值、分区演进一份兼容性矩阵能区分 4 种兼容性与边界
2. Iceberg Schema EvolutionADD COLUMN / RENAME / UPDATE TYPE / DROP / REPLACE一份 Iceberg Schema 演进演练(3 阶段)能加列/改类型不改下游
3. Hudi Schema EvolutionSchema Provider、Inline/External一份 Hudi Schema 演进能解释 Hudi 与 Iceberg 差异
4. Delta Schema EvolutionschemaTracking、columnMapping一份 Delta Schema 演进能解释 Delta 与 Iceberg 差异
5. 批流一体存储Lakehouse + 流批 SQL 引擎一份 Iceberg 表同时被 Spark 批读 + Flink 流写能用同一份表跑批流两种查询
6. 端到端一致性流处理写入 → 表格式 → 批读回填一份批流结果对比报告能验证批流同一查询产出相同

4. 第一周(每天 1.5~2 小时)

Day 1 约定:本计划使用 Iceberg on MinIO + Spark 3.5 + Flink 1.18。Schema 演进前先写兼容性矩阵(向后 / 向前 / 双向),再用 ALTER TABLE 操作。验证用 Time Travel 读历史数据。所有演进必须可回滚:演进失败时 CALL system.rollback_to_snapshot

任务当天交付自检
Day 1装 Docker Compose + MinIO + Spark + Flink + Trino;起 docker compose up -d;用 spark-sql 建 Iceberg 表 dw.orders 含 2 列(order_id BIGINT, amount DECIMAL(10,2))Iceberg 表 + 初始 schemaDESCRIBE TABLE dw.orders 显示 2 列;SELECT * FROM dw.orders.snapshots 至少 1 个快照
Day 2兼容性矩阵:写一份 4 列 × 4 行的表(操作 × 向后兼容 × 向前兼容 × 双向兼容):加列、改类型(int→bigint)、重命名、删列兼容性矩阵 Markdown矩阵 4×4 填齐;用 ADD COLUMN 后老 Job 读旧数据 OK;改 int→string 故意报失败
Day 3Iceberg Schema Evolution:写 3 步演进脚本(加 user_id 列 → 重命名 amounttotal_amount → 改 user_id 类型 BIGINT);每步用 SELECT * FROM dw.orders FOR SYSTEM_TIME AS OF 验证历史可读3 步演进脚本 + 历史快照截图演进后 DESCRIBE TABLE 显示新 schema;Time Travel 查演进前快照仍能返回旧字段;SELECT snapshot_id, operation FROM dw.orders.snapshots 列出 ADD COLUMN / REPLACE / ALTER
Day 4Hudi Schema Evolution:用 hoodie.datasource.write.schema.provider.class=org.apache.hudi.avro.HoodieAvroSchemaProvider 配 Inline 模式;加列 + 重命名;对比与 Iceberg 的差异Hudi 演进脚本 + 对比表Hudi 加列后 hoodie_schema 同步更新;重命名需要在 schema provider 配映射;对比 Iceberg 的 column_id 与 Hudi 的 schema provider
Day 5Delta Schema Evolution:建 Delta 表;开 delta.columnMapping.mode = 'name';重命名列;用 DESCRIBE HISTORY 查 schema 演进历史Delta 演进脚本 + history 截图Delta 重命名不破坏下游;DESCRIBE HISTORY dw_delta.orders 显示多个 commit(ADD COLUMNS / RENAME COLUMN);对比 Iceberg 与 Delta 的列映射机制
Day 6批流一体:起一个 Flink Job 把 Kafka 订单事件流写入 Iceberg dw.orders;同时用 Spark 跑批读;对同一查询(SELECT user_id, SUM(total_amount) GROUP BY user_id)对比 Flink SQL 与 Spark SQL 结果批流对比报告Flink Job 跑通;批流同查询结果完全一致(行数与金额一致);用 EXPLAIN 看 Flink SQL 计划
Day 7步骤 A:跑通 Iceberg 表端到端 Schema 演进 + 批流一致(演进 3 次 + 批流同查 5 次);步骤 B:补齐 4 类边界(演进失败回滚、批流结果不一致、column mapping 漏配、Time Travel 找不到快照)演进 + 一致性报告 + 故障演练pytest -k test_schema_evo 4 个用例全过;每类故障复现 + 防护措施写到 notes/week1-day7.md

Day 7 执行次序

步骤 A —— 最小可用演进 + 一致性(60~90 分钟)

  1. 起 Iceberg 表 2 列;
  2. 演进 3 次(加列 + 改类型 + 重命名);
  3. 起 Flink Job 写 Iceberg;同时 Spark 批读;
  4. 跑批流同查询 5 次,结果必须完全一致;
  5. notes/batch-stream-parity.md

步骤 B —— 4 类边界用例(30~45 分钟)

#用例期望行为验证命令
B1演进失败(不兼容类型)报错并保留旧 schemapytest -k test_bad_evolution
B2批流结果不一致报警 + 阻塞下游pytest -k test_parity_violation
B3column mapping 漏配列重命名破坏下游pytest -k test_column_mapping
B4Time Travel 找不到快照报错 + 提示用最近快照pytest -k test_snapshot_missing

Day 7 当天必完成步骤 A;步骤 B 至少完成 B1、B2。

5. 阶段通用验收

  1. 不看答案独立重写 Iceberg 3 步 Schema 演进 + 批流一致性对比;
  2. 用自己的话解释”为什么需要兼容性矩阵""为什么 Iceberg 用 column_id 而非列名""为什么批流结果可能不一致”;
  3. 画一张图:兼容性矩阵表 + Iceberg manifest 演进图 + 批流架构图(三件套之一);
  4. 测试演进失败回滚、批流结果不一致、column mapping 漏配、Time Travel 失败、并发演进 5 类边界;
  5. 准备至少 3 组自定义 schema 演进场景并贴出实际 snapshot 链;
  6. 记录演进次数、批流同查询结果差异(应为 0)、查询 P95、并发写入冲突次数 4 个核心指标;
  7. 能修改已有 schema(加列 / 改类型 / 换分区策略)并验证批流同查询结果一致。

交付存放:第 3 项的图、第 5 项的输出、第 6 项的指标,统一存到 week1/notes/week1/<table>/README.md

6. 最终验收

7. 综合项目

首选:批流一体数仓(必做:Kafka → Flink → Iceberg + Spark 批读 + Schema 演进演练)。
备选:跨表格式一致性对比(同一份数据写 Iceberg / Hudi / Delta,验证 Schema 演进与批流一致)。
备选:Schema Registry + 演进门禁(用 Confluent SR + Avro + Iceberg 演进 + 自动回滚)。

批流一体数仓必做要求:

任何综合项目都必须包含:

  1. 需求说明与数据契约(源 schema、目标 schema、SLA、演进策略);
  2. 兼容性矩阵 + 演进 SOP 文档;
  3. 核心代码(Iceberg DDL + Flink Job + Spark 批读 + 演进脚本);
  4. 边界测试(演进失败回滚、批流结果不一致、并发演进、Time Travel 失败、Schema Registry 不匹配);
  5. 可观测(Iceberg snapshot 元数据 + Flink checkpoint + Spark metrics);
  6. README(设计取舍、表格式选型、演进门禁规则);
  7. 复盘记录(模拟 1 次不兼容演进失败 + 1 次批流不一致,写 retrospective.md);
  8. notes/ 规范
文件内容
notes/design.md兼容性矩阵 + 演进 SOP + 架构图
notes/test.md每组测试的输入 / 期望 / 实际 / 通过情况
notes/retrospective.md用时、难点、收获、改进点
notes/runbook.md常见失败 + 处置 SOP(演进失败 / 批流不一致 / 并发演进)
notes/iceberg-snapshots/演进前后 snapshot 截图(snapshots / manifests / history 视图)
notes/batch-stream-parity.md批流同查询对比表(5+ 次查询,0 差异为目标)
notes/compatibility-matrix.md4×4 兼容性矩阵(操作 × 向后 / 向前 / 双向)
notes/evolution-evidence/5 次演进的 DESCRIBE / SELECT 截图 + Time Travel 验证

本主题贡献

Schema 演进的核心矛盾是”读历史数据用旧 schema vs 写新数据用新 schema”。本子主题专门讲清 Delta Lake / Iceberg + Avro + CDC 如何把这条边界守住——以 backward compat 为演进默认路径,Avro 为跨语言契约载体,CDC(Debezium / Flink CDC)为流式变更入口,Delta Lake column mapping 为兼容演进提供列追踪;不重复父主题的 Kafka 基础与 Airflow 编排。

3 职责

  1. 用 Avro + Confluent Schema Registry 定义兼容性矩阵(BACKWARD / FULL / BACKWARD_TRANSITIVE),演进操作(加列 / 删列 / 改类型 / 重命名)必须落在 backward compat 范围内,下游 Avro 反序列化永不报错。
  2. 用 Delta Lake delta.columnMapping.mode='name' + Iceberg column_id 做列追踪,重命名不破坏下游;演进前用 Time Travel 读历史验证,演进失败用 CALL system.rollback_to_snapshot 一键回退。
  3. 用 CDC(Debezium 抓 MySQL binlog → Kafka topic cdc.orders → Flink CDC Source → Iceberg dw.orders)把 OLTP 变更流式接入 Lakehouse;批读用 Spark / Trino / Flink SQL 三引擎,跑同一查询对比批流结果必须 0 差异。

4 交付物

  1. 一份 Avro 契约 + Schema Registry 兼容性配置(BACKWARD / FULL / BACKWARD_TRANSITIVE),含加列 / 改类型 / 重命名 3 步演进脚本,每次演进后老 Job 读旧数据 OK。
  2. 一份 Delta Lake Schema 演进演练(开 columnMapping.mode='name' → 重命名列 → DESCRIBE HISTORY 查演进历史),含 Iceberg ADD COLUMN + RENAME COLUMN 对比表,配 format-version='2'
  3. 一份 CDC 接入管道(Debezium MySQL connector → Kafka cdc.orders → Flink MySqlSource + IcebergSink),含 binlog row 模式配置 + enableCheckpointing(60_000, EXACTLY_ONCE)
  4. 一份批流一致性对比报告(5+ 次同查询 Spark vs Flink SQL,结果 0 差异),含演进 SOP 文档(演进前兼容性检查 → 演进 → 灰度 → 全量 → 回滚预案)。

3 指标

  1. Schema 演进 backward compat 通过率 100%(所有演进操作都通过 Schema Registry 兼容性校验,老消费者零修改)。
  2. 批流同查询 0 差异(5 次查询 Spark vs Flink SQL 结果完全一致,含行数与聚合值)。
  3. 演进回滚时长 < 1 min(CALL system.rollback_to_snapshot 执行到表 schema 与数据恢复)。

8. 推荐开源资料

阶段角色资料链接用法
全部规范Apache Iceberg Spec v2https://iceberg.apache.org/spec/必读 schema section
1~2演进Iceberg Schema Evolution 文档https://iceberg.apache.org/docs/latest/evolution/操作指南必读
3HudiApache Hudi Schema Evolutionhttps://hudi.apache.org/docs/schema_evolutionHudi 演进操作
4DeltaDelta Schema 文档https://docs.delta.io/latest/delta-schema.htmlcolumnMapping 模式
5~6批流一体Iceberg Flink 集成https://iceberg.apache.org/docs/latest/flink/IcebergSink 用法
5~6批流一体Flink SQL 官方文档https://nightlies.apache.org/flink/flink-docs-stable/docs/dev/table/流批 SQL 一致性
全部协议Confluent Schema Registryhttps://docs.confluent.io/platform/current/schema-registry/index.html兼容性与演进门禁
全部工具Apache Avrohttps://avro.apache.org/docs/current/Schema 定义
全部工具Protocol Buffershttps://protobuf.dev/跨语言 Schema
6引擎Spark Structured Streaminghttps://spark.apache.org/docs/latest/structured-streaming-programming-guide.htmlSpark 流读 Iceberg
全部对比Apache Paimonhttps://paimon.apache.org/流式 Lakehouse 新选型
全部哲学Kleppmann《DDIA》Ch.4https://dataintensive.net/编码与演进

许可证提示:Iceberg / Hudi / Delta / Flink / Spark 都是 Apache-2.0;Confluent Schema Registry 有社区版与商业版之分。复制 Apache 项目示例前请保留 LICENSE 与 NOTICE。默认做法是读规范后自己写演进脚本,而不是复制官方 Example。

默认使用顺序:先读 Iceberg Spec v2 schema section → 写兼容性矩阵 → 装 Iceberg on MinIO + Spark + Flink → 跑 Iceberg 3 步演进 → 对照 Hudi / Delta 演进差异 → 配 Flink 写 Iceberg + Spark 批读 → 跑批流同查询 5 次对比 → 用 Time Travel 验证演进前历史可读 → 写 notes/retrospective.md

9. 学习资料汇聚(v0.3 自包含)

本节由本计划生成。链接指向原始材料或作者公开内容。规范会演进,记录时务必写明版本与日期。

9.1 背景与动机

传统数据仓库(Teradata、Snowflake、BigQuery)支持 Schema 演进,但实时性差;Kafka 等消息流支持实时,但不支持复杂查询。Lakehouse 通过 Iceberg/Delta/Hudi 等表格式,把”ACID 表 + 实时流写入 + 批 SQL 查询”三件事合并到同一份 Parquet 数据上。Schema 演进(Schema Evolution)从传统数据库继承而来,但在流处理场景下面临新约束——流读需要兼容历史事件,批读需要兼容新写。今天数据团队用”列添加 / 列重命名 / 列删除”等演进操作时,必须保证”读历史数据 + 写新数据”两种路径都不报错。

9.2 概念地图

flowchart LR
  Schema[Schema] --> Compat[兼容性矩阵]
  Compat --> Back[向后兼容]
  Compat --> Forw[向前兼容]
  Compat --> BackForw[双向兼容]
  Ice[Iceberg] --> SchemaEvo[Schema Evolution]
  Hu[Hudi] --> SchemaEvo
  Delta[Delta] --> SchemaEvo
  Lake[Lakehouse] --> Unify[批流一体]
  Unify --> BatchSQL[批 SQL Spark/Trino]
  Unify --> StreamSQL[流 SQL Flink]
  BatchSQL --> Verify[结果一致性]
  StreamSQL --> Verify

9.3 基础知识讲解

9.3.1 经典论文 / 规范

资料贡献读法
Apache Iceberg Spec v2Iceberg Schema Evolution必读 schema section
Apache Hudi RFC-29Hudi Schema Evolution选读
Delta Lake ProtocolDelta Schema Evolution选读

9.3.2 经典书籍

侧重用法
Reis & Housley, Fundamentals of Data Engineering(O’Reilly 2022)数据工程生命周期选读第 9 章
Kleppmann, Designing Data-Intensive Applications(O’Reilly 2017)数据系统原理选读第 4 章

9.3.3 优秀博客与文档

资料特点用法
Apache Iceberg SpecSchema Evolution必读
Apache Iceberg Schema Evolution操作指南必读
Apache Hudi Schema EvolutionHudi 操作选读
Delta SchemaDelta 操作选读
Apache Paimon流式 Lakehouse选读

9.3.4 核心人物

人物影响材料
Ryan BlueIceberg 起源Netflix 公开演讲
Vinoth ChandarHudi 起源Uber 公开演讲
Michael ArmbrustDelta 起源Databricks 博客

9.3.5 开发方法

方法动作何时用
Compatibility-matrix first先定兼容矩阵,再改 schema任何生产表
Add-only by default加列优于改类型演进默认路径
Read-side cast读时兼容写时新类型历史数据兼容
Batch-stream parity test跑同一查询对比批流任何流批一体表

9.4 经典问题与经典案例(≥5 道)

#问题重要性最简答案
1加列是不是 always-safe兼容性大部分是,但下游需要重新读
2改类型什么时候安全兼容性兼容类型(int→bigint)安全,否则不兼容
3列重命名会破坏下游吗影响Iceberg 用 column id 安全,Parquet 用列名不安全
4分区演进怎么做重写历史Iceberg 用 partition spec evolution
5批流结果为什么不一致验证用相同 Watermark + 相同 SQL
6流处理写能不能加列Flink/Spark可以,但下游要兼容
7Iceberg / Hudi / Delta 演进差异选型Iceberg 用 column id,Hudi 用 schema provider,Delta 用 column mapping
8列删除怎么做影响Iceberg 默认安全,Parquet 需重写
9演进后多久可以回滚安全Iceberg Time Travel 可以
10批流一体的性能瓶颈优化用 partition + z-order / data skipping

9.5 学习难点

概念难点

难点突破路径
兼容性矩阵画一个表,写 4 种兼容性
column id vs name用 Iceberg vs Parquet 对比
流批 schema drift用 Schema Registry

思维难点

难点突破路径
演进路径选错先写兼容性矩阵再操作
批流结果差异跑相同 SQL 对比

工程难点

难点突破路径
Iceberg 并发演进配 commit retry
Flink 写 Iceberg schema 不匹配用 Avro/Protobuf 注册 schema
Delta 列重命名破坏下游用 column mapping mode

9.6 技术标准与接口

Entity

名称版本组织状态许可证
Apache Iceberg1.xApacheGAApache-2.0
Apache Hudi0.14+ApacheGAApache-2.0
Delta Lake2.x / 3.xLinux FoundationGAApache-2.0
Apache Paimon0.xApache活跃Apache-2.0
Apache Flink1.18+ApacheGAApache-2.0

Scope

Iceberg/Hudi/Delta/Paimon 是表格式;它们解决 Schema 演进与 ACID;不替代计算引擎与存储。

Structure

Ecosystem

Depth Tiers

层级能力标准
L0知道存在知道 Schema Evolution 与 Lakehouse
L1看得懂示例能读懂 Iceberg schema 操作
L2能正确调用能用 SQL/API 演进 schema
L3能解释与排错能定位演进失败、批流不一致
L4能设计与扩展能为新业务设计 schema 演进规范与批流一体架构

本子主题目标:L3

Source

9.7 关键代码

9.7.1 Iceberg Schema Evolution(加列 / 改类型 / 重命名)

-- 9.7.1 Iceberg Schema Evolution:三阶段演进
-- 初始 schema
CREATE TABLE dw.orders (
    order_id BIGINT,
    amount DECIMAL(10, 2)
) USING iceberg;

-- 阶段 1:加列
ALTER TABLE dw.orders ADD COLUMN user_id BIGINT;
ALTER TABLE dw.orders ADD COLUMN currency STRING;

-- 阶段 2:重命名(用 column id,旧数据仍可读)
ALTER TABLE dw.orders RENAME COLUMN amount TO total_amount;

-- 阶段 3:改类型(兼容类型:int -> bigint)
ALTER TABLE dw.orders ALTER COLUMN user_id TYPE BIGINT;

-- 验证:Time Travel 看历史数据
SELECT * FROM dw.orders FOR SYSTEM_TIME AS OF '2026-01-01';

9.7.2 Hudi Schema Evolution

# 9.7.2 PyHudi:Schema Evolution(HiveStyle 或 Inline)
from hudi import HudiTable

table = HudiTable(
    base_path="/data/hudi/orders",
    table_name="orders",
    storage_type="COPY_ON_WRITE",
)

# 加列
table.add_column("user_id", "bigint")
table.add_column("currency", "string")

# 重命名(需要 schema provider)
table.rename_column("amount", "total_amount", schema_provider="inline")

# 改类型
table.alter_column_type("user_id", "bigint")
// 9.7.3 Flink DataStream:写 Iceberg + 触发批读对比
DataStream<Order> orders = env.addSource(kafkaSource);
orders.addSink(new IcebergSink.Builder<Order>()
    .setTable(TableIdentifier.of("dw", "orders"))
    .setWriteParallelism(4)
    .build());

env.execute("stream-orders-to-iceberg");

// 同一份表用 Spark 批读,对比结果
// Spark SQL
val df = spark.read.format("iceberg").load("dw.orders")
df.createOrReplaceTempView("orders_batch")
spark.sql("SELECT user_id, SUM(total_amount) FROM orders_batch GROUP BY user_id")
  .show()

// 与 Flink SQL 跑相同查询对比结果

10. 常见误区

11. 所有知识点分类

  1. 编程语言
  2. 数据结构与算法
  3. 计算机基础
  4. 工程技术
  5. Web 与后端
  6. 前端与客户端
  7. 数据与人工智能
  8. 项目与职业能力
  9. 安全与可靠性

本计划归属:数据与人工智能 主 + 工程技术 辅。


直接依赖(3)

查看知识图谱