← 返回博客

数据管道工程实践:从采集到落库的 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 小时内回可行性。

相关阅读

常见问题

批处理和流处理怎么选?

取决于数据延迟要求。批处理适合对实时性要求不高的场景(如日报、月报、离线分析),实现简单、成本低、易于重跑。流处理适合需要秒级或分钟级响应的场景(如实时监控、风控、推荐),实现复杂、成本高。建议:大部分场景从批处理起步,等明确需要实时处理时再引入流处理框架(如 Kafka + Flink)。不要一开始就上流处理,它带来的复杂度通常比想象中大得多。

增量更新和全量更新怎么选?

数据量 < 100 万行时,全量更新最简单——每天全量拉取、全量替换,不需要处理增量变更的边界情况。数据量 > 1000 万行时,增量更新几乎是唯一选择——每次只拉取上次同步后变更的数据,用时间戳或版本号标记。100-1000 万行之间是一个灰色地带,取决于你的数据库性能和更新频率。

数据质量怎么保证?

三层检查:① 源端检查——在数据进入管道之前做 schema 校验,拒绝不符合预期的数据;② 管道内检查——每个转换步骤后做完整性检查(行数对比、空值率、分布偏移);③ 目标端检查——落库后做对账(源端计数 vs 目标端计数),不一致则告警。不要等数据落地了才发现问题,越早发现越容易修复。

ETL 任务失败了怎么办?

幂等性是最重要的设计原则——同一个任务重跑多次应该得到相同的结果。具体做法:① 每个任务生成一个唯一的批次 ID,记录在目标表中;② 任务失败时,用批次 ID 清除已写入的部分数据,然后重跑;③ 设置重试机制(指数退避,最多 3 次),重试仍失败则告警,人工介入。

本文来自 AI Enable Harness 一线交付实践。需要同类系统或优化服务?

📡 本文同步发布平台: CSDN 知乎

订阅博客更新

新文章发布后第一时间邮件通知。不定期发送,不推销。

订阅 →