数据管道工程实践:从采集到落库的 ETL 架构设计
ETL(Extract, Transform, Load)是数据工程中最基础也最容易出问题的环节。本文从数据采集、清洗转换、加载落库三个环节出发,讨论管道设计中的关键决策:增量 vs 全量、批处理 vs 流处理、以及数据质量保障——适合正在搭建或优化数据管道的后端工程师和数据工程师。
先说结论:数据管道是”看起来简单,跑起来全是坑”
ETL 的需求描述起来很简单——“把数据从 A 搬到 B,中间做一下清洗”。但实际跑起来,边界情况、数据质量、失败恢复、性能问题会让简单的事情变得复杂。
本文从三个环节展开:数据采集、清洗转换、加载落库,每个环节给出关键决策和工程实践。
1. 数据采集(Extract)
1.1 采集模式
| 模式 | 适用场景 | 优点 | 缺点 |
|---|---|---|---|
| 全量拉取 | 小表(< 100 万行) | 实现简单,不需要处理增量变更 | 数据量大时性能差 |
| 增量拉取 | 大表(> 1000 万行) | 每次只拉取变更数据,性能好 | 需要处理删除和更新 |
| 日志订阅 | 实时数据(CDC) | 无侵入,实时性高 | 实现复杂,依赖消息队列 |
1.2 增量采集的三种方式
时间戳:源表增加 updated_at 字段,每次拉取 WHERE updated_at > last_sync。最简单,但无法捕获删除操作。
版本号:源表增加 version 字段,每次更新时 version +1。可靠,但需要源系统配合。
CDC(Change Data Capture):通过数据库的 binlog 或 WAL 日志捕获变更,如 Debezium + Kafka。实时性最高,但架构最复杂。
建议:从时间戳起步,需要删除检测时加一个软删除标记(deleted_at),等规模上来了再考虑 CDC。
2. 清洗转换(Transform)
2.1 常见转换类型
// 清洗:去空值、去重、格式统一
function clean(row: RawRow): CleanRow {
return {
id: row.id,
email: row.email?.toLowerCase().trim(),
name: row.name?.trim() || 'unknown',
created_at: new Date(row.created_at).toISOString(),
};
}
// 转换:类型映射、单位换算、枚举映射
function transform(row: CleanRow): TargetRow {
return {
...row,
status: STATUS_MAP[row.status] ?? 'unknown',
amount: Math.round(row.amount * 100), // 元转分
};
}
2.2 数据质量检查
每个转换步骤后做三项检查:
| 检查项 | 方法 | 阈值 |
|---|---|---|
| 行数对比 | 源端行数 vs 当前步骤行数 | 偏差 > 5% 告警 |
| 空值率 | 关键字段的空值比例 | > 1% 告警 |
| 分布偏移 | 数值字段的均值/分位数变化 | 偏差 > 10% 告警 |
2.3 幂等性
每个任务必须可重跑。幂等性的核心是”同一个任务跑多次,结果一样”。
// 用批次 ID 实现幂等
const batchId = uuid();
await db.query(`
DELETE FROM target_table WHERE batch_id = $1
`, [batchId]); // 先清空已写入的批次数据
await db.query(`
INSERT INTO target_table SELECT *, $1 as batch_id FROM staging
`, [batchId]); // 重新写入
3. 加载落库(Load)
3.1 加载策略
| 策略 | 做法 | 适用场景 |
|---|---|---|
| 全量替换 | TRUNCATE + INSERT | 小表,每日全量 |
| 增量合并 | UPSERT(INSERT ON CONFLICT) | 大表,增量更新 |
| 分区切换 | 写入临时表 → 原子切换 | 大表,全量更新 |
3.2 分区切换实现
-- 1. 写入临时表
CREATE TABLE target_20260722 (LIKE target INCLUDING ALL);
INSERT INTO target_20260722 SELECT * FROM staging;
-- 2. 原子切换
BEGIN;
DROP TABLE IF EXISTS target_old;
ALTER TABLE target RENAME TO target_old;
ALTER TABLE target_20260722 RENAME TO target;
COMMIT;
-- 3. 清理旧表
DROP TABLE IF EXISTS target_old;
这种方式的好处是:切换是原子的,切换过程中查询不受影响,切换失败可以快速回滚。
4. 监控与告警
# 每分钟检查管道健康状况
check_pipeline() {
local last_run=$(psql -t -c "SELECT max(created_at) FROM pipeline_log" | xargs)
local now=$(date -u +%s)
local diff=$((now - $(date -d "$last_run" +%s)))
if [[ $diff -gt 3600 ]]; then
echo "管道超过 1 小时未运行" | send_alert
fi
}
总结
| 环节 | 关键决策 | 推荐做法 |
|---|---|---|
| 采集 | 全量 vs 增量 | 小表全量,大表增量(时间戳) |
| 清洗 | 数据质量检查 | 每步检查行数、空值率、分布 |
| 加载 | 幂等性 | 批次 ID + 先清后写 |
| 恢复 | 失败重试 | 指数退避,最多 3 次 |
| 监控 | 管道健康 | 检查最后运行时间,超过阈值告警 |
数据管道的核心不是”搬得快”,而是”搬得稳、搬得准、出问题了能快速恢复”。 幂等性设计比性能优化重要得多——一个能稳定重跑且结果一致的管道,比一个快但坏了修不好的管道值钱一百倍。
需要数据采集或 ETL 管道设计服务?联系我们,说清你的数据源与规模,24 小时内回可行性。
相关阅读
- 合规数据采集实践:爬虫项目怎么做才不踩法律与反爬红线 —— ETL 管道的上游数据采集环节
- 运维自动化脚本模式:从一次性脚本到可维护工具 —— 数据管道运维的自动化实践
常见问题
批处理和流处理怎么选?
取决于数据延迟要求。批处理适合对实时性要求不高的场景(如日报、月报、离线分析),实现简单、成本低、易于重跑。流处理适合需要秒级或分钟级响应的场景(如实时监控、风控、推荐),实现复杂、成本高。建议:大部分场景从批处理起步,等明确需要实时处理时再引入流处理框架(如 Kafka + Flink)。不要一开始就上流处理,它带来的复杂度通常比想象中大得多。
增量更新和全量更新怎么选?
数据量 < 100 万行时,全量更新最简单——每天全量拉取、全量替换,不需要处理增量变更的边界情况。数据量 > 1000 万行时,增量更新几乎是唯一选择——每次只拉取上次同步后变更的数据,用时间戳或版本号标记。100-1000 万行之间是一个灰色地带,取决于你的数据库性能和更新频率。
数据质量怎么保证?
三层检查:① 源端检查——在数据进入管道之前做 schema 校验,拒绝不符合预期的数据;② 管道内检查——每个转换步骤后做完整性检查(行数对比、空值率、分布偏移);③ 目标端检查——落库后做对账(源端计数 vs 目标端计数),不一致则告警。不要等数据落地了才发现问题,越早发现越容易修复。
ETL 任务失败了怎么办?
幂等性是最重要的设计原则——同一个任务重跑多次应该得到相同的结果。具体做法:① 每个任务生成一个唯一的批次 ID,记录在目标表中;② 任务失败时,用批次 ID 清除已写入的部分数据,然后重跑;③ 设置重试机制(指数退避,最多 3 次),重试仍失败则告警,人工介入。