数据工程基础:从 ETL 到数仓与流处理
0. 元信息
- 主题路径:
docs/topics/data-engineering-basics/README.md - 主分类:数据与人工智能
- 辅助分类:工程技术
- 适合对象:会写 Python/SQL、用过 Docker 与 Linux 命令行,希望从「会用数据库」过渡到「能设计并运维一套端到端数据流水线」的工程师
- 建议周期:6~10 周(每周 10~14 小时,至少一半时间用于跑批/读计划/做故障演练)
- 前置知识:
db-and-sql(关系模型、SQL、查询计划、事务隔离)、lang-python(脚本与类型)、docker-basics(镜像与 Compose);熟悉 Linux 终端与 Git - 最终目标:拿到一份业务需求(订单/点击流/日志/IoT 任一),能在 2 周内交付一条可重跑、可观测、可治理的数据流水线——包含分层数仓、批流两种执行模式、Schema 演进与数据质量门禁
1. 学习路线
ETL/ELT 编排与批处理(Airflow / Prefect / Dagster)
→ 数仓建模与 Lakehouse(分层架构 / 维度建模 / Iceberg / Delta Lake / Hudi)
→ 流处理与实时管道(Kafka / Pulsar / Flink / Spark Streaming)
→ 数据治理、质量与血缘(dbt / Great Expectations / Soda / OpenLineage)
→ Schema 演进、版本化与批流一体
→ 综合项目:端到端电商/IoT 数据流水线
按”先批后流、再治理、最后统一”的顺序学。每张子主题对应一段流水线环节,最后用综合项目把链路串起来。
2. 阶段周数分配
| 阶段 | 6 周方案 | 10 周方案 | 备注 |
|---|---|---|---|
| 1. ETL/ELT 编排与批处理 | 1 | 1.5 | 必备:DAG、依赖、调度、重跑 |
| 2. 数仓建模与 Lakehouse | 1 | 2 | 分层架构 + 维度建模 + 表格式 |
| 3. 流处理与实时管道 | 1.5 | 2.5 | Kafka + Flink/Spark 实战 |
| 4. 数据治理、质量与血缘 | 1 | 2 | dbt + GE + OpenLineage |
| 5. Schema 演进与批流一体 | 1 | 1.5 | Iceberg/Hudi 时间旅行 + 统一存储 |
| 6. 综合项目与故障演练 | 0.5 | 0.5 | 端到端流水线 + 复盘 |
6 周方案只跑通最小可用路径;10 周方案有 2 周用来重做实验与故障演练。综合项目占最后 1 周。
3. 九阶段表
| 阶段 | 核心知识 | 实践产出 | 可观察学会标准 |
|---|---|---|---|
| 1. ETL/ELT 编排 | DAG、Task、Operator、Scheduler、Backfill、Idempotency | 一条可重跑的批处理 DAG(5 任务) | 能解释 DAG 的拓扑序与失败重试边界 |
| 2. 数仓分层架构 | ODS / DWD / DWS / ADS / DIM、命名规范、分区 | 一套完整分层 schema(PostgreSQL + Iceberg) | 能讲清每一层的责任与边界 |
| 3. 维度建模 | 事实表、维度表、星型/雪花、缓慢变化维 SCD | 一份事实表 + 3 张维度表的 SQL 与 Iceberg | 能区分事实与维度并选 SCD 类型 |
| 4. Lakehouse 与表格式 | Parquet / Avro / ORC、Iceberg / Delta Lake / Hudi、ACID、Time Travel | 一份 Iceberg 表的快照与回滚演练 | 能用 Time Travel 恢复错误写入 |
| 5. 流处理基础 | Kafka topic / partition / offset、exactly-once、Watermark、Checkpoint | 一条 Kafka → Flink → Iceberg 的管道 | 能解释端到端一致性与迟到数据 |
| 6. 流处理进阶 | Window、State、Process Function、回填、双流 Join | 一份 7 天滑动窗口 UV 报表 | 能复现乱序与迟到对结果的影响 |
| 7. 数据治理 | 元数据、血缘、目录、访问控制、契约 | 一份 OpenLineage 事件流 + 数据契约 YAML | 能用 SQL 查血缘并解释跨表依赖 |
| 8. 数据质量 | dbt tests、Great Expectations、Soda、期望/异常/SLA | 5 条质量门禁规则 + 失败告警 | 能区分语法/语义/分布三类异常 |
| 9. Schema 演进与批流一体 | 向后兼容、版本化、Lakehouse 统一存储、批流 SQL 一致性 | 一份 Schema 演进迁移演练 + 批流同一结果对比 | 能用 Iceberg/Hudi 在同一份表上同时跑批流查询 |
9.x 关键代码(冰山下示例):一份 Iceberg 表的最小写入 + Time Travel 回滚片段——用于把”分层 + 表格式”两件套串起来。
-- 9.x.1 Iceberg 最小写入 + Time Travel 回滚(Spark SQL)
-- 阶段 1:建表
CREATE TABLE dw.orders (
order_id BIGINT,
user_id BIGINT,
amount DECIMAL(10, 2),
order_date DATE
) USING iceberg
PARTITIONED BY (days(order_date));
-- 阶段 2:写入两批
INSERT INTO dw.orders VALUES (1, 100, 99.50, DATE '2026-01-15');
INSERT INTO dw.orders VALUES (2, 101, 199.00, DATE '2026-01-16');
-- 阶段 3:故意写错
UPDATE dw.orders SET amount = -1 WHERE order_id = 1;
-- 阶段 4:回滚到上一个快照
CALL dw.system.rollback_to_snapshot('orders', <snapshot_id>);
-- 阶段 5:验证
SELECT * FROM dw.orders WHERE order_id = 1; -- amount = 99.50
此示例呼应阶段 2(分层)、阶段 4(Lakehouse)、阶段 8(Schema 演进)三段决策;真实综合项目应在完整 dbt + Flink + OpenLineage 链路上扩展。
关键陷阱:把”搭好 Airflow”当数据工程;只做 ELT 不做数据契约;流处理只跑通 happy path 不测乱序与迟到;数据质量只在 dashboard 里看,不挂门禁;Schema 演进靠 ALTER TABLE,不走表格式版本化。
4. 第一周任务
Day 1 约定:本计划以 Linux + Docker + 8G 内存的机器为最小实验环境。所有流水线都在 Docker Compose 内启动,所有代码、SQL、配置都要进 Git。所有任务必须可重跑:清空状态后再跑一次产出相同结果。
| 日 | 任务 | 当天交付 |
|---|---|---|
| Day 1 | 装 Docker Compose、Airflow、MinIO;跑通最小 Airflow(1 个 DAG、2 个 Task) | docker-compose.yml + 第一张 DAG |
| Day 2 | 写一个 ETL DAG:从 CSV → Parquet → PostgreSQL;用 Airflow 调度 + 重跑 | DAG 代码 + 3 个 Task |
| Day 3 | 把 DAG 升级为 ELT:先入仓(Parquet → Iceberg on MinIO),再用 dbt 转换 | Iceberg 表 + dbt 项目 |
| Day 4 | 搭分层架构:ODS → DWD → DWS;写一份命名规范与分区策略 | 分层 SQL + 规范文档 |
| Day 5 | 用 Iceberg 做一次 Time Travel:故意写错,回滚到上一个快照 | 截图 + 前后行数对比 |
| Day 6 | 用 dbt 加 5 条数据测试(unique / not_null / accepted_values / relationships / freshness) | dbt 项目 + 测试报告 |
| Day 7 | 项目步骤 A:跑通订单 ETL 最小流水线(CSV → Iceberg ODS → dbt DWD → DWS → ADS);步骤 B:补齐 4 类边界(源文件缺失、Schema 漂移、分区重复、dbt 测试失败) | 流水线 + 4 类故障演练报告 |
Day 7 步骤 A 跑通业务流;步骤 B 每类故障先在测试环境复现,再讲清防护手段。
5. 阶段通用验收
- 不看答案独立重写一条 5 任务的 ETL DAG + 一份 Iceberg 表的 Time Travel 演练;
- 用自己的话解释”为什么分层""为什么流处理需要 Watermark""为什么 Schema 演进要走表格式”;
- 画一张图:数据流水线架构图 + 血缘图 + 一致性状态机图(至少三件套之一);
- 测试空数据、迟到数据、乱序数据、Schema 漂移、节点下线 5 类边界;
- 准备至少 3 组自定义数据画像并贴出实际产出与耗时;
- 记录端到端 P50 / P99 延迟、吞吐、checkpoint 时长、数据质量失败率 4 个核心指标;
- 能修改已有流水线(加一层、加一个数据测试、换一份表格式)并用同样指标验证。
6. 最终验收
-
独立设计并交付一条端到端数据流水线(订单/点击流/日志/IoT 任选),覆盖分层数仓、批流两种模式、Schema 演进、数据质量门禁;
-
至少完成 30 个实战项(分布建议):
子主题 题目数量 难度 平台建议 ETL/ELT 编排 6 Easy / Medium Astronomer Registry / Dagster examples 数仓建模与 Lakehouse 6 Medium Iceberg/Delta Lake 官方 Examples 流处理与实时管道 6 Medium / Hard Flink Forward / Kafka Summit 案例 数据治理与质量 6 Medium dbt + Great Expectations 官方 Examples Schema 演进与批流一体 6 Medium / Hard Iceberg Spec / Hudi RFC 约束:至少 18 项达到 Medium,至少 6 项达到 Hard;每项必须留可运行的代码与说明文档。
-
完成 1 个综合项目:端到端电商数据流水线,含分层数仓 + 实时管道 + 数据治理 + 故障复盘;
-
能用 15 分钟讲清整套数据流水线的分层、表格式、流处理语义、Schema 演进与数据质量门禁。
7. 综合项目
首选:电商端到端数据流水线(必做:分层数仓 + Kafka 实时管道 + dbt 转换 + Iceberg 表 + 数据测试)。
- 输入:订单 CSV + Kafka 实时点击流;
- 输出:必输出(1)分层数仓 schema:ODS(原始)→ DWD(明细)→ DWS(汇总)→ ADS(应用);(2)批流两种模式:DAG(Airflow)+ 流(Flink/Spark Streaming);(3)Iceberg 表 + Time Travel 演示;(4)dbt 数据测试 10 条;(5)OpenLineage 血缘;
- 关键指标:批处理 P95 延迟 < 30min、流处理端到端 P95 延迟 < 30s、数据测试失败 0%、Schema 演进 0 阻塞;
- 进阶可选:接入 Soda 做分布异常检测、Monte Carlo 做数据可观测、用 StarRocks/ClickHouse 做 OLAP 查询。
备选:IoT 设备遥测数据流水线(边缘采集 → Kafka → Flink → Iceberg → 告警)。
备选:日志分析平台(Filebeat → Kafka → Flink → Elasticsearch / Doris)。
备选:对一套遗留 cron 脚本做 Airflow 化重构并验证可观测与重跑。
任何综合项目都必须包含:
- 需求说明与数据契约(源 schema、目标 schema、SLA);
- 分层架构图与命名规范;
- 核心代码(Airflow DAG + Flink Job + dbt models + GE/Soda 检查);
- 边界测试(源缺失、Schema 漂移、迟到数据、节点下线、回填);
- 血缘与可观测(OpenLineage 事件 + Prometheus 指标);
- README(设计取舍、已知限制、下一步);
- 复盘记录(模拟 1 次 P0 故障,写 retrospective.md)。
8. 推荐开源资料
| 阶段 | 角色 | 资料 | 链接 | 用法 |
|---|---|---|---|---|
| 全部 | 主线书 | Martin Kleppmann《Designing Data-Intensive Applications》 | https://dataintensive.net/ | 数据工程「圣经」,必读第 5~12 章 |
| 全部 | 主线书 | Joe Reis & Matt Housley《Fundamentals of Data Engineering》 | https://www.oreilly.com/library/view/fundamentals-of-data/9781098108298/ | 数据工程生命周期全景 |
| 1 | 编排 | Apache Airflow 官方文档 | https://airflow.apache.org/docs/ | DAG、Operator、Executor |
| 1 | 编排 | Dagster 官方文档 | https://docs.dagster.io/ | Software-defined assets |
| 1 | 编排 | Prefect 官方文档 | https://docs.prefect.io/ | 动态工作流 |
| 2 | 数仓 | 《The Data Warehouse Toolkit》Ralph Kimball | https://www.kimballgroup.com/ | 维度建模必读 |
| 2 | Lakehouse | Apache Iceberg 官方文档 | https://iceberg.apache.org/docs/latest/ | 表格式规范与 API |
| 2 | Lakehouse | Delta Lake 官方文档 | https://docs.delta.io/ | 与 Spark 深度集成 |
| 2 | Lakehouse | Apache Hudi 官方文档 | https://hudi.apache.org/ | 记录级索引与 CDC |
| 3 | 流处理 | Apache Kafka 官方文档 | https://kafka.apache.org/documentation/ | Producer/Consumer/Connect/Streams |
| 3 | 流处理 | Apache Flink 官方文档 | https://nightlies.apache.org/flink/flink-docs-stable/ | DataStream / Table API |
| 3 | 流处理 | 《Streaming Systems》Tyler Akidau 等 | https://www.oreilly.com/library/view/streaming-systems/9781491983867/ | 流处理概念必读 |
| 4 | 治理 | dbt 官方文档 | https://docs.getdbt.com/ | 转换 + 测试 + 文档 |
| 4 | 治理 | Great Expectations 官方文档 | https://docs.greatexpectations.io/ | 期望即代码 |
| 4 | 治理 | OpenLineage 规范 | https://openlineage.io/ | 血缘标准 |
| 5 | 可观测 | Monte Carlo / Datafold 博客 | https://www.montecarlodata.com/blog/ | 数据可观测理念 |
| 5 | 可观测 | Apache Doris / StarRocks 文档 | https://doris.apache.org/ | OLAP 查询层 |
默认使用顺序:先读 Reis《Fundamentals of Data Engineering》第 1~6 章建立全景 → 用 Airflow 跑通第一批 DAG → 读 Iceberg 文档做 Lakehouse → 读 Flink 文档做流处理 → 用 dbt + Great Expectations 加质量门禁 → 用 OpenLineage 接血缘 → 写复盘到
notes/retrospective.md。
9. 学习资料汇聚(v0.3 自包含)
本节由本计划生成。链接指向原始材料或作者公开内容。规范会演进,记录时务必写明版本与日期。
9.1 背景与动机
数据工程的诞生是为了回答”数据从产生到产生价值中间发生了什么”。2003 年 Google 发表 GFS / MapReduce / Bigtable 三篇论文,把”分布式存储 + 计算 + 结构化”做成工业可用的三层。2006 年 Hadoop 开源,把这套思想带入主流。2010 年代初 Lambda 架构(Nathan Marz)提出批流双路,把”实时性”与”准确性”分开处理;2014 年 Kappa 架构(Jay Kreps)提出用单一日志流统一两者。2017 年 Delta Lake(Databricks)、2018 年 Iceberg(Netflix 开源)、2019 年 Hudi(Uber 开源)让 Lakehouse 成为现实——同一份存储既支持批 SQL 又支持流式更新。Martin Kleppmann 2017 年 Designing Data-Intensive Applications 把这一切讲透;Joe Reis 与 Matt Housley 2022 年 Fundamentals of Data Engineering 把数据工程生命周期定义为”生成 → 摄取 → 存储 → 转换 → 服务”五段。今天每一个数据团队都把”分层 → 表格式 → 流处理 → 治理 → 可观测”五个字拆成具体技术选型。
一句话总结:数据工程不是搭一条 Airflow,是把数据从产生到消费的全链路上每一段都能被调度、被治理、被观测。
9.2 概念地图
flowchart LR
Source[数据源 DB/日志/IoT/API] --> Ingest[摄取 Airflow/Kafka/CDC]
Ingest --> ODS[ODS 原始层]
ODS --> DWD[DWD 明细层]
DWD --> DWS[DWS 汇总层]
DWS --> ADS[ADS 应用层]
ADS --> Serve[服务 BI/API/ML]
Ingest --> Stream[流处理 Kafka/Flink/Spark]
Stream --> Lake[Lakehouse Iceberg/Delta/Hudi]
Lake --> SQL[OLAP 查询 StarRocks/ClickHouse/Doris]
Meta[元数据] --> Catalog[数据目录 Hive/Glue/Polaris]
Quality[数据质量 dbt/GE/Soda] --> Gate[质量门禁]
Lineage[血缘 OpenLineage/Marquez] --> Gate
Lake --> Gate
Serve --> Gate
关系说明:摄取决定分层入口;分层决定每层责任与命名;流处理与 Lakehouse 让批流共享存储;治理与血缘贯穿全链路;可观测提供反馈环。每一层都对应一个失败模式,必须在设计里写明应对。
9.3 基础知识讲解
9.3.1 经典论文 / 规范
| 资料 | 影响 | 建议读法 |
|---|---|---|
| Jeffrey Dean & Sanjay Ghemawat, MapReduce: Simplified Data Processing on Large Clusters(OSDI 2004) | 分布式计算范式起点 | 重点看容错与数据本地化 |
| Ghemawat et al., The Google File System(SOSP 2003) | 分布式文件系统 | 理解分块、副本、Master |
| Chang et al., Bigtable: A Distributed Storage System for Structured Data(OSDI 2006) | 列存宽表 | 看 SSTable 与 MemTable |
| Nathan Marz, Lambda Architecture(2011) | 批流双路 | 与 Kappa 对照读 |
| Jay Kreps, Questioning the Lambda Architecture(2014) | Kappa 单一日志流 | 看”Why Kappa” |
| Kleppmann, Designing Data-Intensive Applications(O’Reilly 2017) | 数据系统原理圣经 | 必读第 5~12 章 |
| Reis & Housley, Fundamentals of Data Engineering(O’Reilly 2022) | 数据工程生命周期 | 第 1~6 章建立全景 |
| Akidau et al., The Dataflow Model(VLDB 2015) | 流处理语义 | 看 Watermark、Trigger、Accumulation |
| Armbrust et al., Delta Lake: High-Performance ACID Table Storage(VLDB 2020) | Lakehouse 起点 | 对照 Iceberg/Hudi 看 |
| Schmidt et al., Apache Iceberg: A Table Format for Hadoop(2020) | 表格式规范 | 看快照与清单文件 |
| Kiran et al., Hudi: A Storage System for Modern Analytics(2021) | 记录级索引 | 看 Copy-on-Read vs Merge-on-Read |
| OpenLineage Specification 1.x | 血缘事件标准 | 与 Marquez 对照看 |
9.3.2 经典书籍
| 书 | 影响 | 用法 |
|---|---|---|
| Martin Kleppmann, Designing Data-Intensive Applications(O’Reilly 2017) | 数据系统原理 | 必读 5~12 章 |
| Joe Reis & Matt Housley, Fundamentals of Data Engineering(O’Reilly 2022) | 数据工程生命周期 | 第 1~6 章建立全景 |
| Ralph Kimball & Margy Ross, The Data Warehouse Toolkit(3rd ed., Wiley 2013) | 维度建模 | 必读 1~5 章 |
| Bill Inmon, Building the Data Warehouse(4th ed., Wiley 2005) | 第三范式建模 | 与 Kimball 对照 |
| Tyler Akidau et al., Streaming Systems(O’Reilly 2018) | 流处理三大问题(What/Where/When) | 必读 1~5 章 |
| Holden Karau et al., Spark: The Definitive Guide(O’Reilly 2018) | Spark 全景 | 选读 |
| Wes McKinney, Python for Data Analysis(3rd ed., O’Reilly 2022) | 数据处理 Python 生态 | 选读 |
| Richard Stevens, TCP/IP Illustrated(Addison-Wesley 2011) | 网络 | 网络类子主题参考 |
9.3.3 优秀博客
| 资料 | 特点 | 用法 |
|---|---|---|
| Apache Airflow Blog | 编排器演进与最佳实践 | 看 2023~2025 TaskFlow API |
| dbt Blog | 转换层理念与案例 | 看 dbt mesh 与契约 |
| Apache Iceberg Blog | 表格式规范解读 | 看 Hidden Partitioning |
| Apache Flink Blog | 流处理引擎演进 | 看 Checkpoint 与 State |
| Databricks Blog | Delta Lake 与 Lakehouse | 看 Photon 与 Unity Catalog |
| Uber Engineering | Hudi 起源与实践 | 看 CDC 与回填 |
| Netflix Tech Blog | Iceberg 起源与实践 | 看 Meta Catalog |
| Confluent Blog | Kafka 生态 | 看 Schema Registry 与 ksqlDB |
| AWS Big Data Blog | 托管数据服务案例 | 选 5 篇 Glue/Lake Formation |
| Martin Kleppmann 主页 | DDIA 与 CRDT | 选读 |
| Maxime Beauchemin | Airflow 起源 | 必读 “Functional Data Engineering” |
9.3.4 核心人物
| 人物 | 主要影响 | 建议追踪的材料 |
|---|---|---|
| Jeffrey Dean | MapReduce / GFS / Bigtable / Spanner | Google Research 主页 |
| Sanjay Ghemawat | MapReduce / GFS | 同上 |
| Nathan Marz | Lambda 架构 | 2011 博客 |
| Jay Kreps | Kappa 架构 / Kafka | Confluent 博客 |
| Martin Kleppmann | DDIA / CRDT | Cambridge 公开课 |
| Joe Reis / Matt Housley | Fundamentals of Data Engineering | 播客与博客 |
| Ralph Kimball | 维度建模 | Kimball Group |
| Tyler Akidau | Dataflow Model / Beam | Google 与 Apache 博客 |
| Reynold Xin | Spark / Iceberg 推动者 | 公开演讲 |
| Ryan Blue | Iceberg 起源 | Netflix 公开演讲 |
| Vinoth Chandar | Hudi 起源 | Uber 公开演讲 |
| Maxime Beauchemin | Airflow 起源 | Functional Data Engineering |
9.3.5 开发方法
| 方法 | 具体动作 | 何时用 |
|---|---|---|
| Schema-first | 先写数据契约(源 schema、目标 schema、SLA),再写代码 | 任何数据流水线第一步 |
| Idempotent task | 同一 Task 多次执行产出相同结果 | 任何 Airflow DAG Task |
| Partition by date | 按日期分区,避免重写历史 | 任何 ODS/DWD 表 |
| Backfill as first-class | 重跑历史是日常需求,不是异常 | 设计 DAG 时必考虑 |
| ELT over ETL | 入仓后再转换,利用仓内计算 | 数据量小、计算便宜时 |
| Lakehouse over warehouse + lake | 同一份存储支持批流 SQL | 中等规模团队首选 |
| Table format mandatory | 表格式(Iceberg/Delta/Hudi)替代裸 Parquet | 任何需要 ACID 的场景 |
| Quality as gate | dbt/GE 测试失败阻塞下游 | 上线前必加 |
| Lineage everywhere | 每个 Task 输出 OpenLineage 事件 | 任何生产流水线 |
| Observability triplet | Metrics + Logs + Traces 三件套 | 数据流水线也适用 |
| Disaster rehearsal | 故意写错、看能不能恢复 | 季度必做 |
9.3.6 重点训练材料
- DDIA 案例:第 5~9 章每章末尾的”案例研究”是顶级题源(批处理用 MapReduce、流处理用 Flink、事务用 Spanner)。
- Astronomer Registry:Airflow 官方 DAG 模板库,覆盖 S3 → Snowflake、Salesforce → BigQuery 等。
- Apache Iceberg Examples:官方仓库的 Spark/Flink/Hive 集成示例。
- dbt 项目模板:Jaffle Shop、Netflix dbt 示例、GitHub Archive 项目。
- Flink Forward Talks:每年大会的视频,覆盖 Exactly-Once、State、Savepoint 等。
- Reis《Fundamentals of Data Engineering》案例:第 7~12 章每章末尾的”案例研究”覆盖 Uber / Netflix / Airbnb 数据流水线。
9.4 经典问题与经典案例
| # | 问题 | 为什么会重要 | 最简答案或图示 |
|---|---|---|---|
| 1 | Airflow Task 怎么写才 idempotent | 重跑时不重复写数据 | 按日期分区 + 删除目标分区后写入 |
| 2 | ELT 与 ETL 怎么选 | 影响存储与计算分离 | 数据便宜时 ELT、计算重时 ETL |
| 3 | 维度建模星型 vs 雪花 | 决定查询性能与可维护性 | 星型优先、雪花慎用 |
| 4 | SCD Type 1/2/3/4/6 怎么选 | 历史追溯需求 | 大部分场景 Type 2 |
| 5 | Iceberg / Delta / Hudi 怎么选 | 表格式决定 Lakehouse 能力 | Spark 生态选 Delta、Hive 生态选 Iceberg、CDC 强选 Hudi |
| 6 | Kafka 怎么保证 exactly-once | 流处理一致性 | 幂等 Producer + 事务 + 幂等 Consumer |
| 7 | Flink Watermark 怎么设 | 乱序与迟到数据 | 用 bounded out-of-orderness 策略 |
| 8 | 流处理 State 后端选哪个 | RocksDB vs HashMapState | 大状态选 RocksDB |
| 9 | dbt incremental model 怎么写 | 避免每次全量重算 | 用 is_incremental() + unique_key |
| 10 | OpenLineage 怎么接 | 血缘可视化 | 用 Marquez + Airflow OpenLineage provider |
| 11 | 数据质量 5 类测试 | 防止脏数据进仓 | unique / not_null / accepted_values / relationships / freshness |
| 12 | Schema 演进怎么兼容 | 加列删列不阻塞 | 用表格式的 schema evolution + 兼容性矩阵 |
| 13 | 批流一体怎么落地 | 同一份存储两种模式 | Lakehouse + 流批 SQL 引擎(Spark/Flink/Trino) |
| 14 | 迟到数据怎么补 | 流处理常见痛点 | 用 Watermark + Allowed Lateness + Side Output |
| 15 | 数据契约怎么写 | 防止上下游脱节 | 用 Protobuf / Avro + Schema Registry |
9.5 学习难点
概念难点
| 难点 | 为什么会卡 | 突破路径 |
|---|---|---|
| 批处理 vs 流处理语义 | 边界混淆 | 用 Lambda/Kappa 对照实验 |
| Exactly-Once 处理 | 概念阶梯 | 用 Flink Two-Phase Commit Sink 做对照 |
| 维度建模与关系建模 | 思维切换 | 用 Kimball 案例 vs 3NF 对照 |
| 表格式的快照机制 | 概念抽象 | 读 Iceberg manifest list + Parquet 文件 |
| 数据契约 vs 数据 schema | 边界模糊 | 用 Protobuf schema + SLA 文档 |
思维难点
| 难点 | 为什么会卡 | 突破路径 |
|---|---|---|
| 从需求到分层 schema | 不会起步 | 按 ODS/DWD/DWS/ADS 四层映射 |
| 把失败映射到组件 | 慢、错、丢都是分层结论 | 拆成摄取/转换/存储/服务四段 |
| 批流结果不一致 | 容易忽略 | 跑同一查询对比批流输出 |
| 数据质量门禁阈值 | 拍脑袋 | 用历史分布 + 业务 SLA 设 |
工程难点
| 难点 | 为什么会卡 | 突破路径 |
|---|---|---|
| Airflow Task 失败重试 | 容易漏写 | 用 retry + exponential_backoff |
| Iceberg 并发写入冲突 | Optimistic Concurrency | 配 commit.retry.num-retries |
| Flink State 太大 OOM | RocksDB 配置 | 调 write buffer + block cache |
| dbt 全量重算慢 | 增量策略错 | 用 is_incremental() + 物化策略 |
| OpenLineage 事件丢失 | 网络问题 | 加 retry + DLQ |
9.6 技术标准与接口
9.6.1 Entity
| 名称 | 版本 / 文档 | 发布组织 | 状态 | 许可证 / 可访问性 |
|---|---|---|---|---|
| Apache Airflow | 2.x | Apache | GA | Apache-2.0 |
| Apache Kafka | 3.x | Apache | GA | Apache-2.0 |
| Apache Flink | 1.18+ | Apache | GA | Apache-2.0 |
| Apache Spark | 3.5+ | Apache | GA | Apache-2.0 |
| Apache Iceberg | 1.x | Apache | GA | Apache-2.0 |
| Delta Lake | 2.x / 3.x | Linux Foundation | GA | Apache-2.0 |
| Apache Hudi | 0.14+ | Apache | GA | Apache-2.0 |
| dbt Core / Cloud | 1.x / 最新 | dbt Labs | 活跃 | Apache-2.0 / 商业 |
| Great Expectations | 0.18+ | Great Expectations | 活跃 | Apache-2.0 |
| Soda Core | 3.x | Soda Data | 活跃 | Apache-2.0 |
| OpenLineage | 1.x | OpenLineage | 活跃 | Apache-2.0 |
| Marquez | 1.x | OpenLineage | 活跃 | Apache-2.0 |
| Parquet | 2.x | Apache | 成熟 | Apache-2.0 |
| Avro | 1.11+ | Apache | 成熟 | Apache-2.0 |
| Arrow / Flight | 14+ | Apache | 活跃 | Apache-2.0 |
| StarRocks / ClickHouse / Doris | 最新 | 各社区 | 活跃 | Apache-2.0 / 商业 |
9.6.2 Scope
- Airflow / Prefect / Dagster:编排层,不替代存储与计算。
- Iceberg / Delta / Hudi:表格式,不替代文件系统与计算引擎。
- Kafka / Pulsar:消息层,不替代流处理引擎。
- Flink / Spark Streaming:流处理引擎,不替代存储。
- dbt / Great Expectations / Soda:转换与质量,不替代编排。
- OpenLineage / Marquez:血缘标准与实现,不替代元数据目录。
9.6.3 Structure
- Airflow DAG:
DAG → Task → Operator → Hook → Sensor,Task 间用 XCom 传值。 - Kafka topic:
partition × replication_factor,offset 由 consumer group 管理。 - Flink Job:
Source → Transform → Sink,State 后端用 RocksDB/HashMap。 - Iceberg:
metadata.json + manifest list + manifest file + Parquet data file,支持 hidden partitioning。 - dbt model:
sources → staging → intermediate → marts,测试用 schema.yml。 - OpenLineage event:
run / job / dataset / facet四元组,Facet 描述属性。 - 数据契约:源 schema + 目标 schema + SLA + owner + 版本。
9.6.4 Ecosystem
- 编排:Airflow、Prefect、Dagster、Argo Workflows、Temporal。
- 摄取:Kafka Connect、Flink CDC、Debezium、Airbyte、Singer。
- 存储:HDFS、S3、GCS、MinIO、Azure Blob。
- Lakehouse:Iceberg、Delta Lake、Hudi、Apache Paimon。
- 流处理:Flink、Spark Streaming、Beam、ksqlDB。
- OLAP:StarRocks、ClickHouse、Doris、Trino、Presto、DuckDB。
- 转换:dbt、SQLMesh、Coalesce。
- 质量:Great Expectations、Soda、Monte Carlo、Datafold、Bigeye。
- 血缘:OpenLineage、Marquez、DataHub、Apache Atlas、Unity Catalog。
- 调度监控:Airflow UI、Prometheus、Grafana、Dagster UI。
9.6.5 Depth Tiers
| 层级 | 能力 | 数据工程基础的可观察标准 |
|---|---|---|
| L0 | 知道存在 | 知道 ETL / 数仓 / 流 / 治理各自解决什么问题 |
| L1 | 看得懂示例 | 能读懂 Airflow DAG、Iceberg 表、dbt model |
| L2 | 能正确调用 | 能用 Airflow 跑通 DAG、用 Iceberg 写入读取、用 dbt 转换测试 |
| L3 | 能解释与排错 | 能在 2 周内交付端到端流水线,能定位 DAG 失败、流处理乱序、数据质量失败 |
| L4 | 能设计与扩展 | 能为一个新业务选型分层、表格式、治理、可观测 |
本计划目标:L3。综合项目可触及局部 L4,但不作为 6~10 周的硬门槛。
9.6.6 Source
- Apache Airflow 官方文档:DAG、Operator、Executor。引用快照日期:2026-07-30。
- Apache Iceberg 官方文档:表格式规范。引用快照日期:2026-07-30。
- Apache Kafka 官方文档:Producer/Consumer/Streams。引用快照日期:2026-07-30。
- Apache Flink 官方文档:DataStream / Table API。引用快照日期:2026-07-30。
- dbt 官方文档:转换与测试。引用快照日期:2026-07-30。
- OpenLineage 规范:血缘标准。引用快照日期:2026-07-30。
- Reis《Fundamentals of Data Engineering》:数据工程生命周期。引用快照日期:2026-07-30。
10. 常见误区
- 把”搭好 Airflow”当数据工程;缺分层、缺质量门禁、缺血缘;
- 只做 ETL 不做 ELT;数据便宜时先入仓再转换;
- 流处理只跑通 happy path,不测乱序与迟到;
- 数据质量只在 dashboard 看,不挂门禁;
- Schema 演进靠 ALTER TABLE,不走表格式版本化;
- 维度建模用雪花型,雪花型难维护;
- Iceberg / Delta / Hudi 当裸 Parquet 用,不配 schema evolution;
- Kafka 用 at-most-once 当 exactly-once;
- Flink State 后端用 HashMap,处理大状态 OOM;
- dbt 全量重算不开 incremental;
- 数据契约只在文档里写,不入 Schema Registry;
- 血缘只画表级不看列级;
- 批流结果不一致当”小问题”,不验证;
- 流处理 checkpoint 配置不当,状态恢复丢数据;
- 数据测试只测 unique / not_null,不测分布异常;
- 综合项目只跑通 DAG 不做故障演练;
- 学完不复盘,下次面试依然答不上来。
11. 所有知识点分类(统一规则)
- 编程语言
- 数据结构与算法
- 计算机基础
- 工程技术
- Web 与后端
- 前端与客户端
- 数据与人工智能
- 项目与职业能力
- 安全与可靠性
本计划归属:数据与人工智能 主 + 工程技术 辅。