Schema 演进与批流一体:从表格式版本化到统一存储
0. 元信息
- 主题路径:
docs/topics/data-engineering-basics/subtopics/schema-evolution-and-batch-stream-unified/README.md - 父主题:
data-engineering-basics - 适合对象:会写 SQL/Python、跑过批流两条流水线、希望用 Lakehouse 把两者合一并支持 Schema 演进的工程师
- 建议周期:1.5~2 周,每周 10~14 小时
- 前置知识:SQL、Python、Airflow;建议先完成父主题第 5 阶段
- 最终目标:能用 Iceberg/Hudi 在同一份表上同时跑批流查询并保证结果一致;能管理 Schema 演进(加列、改类型、改分区)不阻塞读写
1. 学习路线
Schema 演进兼容性矩阵 → Iceberg Schema Evolution → Hudi Schema Evolution → Delta Schema Evolution → 批流一体存储 → 流批 SQL 引擎 → 端到端一致性验证
2. 阶段周数分配
| 阶段 | 1.5 周方案 | 2 周方案 | 备注 |
|---|---|---|---|
| 1. 兼容性矩阵 | 0.5 天 | 1 天 | 向后 / 向前 / 双向兼容 |
| 2. Iceberg Schema Evolution | 1.5 天 | 2 天 | 加列 / 改类型 / 重命名 |
| 3. Hudi Schema Evolution | 1 天 | 1.5 天 | Schema Provider + Inline |
| 4. Delta Schema Evolution | 1 天 | 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 Evolution | ADD COLUMN / RENAME / UPDATE TYPE / DROP / REPLACE | 一份 Iceberg Schema 演进演练(3 阶段) | 能加列/改类型不改下游 |
| 3. Hudi Schema Evolution | Schema Provider、Inline/External | 一份 Hudi Schema 演进 | 能解释 Hudi 与 Iceberg 差异 |
| 4. Delta Schema Evolution | schemaTracking、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 表 + 初始 schema | DESCRIBE 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 3 | Iceberg Schema Evolution:写 3 步演进脚本(加 user_id 列 → 重命名 amount 为 total_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 4 | Hudi 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 5 | Delta 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 分钟)
- 起 Iceberg 表 2 列;
- 演进 3 次(加列 + 改类型 + 重命名);
- 起 Flink Job 写 Iceberg;同时 Spark 批读;
- 跑批流同查询 5 次,结果必须完全一致;
- 写
notes/batch-stream-parity.md。
步骤 B —— 4 类边界用例(30~45 分钟)
| # | 用例 | 期望行为 | 验证命令 |
|---|---|---|---|
| B1 | 演进失败(不兼容类型) | 报错并保留旧 schema | pytest -k test_bad_evolution |
| B2 | 批流结果不一致 | 报警 + 阻塞下游 | pytest -k test_parity_violation |
| B3 | column mapping 漏配 | 列重命名破坏下游 | pytest -k test_column_mapping |
| B4 | Time Travel 找不到快照 | 报错 + 提示用最近快照 | pytest -k test_snapshot_missing |
Day 7 当天必完成步骤 A;步骤 B 至少完成 B1、B2。
5. 阶段通用验收
- 不看答案独立重写 Iceberg 3 步 Schema 演进 + 批流一致性对比;
- 用自己的话解释”为什么需要兼容性矩阵""为什么 Iceberg 用 column_id 而非列名""为什么批流结果可能不一致”;
- 画一张图:兼容性矩阵表 + Iceberg manifest 演进图 + 批流架构图(三件套之一);
- 测试演进失败回滚、批流结果不一致、column mapping 漏配、Time Travel 失败、并发演进 5 类边界;
- 准备至少 3 组自定义 schema 演进场景并贴出实际 snapshot 链;
- 记录演进次数、批流同查询结果差异(应为 0)、查询 P95、并发写入冲突次数 4 个核心指标;
- 能修改已有 schema(加列 / 改类型 / 换分区策略)并验证批流同查询结果一致。
交付存放:第 3 项的图、第 5 项的输出、第 6 项的指标,统一存到
week1/notes/或week1/<table>/README.md。
6. 最终验收
-
独立设计并交付一套含 Schema 演进 + 批流一体的 Lakehouse,覆盖 Iceberg / Hudi / Delta 至少一种表格式;
-
至少完成 12 个实战项(分布建议):
子阶段 题目数量 难度 平台建议 兼容性矩阵 2 Easy / Medium Iceberg Spec + Confluent Blog Iceberg 演进 3 Medium Iceberg 官方 Examples Hudi / Delta 演进 3 Medium Apache Hudi / Delta 官方 Examples 批流一体 4 Medium / Hard Flink + Iceberg 案例 约束:至少 6 项达到 Medium,至少 2 项达到 Hard;每项必须留可运行 SQL + 演进前后 snapshot 截图。
-
完成 1 个综合项目:批流一体数仓(Kafka → Flink → Iceberg + Spark 批读 + 5 次 Schema 演进演练);
-
能用 15 分钟讲清兼容性矩阵、Iceberg column_id 设计、批流结果一致性验证、Schema 演进回归测试。
7. 综合项目
首选:批流一体数仓(必做:Kafka → Flink → Iceberg + Spark 批读 + Schema 演进演练)。
备选:跨表格式一致性对比(同一份数据写 Iceberg / Hudi / Delta,验证 Schema 演进与批流一致)。
备选:Schema Registry + 演进门禁(用 Confluent SR + Avro + Iceberg 演进 + 自动回滚)。
批流一体数仓必做要求:
- 输入:Kafka topic
orders实时事件 + 历史 CSV 批数据; - 输出:必输出(1)Iceberg 表
dw.orders同时被 Flink 流写与 Spark 批读;(2)5 次 Schema 演进(加列 2 次 + 改类型 1 次 + 重命名 1 次 + 改分区 1 次);(3)批流同查询对比报告(至少 5 次查询);(4)Time Travel 读历史数据截图;(5)兼容性矩阵 + 演进 SOP 文档; - 算法 / 工程:Flink Job 配
enableCheckpointing(60_000, EXACTLY_ONCE);Iceberg 配TBLPROPERTIES ('format-version'='2');批流对比用EXCEPT找差异; - 进阶可选:用 Paimon 做流式 Lakehouse 对比;接 OpenLineage 暴露列级血缘。
任何综合项目都必须包含:
- 需求说明与数据契约(源 schema、目标 schema、SLA、演进策略);
- 兼容性矩阵 + 演进 SOP 文档;
- 核心代码(Iceberg DDL + Flink Job + Spark 批读 + 演进脚本);
- 边界测试(演进失败回滚、批流结果不一致、并发演进、Time Travel 失败、Schema Registry 不匹配);
- 可观测(Iceberg snapshot 元数据 + Flink checkpoint + Spark metrics);
- README(设计取舍、表格式选型、演进门禁规则);
- 复盘记录(模拟 1 次不兼容演进失败 + 1 次批流不一致,写
retrospective.md); - 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.md | 4×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 职责
- 用 Avro + Confluent Schema Registry 定义兼容性矩阵(
BACKWARD/FULL/BACKWARD_TRANSITIVE),演进操作(加列 / 删列 / 改类型 / 重命名)必须落在 backward compat 范围内,下游 Avro 反序列化永不报错。 - 用 Delta Lake
delta.columnMapping.mode='name'+ Icebergcolumn_id做列追踪,重命名不破坏下游;演进前用 Time Travel 读历史验证,演进失败用CALL system.rollback_to_snapshot一键回退。 - 用 CDC(Debezium 抓 MySQL binlog → Kafka topic
cdc.orders→ Flink CDC Source → Icebergdw.orders)把 OLTP 变更流式接入 Lakehouse;批读用 Spark / Trino / Flink SQL 三引擎,跑同一查询对比批流结果必须 0 差异。
4 交付物
- 一份 Avro 契约 + Schema Registry 兼容性配置(
BACKWARD/FULL/BACKWARD_TRANSITIVE),含加列 / 改类型 / 重命名 3 步演进脚本,每次演进后老 Job 读旧数据 OK。 - 一份 Delta Lake Schema 演进演练(开
columnMapping.mode='name'→ 重命名列 →DESCRIBE HISTORY查演进历史),含 IcebergADD COLUMN+RENAME COLUMN对比表,配format-version='2'。 - 一份 CDC 接入管道(Debezium MySQL connector → Kafka
cdc.orders→ FlinkMySqlSource+ IcebergSink),含 binlog row 模式配置 +enableCheckpointing(60_000, EXACTLY_ONCE)。 - 一份批流一致性对比报告(5+ 次同查询 Spark vs Flink SQL,结果 0 差异),含演进 SOP 文档(演进前兼容性检查 → 演进 → 灰度 → 全量 → 回滚预案)。
3 指标
- Schema 演进 backward compat 通过率 100%(所有演进操作都通过 Schema Registry 兼容性校验,老消费者零修改)。
- 批流同查询 0 差异(5 次查询 Spark vs Flink SQL 结果完全一致,含行数与聚合值)。
- 演进回滚时长 < 1 min(
CALL system.rollback_to_snapshot执行到表 schema 与数据恢复)。
8. 推荐开源资料
| 阶段 | 角色 | 资料 | 链接 | 用法 |
|---|---|---|---|---|
| 全部 | 规范 | Apache Iceberg Spec v2 | https://iceberg.apache.org/spec/ | 必读 schema section |
| 1~2 | 演进 | Iceberg Schema Evolution 文档 | https://iceberg.apache.org/docs/latest/evolution/ | 操作指南必读 |
| 3 | Hudi | Apache Hudi Schema Evolution | https://hudi.apache.org/docs/schema_evolution | Hudi 演进操作 |
| 4 | Delta | Delta Schema 文档 | https://docs.delta.io/latest/delta-schema.html | columnMapping 模式 |
| 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 Registry | https://docs.confluent.io/platform/current/schema-registry/index.html | 兼容性与演进门禁 |
| 全部 | 工具 | Apache Avro | https://avro.apache.org/docs/current/ | Schema 定义 |
| 全部 | 工具 | Protocol Buffers | https://protobuf.dev/ | 跨语言 Schema |
| 6 | 引擎 | Spark Structured Streaming | https://spark.apache.org/docs/latest/structured-streaming-programming-guide.html | Spark 流读 Iceberg |
| 全部 | 对比 | Apache Paimon | https://paimon.apache.org/ | 流式 Lakehouse 新选型 |
| 全部 | 哲学 | Kleppmann《DDIA》Ch.4 | https://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 v2 | Iceberg Schema Evolution | 必读 schema section |
| Apache Hudi RFC-29 | Hudi Schema Evolution | 选读 |
| Delta Lake Protocol | Delta 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 Spec | Schema Evolution | 必读 |
| Apache Iceberg Schema Evolution | 操作指南 | 必读 |
| Apache Hudi Schema Evolution | Hudi 操作 | 选读 |
| Delta Schema | Delta 操作 | 选读 |
| Apache Paimon | 流式 Lakehouse | 选读 |
9.3.4 核心人物
| 人物 | 影响 | 材料 |
|---|---|---|
| Ryan Blue | Iceberg 起源 | Netflix 公开演讲 |
| Vinoth Chandar | Hudi 起源 | Uber 公开演讲 |
| Michael Armbrust | Delta 起源 | 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 | 可以,但下游要兼容 |
| 7 | Iceberg / 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 Iceberg | 1.x | Apache | GA | Apache-2.0 |
| Apache Hudi | 0.14+ | Apache | GA | Apache-2.0 |
| Delta Lake | 2.x / 3.x | Linux Foundation | GA | Apache-2.0 |
| Apache Paimon | 0.x | Apache | 活跃 | Apache-2.0 |
| Apache Flink | 1.18+ | Apache | GA | Apache-2.0 |
Scope
Iceberg/Hudi/Delta/Paimon 是表格式;它们解决 Schema 演进与 ACID;不替代计算引擎与存储。
Structure
- Iceberg:
column_id标识列,支持 schema evolution by ID。 - Hudi:
SchemaProvider维护 schema,支持 inline / external。 - Delta:
columnMapping模式,重命名列安全。 - Paimon:流式 Lakehouse,schema evolution 与流写入深度集成。
Ecosystem
- 表格式:Iceberg、Hudi、Delta、Paimon。
- 计算引擎:Spark、Trino、Flink、Presto、Hive。
- Schema 注册:Schema Registry、Avro、Protobuf。
- 演进测试:Schema Compatibility Checker。
Depth Tiers
| 层级 | 能力 | 标准 |
|---|---|---|
| L0 | 知道存在 | 知道 Schema Evolution 与 Lakehouse |
| L1 | 看得懂示例 | 能读懂 Iceberg schema 操作 |
| L2 | 能正确调用 | 能用 SQL/API 演进 schema |
| L3 | 能解释与排错 | 能定位演进失败、批流不一致 |
| L4 | 能设计与扩展 | 能为新业务设计 schema 演进规范与批流一体架构 |
本子主题目标:L3。
Source
- Apache Iceberg Spec:引用快照 2026-07-30。
- Apache Iceberg Schema Evolution:引用快照 2026-07-30。
- Apache Hudi Schema Evolution:引用快照 2026-07-30。
- Delta Schema:引用快照 2026-07-30。
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 写 Iceberg + Spark 批读对比
// 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. 常见误区
- 加列后下游不重新读,新列缺失;
- 改类型用不兼容类型(int → string);
- 列重命名直接改 Parquet 文件名,破坏下游;
- 分区演进后旧分区失效;
- 批流结果差异当成”流处理特性”,不验证;
- 不写演进测试就直接生产演进;
- Iceberg column mapping 没开就重命名;
- Hudi schema provider 选错(inline vs external);
- Delta schemaTracking 没启用就演进;
- 演进后没 Time Travel 验证;
- 演进没兼容矩阵文档;
- 用裸 Parquet 做流批一体。
11. 所有知识点分类
- 编程语言
- 数据结构与算法
- 计算机基础
- 工程技术
- Web 与后端
- 前端与客户端
- 数据与人工智能
- 项目与职业能力
- 安全与可靠性
本计划归属:数据与人工智能 主 + 工程技术 辅。