很多人一提大数据,脑子里立马浮现出各种炫酷的可视化大屏、实时跳动的数据流,还有那个号称能预测一切的人工智能。但我干了这么多年的数据工程,说实话,真正花掉我最多时间、也最影响项目成败的,恰恰是最不起眼、甚至有点枯燥的三个字母——ETL。
ETL是Extract、Transform、Load的缩写,翻译过来就是抽取、转换、加载。三个词连起来看,就是把数据从各种数据源里搬回来,洗干净成能用的样子,再放到数据仓库或者数据湖里,供下游分析、报表、算法去消费。这篇文章我会把大数据场景下的ETL从理论讲到实战,包括分层设计、工具选型、常见调优手段,还会拿一个实际项目把整条链路走一遍。不管是做毕设的学生、准备大数据面试的求职者,还是刚接手数仓建设的工程师,我相信都能从这里拿走一些能直接用的东西。
1. 先讲清楚:为什么在大数据领域还要专门谈 ETL
1.1 从一条日志到一张报表,ETL 到底干了什么
先看一个最简单的例子。用户在前端点了一下商品链接,这个动作如果被埋点系统捕获,会变成一条JSON日志,丢进消息队列。这条日志本身没有任何意义,一堆七零八落的日志也没法直接做决策,它要变成报表上的“UV上涨了5%”这类指标,中间隔着一整套加工工序。
第一步是把日志从Kafka里读出来,落成一个最接近原始格式的表,这步是抽取和初步加载;第二步是把JSON里的字段拆开,把时间统一成北京时间、把设备号脱敏、把明显异常的数据过滤掉,这叫转换;第三步是把处理好的明细数据按天分区写进数仓的明细层,再汇总成汇总层,最后供报表或者大屏查询,这又是加载。
整条链路走下来,几乎每一步都在做ETL。从一开始粗放的原始数据,到后面越来越规整、越来越贴近业务口径的数据,能力就一层一层叠起来了。经常有刚入行的朋友问我,做数据工作到底难在哪,我通常会回答:难就难在“让数据从不可用变成可用”这一路,而这一路的主要工作就是ETL。
1.2 大数据 ETL 和传统数据库 ETL 的本质差别
在十年前,企业里的数据量还算可控,ETL一般怎么做呢?用Oracle或者MySQL,开一个存储过程或者定时SQL,insert into select ... where ...,数据就在数据库内部流转,抽一张业务表、清洗一下、再插入汇总表,一台性能好点的服务器就够了。这种模式在GB级别数据量下跑得很顺,问题也不多。
但到了大数据阶段,情况完全变了。数据量从GB跳到TB、PB,数据源不再只有业务库,还有日志文件、埋点事件、第三方接口、物联网设备数据;数据格式也更加复杂,JSON、Avro、Parquet、协议二进制混合出现;更关键的是单机算不动了,必须把计算任务分布到一堆机器上,由框架去调度、编排、容错。
我把这两代ETL的区别整理成了一个表:
| 对比维度 | 传统数据库 ETL | 大数据 ETL |
|---|---|---|
| 数据规模 | GB 级,单机可处理 | TB/PB 级,分布式集群 |
| 数据源 | 关系型数据库为主 | 数据库、日志、消息队列、对象存储、API 都有 |
| 数据格式 | 结构化行数据 | 结构化、半结构化、非结构化混合 |
| 计算引擎 | 数据库 SQL/存储过程 | Spark、Flink、Hive 等分布式计算框架 |
| 调度方式 | 简单的定时任务 | 有依赖关系的复杂工作流调度 |
| 失败处理 | 重跑一下 SQL | 任务级重跑、断点续跑、幂等性设计 |
传统ETL的思路是一台机器单打独斗,大数据ETL的核心则是把“数据搬运和加工”这件事拆成可并行、可重试的分布式作业。举个例子,原来要处理100GB数据,我觉得已经是很大的工程了,但在Spark里,100GB可能只是几分钟的活,前提是你把数据分区分得够好、把join写得够聪明。
还有一个容易忽略的差别:传统数据库里,ETL容易做到强一致,事务回滚也很顺手;大数据场景下,很多写操作是“先写临时目录再原子rename”这种方式,或者靠分区来保证可重跑。设计ETL任务时,能不能重跑、重跑会不会产生重复数据,这两个问题比单机时代重要得多。
1.3 ETL 与 ELT:选错顺序可能要付出不小代价
ETL叫抽取、转换、加载,ELT则是抽取、加载、转换,多了个顺序差别,但工程上差异很大。ETL是先把数据洗干净,再放进目标系统;ELT是先把原始数据一股脑倒进数仓或数据湖,等要用了再在数仓内部做转换。
很多人第一次意识到“顺序也能聊出花样”,是在接触云数仓或者湖仓一体架构的时候。为什么ELT越来越流行?首先,分布式存储太便宜了,原始数据多放几天也不心疼;其次,现代数仓计算能力强,临时想做哪个维度的分析,直接在仓内写SQL就能处理,不需要重新走一遍复杂的搬运链路;还有就是可回溯性好,原始数据永远在那里,规则变了重算就行。
但ELT也不是万能解药。如果下游系统对数据质量极其敏感,或者目标库存储和计算成本都很高,比如某个指标报表每天要直接对接业务系统,那还是老老实实先ETL,把脏数据挡在门外。还有一个我自己的经验之谈,很多公司表面上叫数据分析,实际做的就是“把几十张Excel和旧系统的导出文件整合起来”,这种高度异构的场景,直接用ELT把文件加载进来再清洗,反而比痛苦的转换过程更现实。
所以我的建议是,不要迷信某一个模式。架构设计的时候问自己三个问题:下游需要多干净的数据、重算历史数据的频率高不高、目标系统的计算和存储成本扛不扛得住。
需要模型API调用? 免费领10W Token,多模型网关一键接入 Claude、DeepSeek 等主流模型。
2. 核心环节拆解:抽取、转换、加载每一步的坑与思路
2.1 数据抽取:全量、增量与 CDC 的取舍
抽取是整个ETL的入口,也是问题最容易被忽视的地方。源端数据一分钱没少,但因为抽取方式选得不对,下游收到的东西可能差得十万八千里。
常见的抽取方式有三种。
第一种是全量抽取。适合那些数据量常年稳定的小表,比如几万行的省份表、商品类目表,每天快照一次很划算。大表绝对不建议每天全量,否则一次同步几个亿行,再快的通道也扛不住,而且源库压力非常大。
第二种是增量抽取。最常见的是依靠源表里的时间字段,比如update_time,每次取“大于上次同步时间”的数据。这个方案简单,但坑非常明确:如果源系统在更新记录时没时间去刷新时间字段,或者有人直接删了行,那增量抽到的数据就会漏。为了兜底,我一般会定期做一次全量对账,或者把增量保留时间拉长。
第三种是CDC,也就是变更数据捕获。主流做法是解析MySQL的binlog,把数据表上的insert、update、delete操作变成一条条变更事件,实时流到下游。比如Flink CDC就是干这个的。CDC的好处是实时性高、对源库压力小,也不会漏掉删除操作,但代价是要理解binlog的格式,维护一些额外配置,比如binlog的保留时间、大事务带来的binlog体积激增问题。
三种抽取方式选择上,我的经验是:维度表、配置表这类小表,直接全量;业务大表优先增量,如果业务实时性要求高,再加一层CDC同步链路,用Kafka把变更事件送出去;两种方式同时存在时,注意两边数据最终要能对得上。选型不是越高端越好,很多项目其实用DataX配一个增量定时任务就够了,没必要一上来就堆组件。
2.2 数据转换:清洗、标准化与宽表构建
抽取只是把数据搬回来,转换才是ETL里最考验“对业务的理解”的地方。很多人一提到转换就以为是写一堆replace、drop duplicates,其实远不止这些。
清洗要做的首先是识别哪些数据不可信。比如日志里的user_id出现负数或者空串,时间戳溢出,价格字段出现字母,这些都要有明确的分流策略。我习惯在清洗阶段设定一个原则:能修正的修正,不能修正的先落到一个异常表,绝不直接把这个脏值继续往下传,因为下游每个环节都会放大脏数据的影响。
标准化也容易忽略。一个项目里最好统一时区、统一日期格式、统一枚举值。比如业务库里的性别字段有用0/1的,有用male/female的,还有用中文“男”“女”的,不做标准化映射,后面做报表或者算法特征都得痛一次。
再往上走是宽表构建。分析业务方要的指标往往是多个事实表join出来的,比如“某品类商品在华东地区的销售情况”,每次让他们从明细层自己join,既不高效也容易算错口径。ETL阶段提前把这些常用维度组合成宽表,相当于把分析口径固化下来。但宽表也切忌贪多,一个宽表几百个字段,写入成本高、线上改起来痛苦,真正合理的做法是按业务域切,用户域、商品域、交易域分开建。
2.3 数据加载:写 ODS、DWD、DWS 的节奏把控
加载环节的重点是“写到哪儿、怎么覆盖”。大多数数仓项目都会做分层,ODS层原样保留上游数据,DWD层做清洗和标准化,DWS层做业务域的汇总,再往上还有ADS应用层。每一层的数据加载都有讲究。
ODS层要求的是“忠实还原”,上游什么格式,这里就什么格式,最多加一个抽取时间字段,方便追踪数据落地时间。这一层通常按天分区,也允许存在部分重复数据,因为后面还会去重。
DWD层的关键是幂等性。因为不少任务可能会因为上游迟到数据而重跑,如果重复跑一次就产生双份结果,那数据就对不上了。我常用的做法是先把当天目标分区drop掉,再用Spark或者Hive重新覆盖写,这样无论跑几次结果都是一致的。
DWS层则要面对小文件问题。如果你用Spark往Hive写聚合结果,默认可能产生成千上万个几十KB的小文件,后续查询会被文件个数拖死。写之前根据数据量做一次repartition或者coalesce,把输出文件数量控制在一个合理的范围,比如单文件50MB到200MB之间,这是写DWS层时非常重要的一个动作。
3. 工具选型与架构设计:别一上来就上全家桶
3.1 离线批处理与实时流处理怎么选
做大数据ETL,绕不开一个问题:用离线的Spark/Hive,还是上实时的Flink?如果只看趋势,很多人会觉得实时流处理未来会替代批处理,但从我接触的项目看,两者很长一段时间会并存。
有个判断方法我一直在用:如果业务要的是“今天看昨天的报表”,T+1完全能接受,那就用离线批处理,稳定、好排查、成本低。如果业务要的是“刚才发生的异常现在就要盯住”,比如风控、秒级监控、实时大屏,那才需要上Flink做流式ETL。
最怕的是为了展示技术栈而强行实时化。曾经有个同学组的毕设,明明是一个日更数据分析平台,非要把所有环节都改成流处理,结果数据源一天才提供一次文件,实时链路天天空转,运维成本还高。工具是服务业务的,不是用来炫技的。
3.2 任务调度与集群部署策略的配合
ETL任务跑起来还只是第一步,怎么调度、依赖关系怎么编排同样重要。常见的工作流调度工具有Airflow和DolphinScheduler,前者偏向开发,用Python写DAG很灵活;后者是国产开源项目,Web界面可视化定义依赖,中文资料多,对团队更友好。
调度设计要解决的问题是“先抽哪个、再算哪个、失败怎么办”。最典型的一种依赖是:业务库抽取任务完成之后,才启动ODS到DWD的清洗任务,清洗任务跑完再触发汇总任务。这种依赖关系需要在调度系统里显式表达,而不是靠几个cron盲跑——否则你永远不知道报表迟到是因为数据抽取卡住还是转换任务挂了。
再说集群部署策略。我见过不少团队所有的任务都堆在一个默认队列里,跑起来之后大任务占光了资源,小任务全部排队,整条ETL链路动不动就超时。比较合理的做法是给不同的任务划分YARN队列:离线ETL队列、实时计算队列分开;小批量任务和大任务也分开,避免互相干扰。任务调度和资源隔离是两件事,但两者要配合着设计。
3.3 常见工具矩阵与适用边界
数据工程领域工具非常多,很多人刚接触时容易挑花眼。我的建议是,先根据自己手里的数据源和下游使用方式,把工具按“数据同步、离线计算、实时计算、工作流调度”四个位置填进去,缺哪儿补哪儿。
| 环节 | 常用工具 | 适用场景 |
|---|---|---|
| 数据同步 | DataX、Sqoop、Kettle | 关系型数据库与大数据存储之间的批量搬迁 |
| 离线计算 | Spark SQL、Hive、Hadoop MR | 海量数据的清洗、join、聚合,T+1周期为主 |
| 实时计算 | Flink、Flink CDC | 秒级延迟的流式ETL、增量变更捕获 |
| 工作流调度 | Airflow、DolphinScheduler、DataWorks | 任务依赖编排、定时触发、失败告警 |
经验上,中小团队的数据同步用DataX就够了,社区活跃、上手简单;如果数据量和并发到了一定的规模,可以引入Spark统一处理同步和转换,把链路简化;实时计算不是必选项,只有当需求明确要秒级或者分钟级延迟时才引入Flink。
工具选型还有一个容易犯错的地方:看到别人用什么组件就觉得自己的项目也要用什么。组件越多,交接和运维成本越高。我自己会倾向于少而精,能用SQL解决的事情绝不多写一个Java服务,能用一条Spark SQL解决的事情绝不多搭一套Flink作业。
4. 从零到一实战:用户行为日志数仓的 ETL 链路
4.1 需求与分层设计
用一个我比较熟悉的场景来拆解整条链路:一个电商平台的用户行为数据分析项目。这个项目很像很多人做的毕设或者期末大作业,但数据链路是完整生产级的。
需求大概这些:统计每天的活跃用户数、各品类浏览热度、核心漏斗转化率,还要把这些指标输出到大屏和日报表上。数据源主要有三块:业务库MySQL里的用户表和订单表、前端埋点打到Kafka里的行为日志、第三方渠道提供的一份CSV商品信息文件。
按前面讲的思路,我把数仓分为四层。ODS层先把三块数据原样落地;DWD层做清洗和标准化,解析JSON日志、统一时间字段、过滤无效数据;DWS层按用户和商品两个主题做聚合;ADS层按大屏需求计算指标。每层之间通过调度任务串起来,形成完整链路。这个设计的核心考虑是让每一层职责单一,上层不用关心下层是怎么洗数据的,下层也不用关心上层怎么消费,将来不管是加指标还是新增数据源,改动范围都尽量集中。
4.2 用 Flink CDC + DataX 抽业务库,用 PySpark 处理日志
先解决业务库和离线文件怎么进来。MySQL业务表和订单表,我采用了两种抽取方式组合:订单表数据量大且更新频繁,用Flink CDC同步到Kafka,再由Spark Streaming落Hive;用户表量级不大,用DataX每天全量抽取成Parquet文件,一次性写进ODS和DWD。
日志部分用PySpark从Kafka读JSON并解析。下面是一段真实可跑的示例代码,演示了最常见的日志ETL写法:
python复制from pyspark.sql import SparkSession
from pyspark.sql.functions import col, get_json_object
spark = SparkSession.builder \
.appName("ods_user_behavior_log") \
.config("spark.sql.adaptive.enabled", "true") \
.getOrCreate()
# 读取Kafka日志,测试环境可以从小文件代替
df_stream = spark.read \
.format("kafka") \
.option("kafka.bootstrap.servers", "kafka-node1:9092") \
.option("subscribe", "user_behavior") \
.option("startingOffsets", "earliest") \
.load()
# 把value字节流转成字符串
df_json = df_stream.selectExpr("CAST(value AS STRING) as json_str")
# 从JSON中抽取核心字段
df_parsed = df_json \
.withColumn("user_id", get_json_object(col("json_str"), "$.user_id")) \
.withColumn("item_id", get_json_object(col("json_str"), "$.item_id")) \
.withColumn("behavior", get_json_object(col("json_str"), "$.behavior")) \
.withColumn("event_time_str", get_json_object(col("json_str"), "$.event_time"))
df_parsed.show(5, truncate=False)
这段代码只是一个起点,真实场景里JSON的嵌套层级很深,解析逻辑可能要写几十个字段。用from_json配合schema定义会更灵活,字段缺省、嵌套数组都能处理。我经常跟同事说,日志ETL里80%的代码就是这段解析逻辑的对齐工作,想清楚字段映射再动手,比写代码本身更能省时间。
4.3 增量清洗去重与分区写入 Hive
日志处理通常比业务库同步更复杂,因为日志是流式的,天然有重复,而且不能靠主键约束。所以DWD层的清洗去重一定要做。
下面这段代码延续上面的例子,完成DWD层的输出:
python复制from pyspark.sql.functions import col, row_number, to_date
from pyspark.sql.window import Window
# 1. 过滤明显异常数据
df_clean = df_parsed \
.filter(col("user_id").isNotNull() & (col("user_id") != "")) \
.filter(col("behavior").isin("pv", "cart", "fav", "buy")) \
.withColumn("date", to_date(col("event_time_str"), "yyyy-MM-dd HH:mm:ss"))
# 2. 按用户+商品+时间窗口去重,保留最早一条
window_spec = Window.partitionBy("user_id", "item_id", "date", "behavior") \
.orderBy(col("event_time_str").asc())
df_dedup = df_clean.withColumn("rn", row_number().over(window_spec)) \
.filter(col("rn") == 1) \
.drop("rn")
# 3. 按天分区写入Hive,覆盖本次分区保证幂等
df_dedup.repartition(4) \
.write \
.mode("overwrite") \
.partitionBy("date") \
.format("parquet") \
.saveAsTable("dwd.dwd_user_behavior")
这里面有几个细节非常值得注意。repartition(4)不是随便写的,它控制输出文件数量,避免写Hive时产生几百个小文件;mode("overwrite")配合分区字段,保证同一批数据重跑不会重复叠加;row_number去重逻辑要选好排序字段,否则去重结果不稳定。
写完之后我还会做一个快速校验,比如对比DWD表行数和ODS层行数,检查空值率是否突然升高。不要小看这一步,很多线上事故都是从一次“没校验就上线”开始的。
4.4 数据可视化大屏前的最后一道工序
DWD层数据准备好了,并不意味着可以直接画大屏。大屏上展示的往往是“今天的成交总额”“各品类热度排行”“近7日活跃趋势”,这些指标如果每次查询都从DWD明细去count、group by,数据库压力大,查询还慢。
所以在大屏展示之前,还要有一层ADS层或者指标表。它在ETL中通常是定时任务,每天凌晨把DWD层明细汇总成一张小的指标表,比如这样:
sql复制INSERT OVERWRITE TABLE ads.ads_gmv_daily
SELECT
date,
SUM(IF(behavior='buy', 1, 0)) AS order_cnt,
COUNT(DISTINCT user_id) AS uv,
COUNT(*) AS pv
FROM dwd.dwd_user_behavior
GROUP BY date;
这个表只有几十行或者几百行,大屏接口查询时毫秒级返回。所谓的免费数据可视化大屏、ECharts大屏,其实真正复杂的部分都不是前端渲染,而是前端之前这张ADS表有没有准备好。ETL把数据做成“拿来即用”,大屏自然就快了。
5. 调优经验与面试高频考点:这些年踩过的坑
5.1 数据倾斜:ETL 任务最常见的性能杀手
做过一段时间大数据后,你会发现大多数任务慢,不是因为集群不行,而是因为数据倾斜。所谓数据倾斜,就是数据在分布式节点之间分配不均,少数几个节点几乎承载了全部压力,其他节点闲着没事干。
典型场景发生在group by某个维度或者join的时候。比如电商大促当天,“爆款商品ID”的日志条数是其他商品的几十上百倍,按商品ID做聚合时,那个热点key所在的task要处理海量数据,其他task几分钟跑完,唯独这一个卡了一个多小时。
解决数据倾斜有很多种思路,我按优先级排一个:
- 先看能不能过滤热点key,比如对“其他”类做单独处理
- 再看能不能广播小表,让join不发生shuffle
- 如果还是倾斜,给热点key加随机盐,分成两步聚合,先按加盐后的key做group by,再去掉盐做最终聚合
- 最后才考虑调整并行度和内存参数,这类手段只能缓解,不能根治
值得强调一下加盐的思路。假设商品ID为“10001”这一条就占了全表30%,直接group by会卡死。可以先给这个商品ID拼一个1到10的随机后缀,把它分成10个子组,每组的量立刻变成原来的十分之一,先各自聚合,再合并结果。原理简单,但很实用。
5.2 从 N+1 问题谈大数据场景下的反模式
N+1问题这个词,最早是在讲关系型数据库ORM的时候出现的:先查一次主表得到N条记录,再循环N次查询子表,导致查询次数变成1+N。在大数据ETL中,我经常看到类似的错误被以“脚本化”的方式复现。
有个真实的例子,同事把几万条用户ID从数据库导出,然后用Python写了个for循环,逐条去调另一个系统的HTTP接口拿用户画像字段,再拼装成结果文件。几万条数据跑了快一天,后半段还经常因为连接超时而中断。这种模式单看每一行代码都没问题,但整体执行方式完全违背了分布式计算的思路。
正确的做法是至少做到两点:第一,把目标数据源的数据全部或者按分区批量导出,不要一条条取;第二,用Spark或者Hive做关联,而不是在Python循环里做关联。一次分布式join,几万条和几亿条的差距无非是几分钟和几十分钟,但如果用循环脚本,几万条就能跑一天。
识别反模式也有个笨办法:看运行时间是不是和行数呈线性增长,以及看数据库连接数是否异常高。如果任务运行30分钟,有25分钟都在源库查询,那基本就是踩了N+1的坑。
5.3 排障与体检清单
ETL任务出问题不可怕,可怕的是没有排查思路。我总结了一个相对固定的排查顺序,每次任务失败都按这个来:
- 第一步看日志定位,是任务调度没触发,还是抛了异常。调度平台和计算引擎日志要看细一点,别只看最后几行。
- 第二步看资源,任务失败或者变慢,先查集群资源有没有打满,是不是被其他任务占了队列。
- 第三步看数据,行数突然波动、空值率异常,往往是上游数据源变了,而不是ETL逻辑出了问题。
- 第四步看源端,增量同步的坑大多出在源端,binlog没开、时间字段没更新、字段类型变了,都可能造成静默失败。
这里再给一个数据质量的体检清单,每天定时任务跑完后执行一遍,会省掉不少半夜被电话叫醒的痛苦。
| 检查项 | 检查方式 | 异常处理建议 |
|---|---|---|
| 主键/去重键重复率 | count与distinct count对比 | 重复率突增,优先检查去重逻辑 |
| 空值率 | 关键字段is null占比 | 突增则检查清洗规则是否覆盖新数据 |
| 行数波动 | 与昨日/上周同日对比 | 波动超过阈值,触发上游数据源排查 |
| 时间分区完整性 | 检查Hive分区/目录是否存在 | 缺分区则触发补数任务 |
5.4 面试现场最常被追问的 ETL 问题
这几年带过不少新人,也帮朋友做过模拟面试,发现大数据岗位面试中ETL相关问题的出现频率比想象中高得多。这里把高频问题简单梳理一下,供准备面试的朋友参考。
第一类问题是概念题:ETL和ELT有什么区别?回答时先讲顺序差异,再补一句“大数据场景因为存储和计算成本下降、可回溯性要求提高,ELT使用越来越普遍”,基本就是加分项。第二类是架构题:数仓为什么分层?核心答三点:职责分离、口径统一、可回溯。第三类是实战题:你遇到过数据倾斜吗,怎么解决的?这种问题要讲具体案例,不要背理论,能说出“某个热点key加盐分两步聚合”基本能过。第四类是场景题:一张千万级表和一个亿级表join,你怎么优化?可以从数据量、倾斜情况、join类型多个维度展开。
还有一个几乎必问的问题是:Spark写Hive为什么会有小文件?这个问题不会直接问小文件,而是考你有没有真实操作的体会。从RDD分区到文件数一一对应这个角度去答,再说一个coalesce控制和合并输出文件的方法,面试官就知道你是真的写过生产代码的。
做数据这行时间越长,越发现ETL这个看似基础的东西,其实是整个数据工程的基本功。很多后来让我头疼的问题,根源都出在当初某一层ETL没有认真设计——要么没有分区策略,要么没有幂等保证,要么顺手写了个循环脚本。
对我个人而言,最看重的不是用了多新的框架,而是这条数据链路的稳定性。每次看到新的数据工具出来,我也会去试用,但真正留到最后的生产链路,永远是设计得最简单、最容易排查、最便于重跑的那一套。
如果你是正在做大数据毕设的学生,我建议把ETL当成展示你工程能力的最佳舞台,一个能跑通、有校验、能可视化展示的数据平台比任何华丽的概念都更能打动老师;如果你想转行做数据,ETL也是相对友好的切入点,学历或者专业背景不是不能进这个行业,动手处理过真实脏数据、真被数据倾斜坑过、真的画过一张能汇报的大屏,这些才是在面试里最值钱的东西。
最后再分享一个小技巧:以后设计任何一条ETL任务,先写一份“这份数据的校验规则”,哪怕只有三条,写清楚“我如何判断这次抽取成功、这次转换正确、这次加载不重复”,你会少踩无数坑。这也是我这些年做数据项目最深的体会。
