第 3 课
← 返回系列列表

ETL 与数据转换

Data Engineering — 数据工程基础

批处理与流处理 ETL、dbt、Airflow 等工具链。

ETL 与数据转换

课程简介

批处理与流处理 ETL、dbt、Airflow 等工具链。

🎬 本课程视频:Data Engineering — 数据工程基础


ETL 与数据转换:批处理、流式处理与工具链

一、ETL 的演进

ETL(Extract-Transform-Load)是现代数据工程的核心范式。传统 ETL 很简单:从源系统抽取数据 → 在中间层做清洗、聚合、类型转换 → 加载到目标数仓。但今天,ETL 已经演变成为覆盖批处理、流式处理和数据转换的复杂体系。

更重要的是,现代数据工程正在经历从 ETL 到 ELT 的转变——先把数据加载到目标系统,再在目标系统内做转换。数据湖和数据仓库的计算能力越来越强,使得"先加载后转换"成为更高效的模式。但这并不意味着 ETL 过时了,而是我们在不同的场景下需要选择不同的范式。

二、批处理 ETL

2.1 传统 ETL 流程

Extract(抽取):从源系统读取数据。源可能包括关系型数据库(通过 JDBC/ODBC 连接),API 接口(REST 或 GraphQL),以及文件系统(CSV、JSON、Parquet 文件)。抽取分为全量抽取(第一次加载或全量刷新时使用)和增量抽取(只抽取自上次以来变化的数据)。

Transform(转换):这是 ETL 的核心环节。清洗操作包括处理缺失值(删除、均值填充、默认值填充)、标准化格式(统一日期格式、货币单位)、去重、数据质量校验。聚合操作包括按时间维度汇总(按天/周/月聚合销售额)、按业务维度汇总(按地区/品类聚合)。业务逻辑计算包括衍生字段(计算客户生命周期价值)、编码转换(类别特征独热编码)、数据脱敏(手机号/身份证号掩码)。

Load(加载):将转换后的数据写入目标系统。全量加载覆盖写入整个表,增量加载通过 upsert(存在则更新,不存在则插入)或分区覆盖来实现。

2.2 核心工具

dbt(data build tool):dbt 是当前最受欢迎的转换层工具。它的核心理念是"用 SQL 做转换,用 Git 做版本控制"。dbt 的核心优势包括:
- 版本控制:所有转换逻辑都在 Git 中管理,支持代码审查和回滚。
- 血缘追踪:dbt 自动生成从原始表到最终模型的完整数据血缘图。
- 测试框架:内置数据质量测试,可以定义每个字段的唯一性、非空、外键约束等。
- 文档生成:自动从 SQL 注释生成数据目录文档。

一个简单的 dbt 模型示例:

-- models/marts/finance/daily_revenue.sql
WITH orders AS (
    SELECT * FROM {{ ref('stg_orders') }}
)
SELECT
    DATE(order_created_at) AS order_date,
    COUNT(DISTINCT order_id) AS total_orders,
    SUM(total_amount) AS total_revenue
FROM orders
WHERE order_status != 'cancelled'
GROUP BY 1

Apache Spark:当数据量达到 TB 级以上时,dbt 可能不够用了。Spark 是分布式计算引擎,可以在数十台甚至数百台机器上并行处理数据。Spark 的核心抽象是 DataFrame——类似于一张表,但分布在多台机器上。Spark 的优势在于内存计算——数据在内存中处理而不是反复读写磁盘,这使得 Spark 比传统 MapReduce 快 10-100 倍。

三、流式 ETL

3.1 流式处理的核心概念

流式处理与批处理不同——数据以事件流的方式持续到达,而不是一批一批地到达。核心概念包括:
- 事件时间 vs 处理时间:事件时间是数据实际发生的时间,处理时间是系统处理数据的时间。对流处理来说,两者可能差异很大。
- 窗口:时间窗口(固定窗口、滑动窗口、会话窗口)定义了如何将无限流切分为有限的计算单元。
- Exactly-Once 语义:确保每条数据恰好被处理一次,不会丢失也不会重复。

