做大数据这行的朋友,一定绕不开三个字母:ETL。刚开始接触的时候,我以为 ETL 就是写几个 SQL 把数据从一个表挪到另一个表,简单得很。后来真正接手离线数仓、数据湖、数据中台之类的项目,才意识到 ETL 的体系比想象中大得多——从数据抽取的方式、转换逻辑的设计、加载策略的取舍,到调度、监控、数据质量保障,每一个环节都能决定你下游报表和数据服务的稳定性。
今天这篇不聊虚的,我把从理论到实战整个链路里踩过的坑、总结的方法、可以直接抄走的脚本和参数配置,完整梳理一遍。文章会围绕大数据领域的离线 ETL 场景展开,重点讲清楚每一步为什么这么设计、实际怎么落地,以及那些文档里不会写的问题排查经验。适合刚入门的数据开发、正在做大数据的毕设同学,以及想系统梳理 ETL 思路的同行参考。
1. 大数据场景下,ETL 到底在解决什么问题
1.1 所谓“大数据 ETL”,和传统数仓里的 ETL 有什么不同
ETL 全称是 Extract-Transform-Load,抽取、转换、加载。传统数仓时代,ETL 基本是 Informatica、DataStage、Kettle 这类工具的天下,数据量在 GB 级别,单机跑跑就完事了。
但到了大数据场景,情况完全变了。数据量从 GB 涨到 TB 甚至 PB,数据源也不再只是业务库里的几个表,而是包括了日志文件、埋点数据、第三方接口、消息队列里的实时流等等。这时候你再靠单机工具去抽数,抽个全量可能就要跑一天,完全不可行。
大数据 ETL 的本质变化在于三点:第一,计算引擎换成了分布式框架,比如 Spark、Flink、Hive,靠集群算力来扛数据量;第二,数据存储从关系型库换成了 HDFS、Hive 分区表、HBase、Doris、ClickHouse 这类分布式存储和查询引擎;第三,ETL 不再是“一次性搬数据”,而是变成了一个持续运行的、有调度有监控的数据管道。
你可以把传统 ETL 理解成“搬家”,一次性把东西搬完就结束了;而大数据 ETL 更像“城市供水系统”,建好管道之后每天都要稳定运行,供水质量合格、供水量可预期、管道出问题能快速定位。这也是为什么现在大家更爱讲“数据管道”这个词,其实就是 ETL 在大数据时代的进阶形态。
1.2 ETL 的三种典型场景:离线批处理、实时流、准实时增量
大数据 ETL 落到实际项目里,基本跑不出三种场景。
第一种是离线批处理,也是最常见的。典型链路是:每天凌晨通过调度平台触发任务,把业务库前一天的数据抽到 Hive 的 ODS 层,然后做清洗转换生成 DWD 层,再进一步汇总成 DWS 层的指标表,最后供 BI 报表和数据分析使用。这种场景对时效性要求不高,T+1 就可以,但对稳定性和数据准确性的要求极高。
第二种是实时流处理。数据源是 Kafka 里的消息流,通过 Flink 做实时清洗和聚合,结果写入 Doris、ClickHouse、Redis 等存储,支撑实时大屏、实时风控这类场景。严格来说这不叫 ETL 了,应该叫 ELT 或者流式数据集成,但核心思路是一样的——抽取、转换、加载,只是处理的节奏从“每天一次”变成了“每条消息”。
第三种是准实时增量,介于前两者之间。比如每 5 分钟或每 10 分钟拉取一次业务库的增量数据,同步到数仓或查询引擎。这种场景通常会用 DataX、Flink CDC、Canal 之类的工具来做。我见过不少团队低估了准实时场景的复杂度,因为“增量”听起来比“全量”简单,但实际涉及断点续传、去重、时区、延迟监控一堆问题。
你做项目规划的时候,第一步一定要先明确自己属于哪种场景。不同场景下技术选型和链路设计完全是两码事,别一上来就照着网上的所谓“最佳实践”去搭,大概率水土不服。
需要模型API调用? 免费领10W Token,多模型网关一键接入 Claude、DeepSeek 等主流模型。
2. 链路设计先行:工具选型和架构方案的取舍逻辑
2.1 常见开源 ETL 工具矩阵对比
聊完场景,接下来说工具。大数据领域做 ETL 的开源工具很多,但每个工具的定位和适用边界差别很大。我把常用的几类放一起做个对比,方便你选型。
| 工具 | 类型 | 典型用途 | 优点 | 短板 |
|---|---|---|---|---|
| DataX | 离线同步工具 | 各类数据源之间的批量同步,如 MySQL→HDFS/Hive | 插件丰富、部署简单、支持多种 RDBMS | 单机运行,超大流量下吞吐有限 |
| Sqoop | 离线同步工具 | Hadoop 与传统数据库之间的数据迁移 | 老牌、与 Hive/HDFS 集成度高 | 维护不太活跃,新特性少 |
| Flink CDC | 实时同步工具 | 基于 Binlog 捕获数据库变更,实时入湖入仓 | 实时性高、断点续传机制完善 | 对运维有一定要求,需要配套 Kafka 或 Doris 等目标端 |
| Canal / Debezium | CDC 组件 | 把数据库变更解析成消息流 | 作为中间件很稳定,生态成熟 | 本身不是完整的 ETL 工具,需要自己接下游 |
| Spark | 计算引擎 | 大规模数据的清洗、转换、聚合 | 吞吐量大、代码灵活、既能批也能微批 | 上手门槛比 SQL 工具高,需要调优 |
| Flink | 流计算引擎 | 实时 ETL、实时指标计算 | 毫秒级延迟、精确一次语义 | 状态管理复杂,排错成本高 |
| DolphinScheduler / Airflow | 调度平台 | 编排所有 ETL 任务的依赖和运行时机 | 可视化、支持复杂依赖、有告警 | 调度平台本身不干 ETL 的活,只是“发号施令” |
这里说一个我自己的判断标准:如果只是“把 A 库的表同步到 B 库”,DataX 是性价比非常高的选择,一个 json 配置搞定,不用写代码。但如果数据进了数仓之后要做复杂的清洗、多表关联、聚合,那就必须交给 Spark SQL 或 Hive SQL 来做。实时场景则优先考虑 Flink + Kafka 的组合。不要指望一个工具解决所有问题,ETL 链路里的每一段选最合适的工具,而不是选“最流行”的工具。
2.2 链路设计方案时需要考虑的关键因素
设计一条 ETL 链路,真正决定成败的往往不是技术本身,而是你在设计阶段有没有把下面这些事想清楚。
第一是数据量级和增长趋势。你可以先估算每天要处理的数据量有多大、高峰期是平时的几倍。这个直接决定了你要不要上 Spark、集群资源怎么配。曾见过有团队数据量一天就几百万行,却开了一个几十核的 Spark 集群去跑,资源浪费不说,调度排队的时间比任务本身还长。
第二是数据源的类型和接入方式。对方是 MySQL、PostgreSQL、Oracle,还是文件上传、消息队列、第三方 API?每种方式对应的抽取策略完全不一样。关系型库能做增量抽取,文件往往只能全量扫描,API 则要关注限流和翻页逻辑。
第三是 SLA,也就是你必须保证数据在什么时间点之前可用。比如老板要求每天早上 8 点前看到昨天的数据报表,那你的 ETL 任务最好在凌晨 3 点前全部跑完,预留出重跑和排查的时间。很多团队把任务排得满满当当,任何一个上游任务延迟 10 分钟,下游就全挂了,这种设计从一开始就有问题。
第四是数据质量要求。数仓里的数据是用来做决策的,准确性怎么强调都不为过。我在设计链路的时候一定会留出质量校验的环节,比如行数校验、主键重复校验、空值率校验,校验不过就告警并阻断下游任务。这个后面在实战部分会详细说。
说到底,ETL 链路设计就是在“时效、成本、质量、稳定性”这四者之间找平衡。不要追求每一项都最优,而是要找到符合你业务场景的最优组合。
3. 核心三阶段拆解:抽取、转换、加载各有各的门道
3.1 抽取层的三种方式:全量、增量、CDC,怎么选
抽取是 ETL 的第一步,也是所有问题的起点。抽取方式选错了,后面转换写得再漂亮都白搭。
全量抽取最简单,每次任务启动,把源表整个读一遍。优点是逻辑简单、数据完整,缺点是耗时、耗资源、对源库压力大。我一般只在两种场景下用全量:一是数据量确实小,几十万行以内;二是维表这种“变得慢”的数据,比如用户维表、商品维表,每天全量拉一次完全没毛病。
增量抽取则是在源表有更新时间字段(如 update_time)的前提下,只拉取上次任务之后变化的数据。这是离线数仓里最常用的方式。但增量抽取有两个隐患:第一,如果源表物理删除数据,你永远抽不到被删掉的那部分;第二,如果业务系统把更新时间更新得不准,漏数据就是必然的。做增量抽取前,务必先确认源表有没有可靠的更新字段,以及这个字段是否真的会被业务代码正确维护。
CDC(Change Data Capture)抽取则是通过解析数据库的 Binlog 日志来捕获数据的增删改操作。Flink CDC、Canal 都是这个路子。CDC 的优点是能捕获删除操作、实时性高、对源库压力极小,缺点是链路更长、组件更多,需要处理断点续传、DDL 变更同步、大事务拆解等复杂问题,适合实时或准实时场景。
选型的经验是先看业务需求和数据特征:等不了 T+1 就上 CDC;T+1 能接受就看表有没有可靠 update_time;两者都没有就老实全量。这三种方式不冲突,一条链路里完全可以把订单大表做成增量、维表做成全量、关键表再做 CDC 实时辅助校验。
3.2 转换层的通用处理套路
转换是整个 ETL 里最“百花齐放”的部分,不同业务清洗逻辑千差万别。但把大量项目经验沉淀下来,核心套路其实就那么几种。
数据清洗最常见:处理空值、去空格、格式统一、单位换算。比如业务库里电话号码可能有 +86 前缀,有的没有;时间字段有的是字符串,有的是时间戳;金额有的是元,有的是分。这些脏数据不在 ODS 到 DWD 的过程里清掉,后面所有统计指标都会跟着错。
去重也是高频操作。离线链路中,源表数据重复是家常便饭,尤其是业务库没有做唯一约束或者并发写入导致的问题。我的习惯是在 DWD 层对每个下游需要的业务主键做一次 row_number 去重,取最新的一条。用 SQL 写就是:
sql复制SELECT *
FROM (
SELECT t.*,
ROW_NUMBER() OVER (PARTITION BY order_id ORDER BY update_time DESC) AS rn
FROM ods.ods_order_info t
WHERE dt = '${bizdate}'
) tmp
WHERE rn = 1;
维表关联是最后一个高频套路。ODS 层拿到的往往是业务表的原始外键,比如 order 表里有 user_id、product_id,但报表要的是用户名、商品分类。这时候就需要去关联维表补全维度信息。用 Spark SQL 写就是普通的 JOIN,但要注意维表数据量如果很大,要优先考虑广播变量(Broadcast)或者将维表也做成分区表来减少 Shuffle。
转换层最需要警惕的其实是“过度转换”。我见过不少团队在 DWD 层把各种业务规则全都揉进 SQL 里,一个脚本几百行,最后没人敢动。建议的拆分逻辑是:清洗类操作尽量下沉到 ODS 到 DWD 的过程,指标计算类操作留到 DWS 层,DWD 只做标准化和维度补充,保持每层职责单一。
3.3 加载层的分区、小文件与写入策略
加载阶段,大多数离线 ETL 的目标端是 Hive 分区表。这里有两个问题最让人头疼,一个是分区策略怎么定,另一个是小文件怎么治。
分区策略方面,日期分区是绝对的主流,也就是 dt 字段。但再往细了分就有讲究了:日活千万级以下,天分区就够了;要跑小时级的数据就看业务需求,常见的是还加一个 hour 分区,或者干脆做成时间字段非分区。我个人的建议是,除非真的有小时级分析需求,否则不要轻易按小时分区,分区数量太多会导致元数据压力大、小文件问题严重,运维成本直线上升。
写入策略上,大多数离线场景用“先删除目标分区再写入”或者“覆盖写入”,保证任务重跑之后数据是一致的。这就是常说的“幂等性”——同一个任务无论跑多少次,最终数据结果都一样。增量 merge 的场景则要复杂一些,一般用全量覆盖加去重,或者用 Hive 的 MERGE 语法做 upsert(需要 ACID 表支持),再或者配合 Doris/ClickHouse 的 Unique 模型来解决。
小文件问题是加载阶段最经典的老大难。简单说,HDFS 不适合存储大量小文件,因为每个文件都要占一块元数据,NameNode 压力大,读取时随机 IO 也多。产生小文件的根源通常是:DataX 并发度设置过高、Spark 作业并行度过大、上游 Flume 或 Kafka 写入太碎、分区表每天产生大量新分区。治理思路是“合并 + 控制源头”:写任务时调整并发度或加上 coalesce,定期做文件合并,能合并的文件尽早合并。
4. 一次完整的离线 ETL 实战:MySQL 订单表同步到 Hive 数仓
4.1 场景设定与整体流程
理论说了那么多,下面走一遍完整流程。
假设场景是:业务库 MySQL 里有一张订单表 order_info,每天产生几十万行新数据,部分历史订单状态会更新。需要每天把数据同步到 Hive 数仓,清洗成标准格式后供后续数据分析使用,数据要求在每天早上 7 点前可用。
整体链路设计如下:凌晨 2 点由调度平台触发 DataX 任务,把 MySQL 的增量数据(update_time 大于上次同步时间)抽取到 Hive 的 ODS 层;凌晨 2 点半触发 Spark SQL 任务,对 ODS 层数据做去重、清洗、维表关联,写入 DWD 层;凌晨 3 点触发数据质量校验脚本,对 DWD 层做完整性、唯一性校验;最后输出校验报告,失败则发送告警通知。
这个流程图不用画得很复杂,关键是你心里要把每个环节的依赖关系和数据流向先理清楚。我做任何 ETL 项目,第一步都是先画这个数据流向和依赖关系图,确保每一步的输入输出边界清晰,再开始写具体实现。
4.2 ODS 层建表:先定规范再干活
ODS 层的数据要和源表保持基本一致,但建议还是加上日期分区,这样每次同步只操作当天分区,重跑任务也不会影响历史数据。建表语句如下:
sql复制CREATE TABLE IF NOT EXISTS ods.ods_order_info (
id BIGINT COMMENT '主键ID',
order_no STRING COMMENT '订单号',
user_id BIGINT COMMENT '用户ID',
product_id BIGINT COMMENT '商品ID',
amount DECIMAL(10,2) COMMENT '订单金额',
status TINYINT COMMENT '订单状态',
create_time TIMESTAMP COMMENT '创建时间',
update_time TIMESTAMP COMMENT '更新时间'
)
PARTITIONED BY (dt STRING COMMENT '日期分区')
STORED AS PARQUET
TBLPROPERTIES ('parquet.compression'='SNAPPY');
这里说几个细节。第一,存储格式建议用 Parquet 而不是文本格式,一是压缩比高,二是列式存储对查询友好。第二,字段类型要和源表对齐,金额必须用 DECIMAL 而不是 DOUBLE,DOUBLE 在精度上会出问题。第三,字段注释别省,数仓表动辄十几个字段,没有注释到后面谁都不敢维护。
4.3 DataX 抽取配置与调度集成
接下来用 DataX 完成从 MySQL 到 HDFS 的抽取。DataX 的核心是 json 配置文件,其中 reader 定义了从哪儿读,writer 定义了写到哪儿。下面是一个简化后的配置,实际使用中你还需要调整并发度、超时时间等参数。
json复制{
"job": {
"setting": {
"speed": {
"channel": 4
}
},
"content": [
{
"reader": {
"name": "mysqlreader",
"parameter": {
"username": "data_user",
"password": "xxxxxx",
"connection": [
{
"jdbcUrl": ["jdbc:mysql://10.0.0.1:3306/order_db?useSSL=false&characterEncoding=utf8"],
"querySql": [
"SELECT id, order_no, user_id, product_id, amount, status, create_time, update_time FROM order_info WHERE update_time >= '${last_time}'"
]
}
]
}
},
"writer": {
"name": "hdfswriter",
"parameter": {
"defaultFS": "hdfs://nameservice1",
"fileType": "parquet",
"path": "/warehouse/ods/ods_order_info/dt=${bizdate}",
"writeMode": "append",
"column": [
{"name": "id", "type": "bigint"},
{"name": "order_no", "type": "string"},
{"name": "user_id", "type": "bigint"},
{"name": "product_id", "type": "bigint"},
{"name": "amount", "type": "double"},
{"name": "status", "type": "long"},
{"name": "create_time", "type": "date"},
{"name": "update_time", "type": "date"}
]
}
}
}
]
}
}
实际执行时,通过 Java 命令行工具或脚本传入参数:
bash复制python datax.py -p "-Dlast_time='2025-01-01 00:00:00' -Dbizdate=20250101" job/order_info.json
这里有两个容易踩的坑。第一个是 querySql 方式虽然灵活,但 DataX 分片时没法做到多个 channel 并行读同一个 SQL,所以并发度设置再高实际也只有一个查询线程,量特别大的表建议改成 splitPk 方式。第二个是时区问题,MySQL 的 DATETIME 类型通过 JDBC 读出后,带的时区信息可能和 Hive 里的存储不一致,如果你发现时间字段差了几个小时,多半是这个问题。
调度集成方面,我优先推荐 DolphinScheduler。把 DataX 的 json 文件放到服务器上,在 DolphinScheduler 里创建一个 shell 节点执行上面的启动命令,再把依赖任务串起来。调度平台的好处是,失败自动重跑、依赖配置可视化、历史日志可查,这些是脚本裸跑给不了的。
4.4 用 Spark SQL 完成 DWD 层清洗转换
数据进了 ODS 层,接下来用 Spark SQL 做清洗和转换,写入 DWD 层。这个环节,你可以直接用 spark-sql 命令执行 SQL 文件,也可以写 PySpark 脚本。我个人更推荐前者,因为 SQL 可读性强、维护成本低,性能上也足够。
下面是核心的转换 SQL:
sql复制INSERT OVERWRITE TABLE dwd.dwd_order_detail PARTITION (dt = '${bizdate}')
SELECT
t1.id,
t1.order_no,
t1.user_id,
t2.user_name,
t2.mobile,
t1.product_id,
t3.product_name,
t3.category_name,
t1.amount,
TRIM(CAST(t1.status AS STRING)) AS status,
t1.create_time,
t1.update_time
FROM (
SELECT *,
ROW_NUMBER() OVER (PARTITION BY order_id ORDER BY update_time DESC) AS rn
FROM ods.ods_order_info
WHERE dt = '${bizdate}'
) t1
LEFT JOIN dwd.dim_user t2 ON t1.user_id = t2.user_id
LEFT JOIN dwd.dim_product t3 ON t1.product_id = t3.product_id
WHERE t1.rn = 1;
这个 SQL 里同时处理了两件事:去重和维度补全。ROW_NUMBER() 按订单 ID 分区,按更新时间倒序排序,取每条订单的最新记录;然后将用户维表和商品维表 JOIN 进来,补上用户名、手机号、商品分类等信息。你可以在 WHERE 里再加一些过滤条件,比如过滤掉状态为删除或无效的订单。
在跑这个任务时,有几个参数建议直接写进提交脚本:
bash复制spark-sql \
--master yarn \
--deploy-mode client \
--driver-memory 4g \
--executor-memory 8g \
--executor-cores 4 \
--num-executors 20 \
--conf spark.sql.shuffle.partitions=400 \
--conf spark.sql.adaptive.enabled=true \
--conf spark.sql.adaptive.coalescePartitions.enabled=true \
-f dwd_order_detail.sql
spark.sql.shuffle.partitions 这个参数尤其重要,很多任务性能上不去都是因为 Shuffle 分区数设置不合理。数据量小设大了会生成大量小文件,数据量大设小了会发生数据倾斜。开了 AQE(自适应查询执行)之后,Spark 会在运行时自动调整分区数,建议默认开着。
4.5 数据质量校验怎么做
ETL 跑完不等于任务结束,我最重视的反而是跑完后的那一套校验逻辑。校验不过,宁可把下游任务拦住,也不能让脏数据流下去。
常见的校验我做三类。第一类是行数校验:把源表(或上游 DWD 表)的总行数和目标表比对,差距超过阈值就告警。用 SQL 可以写成这样:
sql复制-- ODS 和 DWD 的总行数对比
SELECT
(SELECT COUNT(*) FROM ods.ods_order_info WHERE dt = '${bizdate}') AS ods_cnt,
(SELECT COUNT(*) FROM dwd.dwd_order_detail WHERE dt = '${bizdate}') AS dwd_cnt;
第二类是唯一性校验:检查业务主键(比如订单号)是否重复。如果 DWD 层的 order_no 还有重复,说明去重逻辑有问题,需要立刻排查。
sql复制SELECT order_no, COUNT(*) AS cnt
FROM dwd.dwd_order_detail
WHERE dt = '${bizdate}'
GROUP BY order_no
HAVING COUNT(*) > 1;
第三类是空值率和异常值检查:比如订单金额不允许为 NULL,用户 ID 不允许为 NULL,金额不允许出现负数(如果业务上不可能出现)等等。你可以把校验 SQL 写成脚本,通过调度平台每天执行,结果写进一张校验结果表,出错就发钉钉或邮件告警。
5. 高频踩坑实录:数据倾斜、小文件、时区问题一次讲清
5.1 数据倾斜的定位与处理
数据倾斜是离线 ETL 里出现频率最高、也最容易让人抓狂的问题。现象很典型:整个任务大部分 Task 几秒跑完,但有一两个 Task 跑了十几分钟甚至几个小时,最终整个作业卡在那。
倾斜的本质是数据分布不均,某些 Key 的数据量远大于其他 Key。比如订单表里某个头部用户下单量占全表的 30%,你 JOIN 用户维表或者做 GROUP BY 的时候,这个 Key 所在的 Reducer 就会特别慢。
定位方法:在 Spark UI 里看 Stage 详情,找到运行时间最长、输入记录数最大的那个 Task,看看它的处理数据量是不是其他 Task 的几十倍。如果是,再检查是不是某个 Key 的数据量特别大。
处理手段分几种。如果只是 GROUP BY 倾斜,可以做两阶段聚合,先给 Key 加随机前缀打散,做一次部分聚合,再去掉前缀做最终聚合。如果是 JOIN 倾斜,优先考虑把大 Key 的数据用 Broadcast Join 方式处理,让 Spark 把小表分发到每个 Executor,避免大 Shuffle。如果倾斜 Key 很少(通常一两个),可以在关联前用 skew join 的 hint 或者手动拆出倾斜 Key 单独处理。
我之前处理过一次极其极端的情况:一个表里有个“未知用户”的 ID 是默认值 0,占全表数据的 60%,导致 JOIN 用户维表时,所有为 0 的数据都落到同一个 Reducer。最终方案是先把 user_id=0 的数据拆出来不参与 JOIN,最后再用 UNION ALL 拼回去。这种特殊 Key 的处理思路网上很少讲,但实际项目里遇到概率很大。
5.2 小文件泛滥的根因与治理
小文件问题是数仓里最常见也最容易积累的问题。一开始你可能不在意,等几个月后 HDFS 上堆了几十万个小文件,再跑任务就会发现 NameNode 都快扛不住了,查询也越来越慢。
小文件产生的根源,多半是 ETL 写入时分区数或并发度过高。比如 DataX 设置了 20 个 channel,每个 channel 写一个小文件,这一个分区就产生了 20 个文件。再比如 Spark 任务里 spark.sql.shuffle.partitions 设成了 2000,每个 Shuffle 输出都写文件,最终写入的文件数就等于分区数。
治本的方法有三层。第一层是在写入前控制文件数:如果目标分区不大,可以用 coalesce 或 repartition 把最终写入的分区数控制在一个合理的范围。比如几百万行的数据,目标文件数量控制在 50 到 100 个之间比较合适。第二层是定期合并:对已经存在的分区做文件合并,可以写一个定时任务,把数据重新加载后覆盖写回,这样文件就合并了。第三层是开启 Hive 的自动合并功能,去避免 INSERT 语句把小文件写入目标表时产生的碎片。
我说一个相对容易落地的合并方案:在每天 ETL 任务结束后,加一个文件合并步骤,使用 Spark SQL INSERT OVERWRITE 从旧分区读取数据再写回,利用动态分区自动控制 Reduce 数量来减少文件数。虽然多跑一步,但换来的查询性能提升非常明显。
5.3 时间与时区不一致的坑
时间问题是那种“平时不炸、一炸就是大事”的坑。我处理过不止一次凌晨被叫起来排查,发现前一天的数据时间整体偏移了几个小时。
最典型的是 MySQL 的 DATETIME 和 Hive/Spark 的 TIMESTAMP 之间时区转换不一致。JDBC 在读取 DATETIME 时,默认会按 JVM 时区转一下,如果你集群的 JVM 时区是 UTC,而业务数据是北京时间,存到 Hive 里就会显示成少 8 个小时。
解决办法,一是统一集群层面的时区配置,把 Spark 和 Hive 的时区都设成和业务时区一致;二是如果是纯离线场景,建议把时间字段统一按字符串处理,避免隐式时区转换带来的不确定性。比如 ODS 层直接存储 yyyy-MM-dd HH:mm:ss 格式的字符串,到 DWD 层再转换成需要的格式。这个方法虽然不够“标准”,但在跨时区、跨系统数据同步时,能省掉大量排查时间。
5.4 幂等性设计的细节
所谓幂等性,就是同一个 ETL 任务,你跑一次和跑十次,最终结果完全一致。这是离线数仓最基本的要求,但实现起来不少团队都会踩坑。
最常见的错误是任务在写目标表时用了追加模式。一旦任务重跑,数据就会重复一份,下游统计直接翻倍。我见过一次事故,就是因为某天凌晨的任务失败后手动重跑,结果目标表里插入了两遍数据,第二天早上报表一大堆数字对不上,查了小半天才定位到这个原因。
所以离线任务写入目标表,尽量用 INSERT OVERWRITE ... PARTITION (dt = 'xxx'),先覆盖目标分区再写入新的数据,保证每次跑完状态一致。如果目标端不支持 OVERWRITE,比如某些数据库,就需要先 DELETE 目标分区再 INSERT,或者使用数据库自带的 upsert 机制。
还有一个容易被忽视的点:调度平台里的重试机制。DolphinScheduler 默认失败后会自动重试,但如果你没有把任务写成幂等,第一次跑了一半失败,重试再次从失败点继续,那么目标表里就会残留第一次写入的部分数据。建议先把任务设计成“失败整体重来”或“OVERWRITE 后重来”,再设置自动重试。
5.5 调度与监控的注意事项
最后说我一直在强调的一点:ETL 不是写完跑通就完事了,没有监控的管道等于一个哑弹,随时可能出问题。
调度平台选型上,我推荐 DolphinScheduler,开源、功能全、中文社区活跃,适合大多数中小团队。Airflow 也不错,但对 Python 能力有要求,而且调度引擎本身对运维的要求更高。阿里云这类云平台也自带调度服务,看你的部署环境来定。
监控方面至少要做到三件事。第一是任务失败告警:失败自动重试,重试三次还失败就发告警。第二是数据时效告警:任务整体延迟超过预期阈值要通知,别等业务方发现数据没出来才来找你。第三是数据质量告警:把前面讲的校验 SQL 注册成监控任务,校验不通过就拦截下游并告警。
有一个心得是我这些年总结出来的:ETL 链路里最贵的不是开发成本,而是排查成本。一个任务写出来可能只要一天,但上线后半年内因为各种问题被反复排查的时间,可能会花掉你几周。所以前期把日志打足、把监控配好、把幂等性做好,看似多花了很多时间,实际上是在帮你省命。
6. 最后聊聊我的一点个人感受
做了这么多年数据,越来越觉得 ETL 这个活,表面上是在搬数据,实际上是在为整个数据体系打地基。地基不稳,上面盖多少层楼都是危楼。
如果你正准备做大数据相关的毕业设计、或者刚进入数据开发这个岗位,我的建议是别急着追新框架、新技术,先把一条最简单的链路完整跑通:MySQL 到 Hive,用 DataX 做抽取、用 Spark SQL 做转换、用调度平台做调度,最后把监控和校验补上。这条路走通之后,你对整个大数据技术栈的理解会有一个质的变化。
最后再分享一个小技巧:所有 ETL 任务的代码和配置,都要做版本管理,每个脚本都要有注释、有负责人、有报警联系人。这些东西在项目初期看着像是“浪费时间”,但等项目变复杂后,你会发现这是让你半夜能睡个安稳觉的唯一保障。
