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

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

分类:数据与人工智能 · 路径:docs/topics/data-engineering-basics/README.md

#data-engineering#etl#warehouse#streaming#airflow#spark

用 6~10 周从 ETL 到能设计分层数仓与流式管道,落地数据治理与可观测

父主题

顶层主题

子主题(5)

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

0. 元信息

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 编排与批处理11.5必备:DAG、依赖、调度、重跑
2. 数仓建模与 Lakehouse12分层架构 + 维度建模 + 表格式
3. 流处理与实时管道1.52.5Kafka + Flink/Spark 实战
4. 数据治理、质量与血缘12dbt + GE + OpenLineage
5. Schema 演进与批流一体11.5Iceberg/Hudi 时间旅行 + 统一存储
6. 综合项目与故障演练0.50.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、期望/异常/SLA5 条质量门禁规则 + 失败告警能区分语法/语义/分布三类异常
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. 阶段通用验收

  1. 不看答案独立重写一条 5 任务的 ETL DAG + 一份 Iceberg 表的 Time Travel 演练;
  2. 用自己的话解释”为什么分层""为什么流处理需要 Watermark""为什么 Schema 演进要走表格式”;
  3. 画一张图:数据流水线架构图 + 血缘图 + 一致性状态机图(至少三件套之一);
  4. 测试空数据、迟到数据、乱序数据、Schema 漂移、节点下线 5 类边界;
  5. 准备至少 3 组自定义数据画像并贴出实际产出与耗时;
  6. 记录端到端 P50 / P99 延迟、吞吐、checkpoint 时长、数据质量失败率 4 个核心指标;
  7. 能修改已有流水线(加一层、加一个数据测试、换一份表格式)并用同样指标验证。

6. 最终验收

7. 综合项目

首选:电商端到端数据流水线(必做:分层数仓 + Kafka 实时管道 + dbt 转换 + Iceberg 表 + 数据测试)。

备选:IoT 设备遥测数据流水线(边缘采集 → Kafka → Flink → Iceberg → 告警)。
备选:日志分析平台(Filebeat → Kafka → Flink → Elasticsearch / Doris)。
备选:对一套遗留 cron 脚本做 Airflow 化重构并验证可观测与重跑。

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

  1. 需求说明与数据契约(源 schema、目标 schema、SLA);
  2. 分层架构图与命名规范;
  3. 核心代码(Airflow DAG + Flink Job + dbt models + GE/Soda 检查);
  4. 边界测试(源缺失、Schema 漂移、迟到数据、节点下线、回填);
  5. 血缘与可观测(OpenLineage 事件 + Prometheus 指标);
  6. README(设计取舍、已知限制、下一步);
  7. 复盘记录(模拟 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 Kimballhttps://www.kimballgroup.com/维度建模必读
2LakehouseApache Iceberg 官方文档https://iceberg.apache.org/docs/latest/表格式规范与 API
2LakehouseDelta Lake 官方文档https://docs.delta.io/与 Spark 深度集成
2LakehouseApache 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 BlogDelta Lake 与 Lakehouse看 Photon 与 Unity Catalog
Uber EngineeringHudi 起源与实践看 CDC 与回填
Netflix Tech BlogIceberg 起源与实践看 Meta Catalog
Confluent BlogKafka 生态看 Schema Registry 与 ksqlDB
AWS Big Data Blog托管数据服务案例选 5 篇 Glue/Lake Formation
Martin Kleppmann 主页DDIA 与 CRDT选读
Maxime BeaucheminAirflow 起源必读 “Functional Data Engineering”

9.3.4 核心人物

人物主要影响建议追踪的材料
Jeffrey DeanMapReduce / GFS / Bigtable / SpannerGoogle Research 主页
Sanjay GhemawatMapReduce / GFS同上
Nathan MarzLambda 架构2011 博客
Jay KrepsKappa 架构 / KafkaConfluent 博客
Martin KleppmannDDIA / CRDTCambridge 公开课
Joe Reis / Matt HousleyFundamentals of Data Engineering播客与博客
Ralph Kimball维度建模Kimball Group
Tyler AkidauDataflow Model / BeamGoogle 与 Apache 博客
Reynold XinSpark / Iceberg 推动者公开演讲
Ryan BlueIceberg 起源Netflix 公开演讲
Vinoth ChandarHudi 起源Uber 公开演讲
Maxime BeaucheminAirflow 起源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 gatedbt/GE 测试失败阻塞下游上线前必加
Lineage everywhere每个 Task 输出 OpenLineage 事件任何生产流水线
Observability tripletMetrics + Logs + Traces 三件套数据流水线也适用
Disaster rehearsal故意写错、看能不能恢复季度必做

9.3.6 重点训练材料

9.4 经典问题与经典案例

#问题为什么会重要最简答案或图示
1Airflow Task 怎么写才 idempotent重跑时不重复写数据按日期分区 + 删除目标分区后写入
2ELT 与 ETL 怎么选影响存储与计算分离数据便宜时 ELT、计算重时 ETL
3维度建模星型 vs 雪花决定查询性能与可维护性星型优先、雪花慎用
4SCD Type 1/2/3/4/6 怎么选历史追溯需求大部分场景 Type 2
5Iceberg / Delta / Hudi 怎么选表格式决定 Lakehouse 能力Spark 生态选 Delta、Hive 生态选 Iceberg、CDC 强选 Hudi
6Kafka 怎么保证 exactly-once流处理一致性幂等 Producer + 事务 + 幂等 Consumer
7Flink Watermark 怎么设乱序与迟到数据用 bounded out-of-orderness 策略
8流处理 State 后端选哪个RocksDB vs HashMapState大状态选 RocksDB
9dbt incremental model 怎么写避免每次全量重算is_incremental() + unique_key
10OpenLineage 怎么接血缘可视化用 Marquez + Airflow OpenLineage provider
11数据质量 5 类测试防止脏数据进仓unique / not_null / accepted_values / relationships / freshness
12Schema 演进怎么兼容加列删列不阻塞用表格式的 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 太大 OOMRocksDB 配置调 write buffer + block cache
dbt 全量重算慢增量策略错is_incremental() + 物化策略
OpenLineage 事件丢失网络问题加 retry + DLQ

9.6 技术标准与接口

9.6.1 Entity

名称版本 / 文档发布组织状态许可证 / 可访问性
Apache Airflow2.xApacheGAApache-2.0
Apache Kafka3.xApacheGAApache-2.0
Apache Flink1.18+ApacheGAApache-2.0
Apache Spark3.5+ApacheGAApache-2.0
Apache Iceberg1.xApacheGAApache-2.0
Delta Lake2.x / 3.xLinux FoundationGAApache-2.0
Apache Hudi0.14+ApacheGAApache-2.0
dbt Core / Cloud1.x / 最新dbt Labs活跃Apache-2.0 / 商业
Great Expectations0.18+Great Expectations活跃Apache-2.0
Soda Core3.xSoda Data活跃Apache-2.0
OpenLineage1.xOpenLineage活跃Apache-2.0
Marquez1.xOpenLineage活跃Apache-2.0
Parquet2.xApache成熟Apache-2.0
Avro1.11+Apache成熟Apache-2.0
Arrow / Flight14+Apache活跃Apache-2.0
StarRocks / ClickHouse / Doris最新各社区活跃Apache-2.0 / 商业

9.6.2 Scope

9.6.3 Structure

9.6.4 Ecosystem

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

10. 常见误区

11. 所有知识点分类(统一规则)

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

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


直接依赖(3)

查看知识图谱