3.2 核心工具

Apache Kafka:作为事件总线,Kafka 连接数据生产和消费两端。数据生产者(Producer)发布消息到 Topic,数据消费者(Consumer)订阅 Topic 并消费消息。Kafka 的日志结构使得消息可以被多个消费者独立消费,且消费者可以回溯历史消息(在保留期内)。Kafka 能够处理每秒百万级消息,延迟在毫秒级。

Apache Flink:目前最强大的流处理引擎。与 Spark Streaming(实际上是微批次处理)不同,Flink 是真正的逐事件流处理引擎。Flink 支持真正的 Exactly-Once 语义、事件时间处理、以及复杂的状态管理。Flink 的核心 API 包括 DataStream API(处理无限流)和 Table API / SQL(声明式处理)。

四、任务编排

无论是批还是流,生产环境都需要调度和编排。

Apache Airflow:最流行的 DAG 调度器。核心概念包括 DAG(有向无环图,定义任务依赖关系)、Task(单个执行单元)、Operator(任务类型,如 PythonOperator、BashOperator)、Sensor(等待外部条件满足的触发器)。Airflow 使用 Python 定义 DAG,调度器负责按依赖关系执行任务。

Dagster:新一代编排平台,内置数据资产和软件定义资产的概念。与 Airflow 相比,Dagster 更强调数据血缘、类型安全和测试性。

选择和编排工具时考虑团队已有的技术栈和运维能力。Airflow 适合大多数团队,Dagster 适合追求新一代数据工程体验的团队。

五、ETL vs ELT:什么时候选择哪个?

维度 ETL ELT
转换位置 中间转换服务器 目标数据仓库
适用场景 源数据质量差,需大量清洗 源数据结构较好,需灵活分析
数据量 适合中小规模 适合大规模
灵活性 低(转换后加载) 高(加载后可按需转换)
工具 dbt(在仓库中转换仍可用) dbt + Snowflake/BigQuery

六、总结

ETL 从传统的批处理流程演变为了批流一体的复杂体系。批处理用 dbt 和 Spark,流式用 Kafka 和 Flink,编排用 Airflow 或 Dagster。选择工具链时优先考虑团队已有的技术能力,而不是追求最新最热的工具。

七、批流一体的架构模式

现代数据处理正在从"批处理 vs 流处理"的二元对立转向"批流一体"——同一套代码既能跑在批模式也能跑在流模式。

Apache Flink 的批流一体设计:Flink 从一开始就是为流处理设计的,但将批处理视为"有界流"——数据有限、有明确起点和终点的流。这意味着同一套 Flink API 可以处理两种场景。

Apache Spark 的 Structured Streaming:将流数据视为"不断地追加数据到无限表",使用与 DataFrame 一致的 API。Spark 将流切分为微批次(Micro-batch),每个微批次是一个小型的批处理作业。延迟通常在秒级。

Kafka Streams:轻量级的流处理库,嵌入到应用进程中运行。不需要单独的集群,部署简单。适合与 Kafka 深度集成的场景。

八、数据转换的最佳实践

数据转换是 ETL 的核心环节,涉及模式映射、清洗、聚合、业务逻辑计算等。

模式映射的最佳实践:
1. 建立源-目标字段映射文档,包括字段名、数据类型、转换规则
2. 使用明确的类型转换,不要依赖隐式转换
3. 日期时间统一为 UTC 存储,展示时按用户时区转换
4. 枚举值使用代码表(Code Table)而非硬编码数值

数据清洗的常见操作:
1. NULL 处理:区分"真正的空值"和"数据未采集"两种情况
2. 去重:基于业务主键去重,记录重复次数用于监控
3. 标准化:统一单位、统一编码(UTF-8)、统一格式
4. 越界检测:数值超过合理范围时报警并标记
5. 格式校验:邮箱、手机号、身份证号按规则校验

九、数据管道测试策略

数据管道也需要测试,但测试方式与传统应用不同:

  1. 数据质量测试(Data Quality Test):用 Great Expectations 验证输出数据是否符合期望——行数、空值率、唯一性、分布特征
  2. 端到端测试(End-to-End Test):用已知的测试数据输入管道,验证输出是否正确
  3. 回滚测试(Replay Test):跑回历史数据,验证新版本和老版本输出一致
  4. 压力测试(Stress Test):使用大容量数据验证管道的吞吐能力和稳定性

数据管道的 CI/CD:将管道代码存储在 Git 中,每次变更触发自动测试。测试通过后自动部署到测试环境,验证通过后部署到生产。

十、CDC 详解:变更数据捕获

CDC(Change Data Capture)是数据摄入领域最重要的技术之一。它的核心思想是:监听源数据库的事务日志,实时捕获数据变更事件。

Debezium 是目前最流行的 CDC 平台,基于 Kafka Connect 构建。支持 MySQL、PostgreSQL、MongoDB、SQL Server 等主流数据库。

CDC 的三种实现方式:
1. 基于日志的 CDC(最常用):直接读取数据库的事务日志(PostgreSQL 的 WAL、MySQL 的 binlog),不影响源系统性能。Debezium 就是这种模式
2. 基于触发器的 CDC:在源数据库上创建触发器,捕获数据变更并写入变更表。性能开销较大
3. 基于查询的 CDC:定期查询源表检查数据变更。最简单但延迟高、效率低——不适合高频变更场景

CDC 的典型应用场景:
- 实时同步到数仓:数据库 → Debezium → Kafka → Flink → ClickHouse/StarRocks
- 缓存失效:数据库变更 → CDC → 通知缓存系统更新
- 搜索索引同步:数据库变更 → CDC → 更新 Elasticsearch 索引
- 微服务间数据同步:一个服务的数据变更通过 CDC 广播给其他服务

十一、流处理中的 Exactly-Once 语义

Exactly-Once 是流处理的"圣杯"——每条数据恰好被处理一次,不会丢失也不会重复。

实现 Exactly-Once 需要三个层面的协调:
1. 源头层(Source):数据源支持数据回放(Replay)——Kafka 支持根据偏移量重新消费
2. 计算层(Processor):Flink 使用分布式快照(Checkpoint)机制——定期保存算子的状态和消费者偏移量。故障时从最近一次 Checkpoint 恢复
3. 输出层(Sink):目标系统支持幂等写入——写入操作具有唯一标识符,重复操作不会产生重复数据

实际生产中,Exactly-Once 有显著的性能开销。对于大多数场景,At-Least-Once + 幂等去重是更实用的选择——性能更好,也能达到接近 Exactly-Once 的效果。

十二、数据管道的版本管理

数据管道也需要像代码一样进行版本管理。但数据管道的版本管理比普通代码更复杂——不仅要管理代码版本,还要管理数据版本和模型版本。

代码版本:ETL 脚本、dbt 模型、Airflow DAG 都在 Git 中管理。每次变更都经过代码审查。

数据版本:数据本身也需要版本控制——你知道一个月前的报表用的是哪个版本的数据吗?dbt 的 snapshot 功能、Iceberg 的时间旅行、以及数据湖的版本化存储都支持数据版本管理。

模型版本:ML 管道的模型需要版本化——包括模型文件、训练数据版本、训练参数、评测指标。MLflow 和 DVC 是常用的模型版本管理工具。

回滚策略:所有变更都必须可回滚。代码变更 → Git revert。数据变更 → 用快照或时间旅行恢复。模型变更 → 切换到上一个版本的模型。回滚测试应该定期演练。

延伸阅读

← 数据存储与建模 数据流水线与编排 →