ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

数据仓库ETL全链路实战:工具选型、增量拉链与幂等治理

数据仓库ETL全链路实战:工具选型、增量拉链与幂等治理 很多人对数据仓库的理解停留在建几张表、写几条SQL的层面直到真正接手一条从业务库到报表的完整链路才发现最耗时间的从来不是建模而是中间那段 ETL。数据仓库的成败八成取决于 ETL 这一环做得稳不稳源库加了字段怎么办任务失败重跑会不会把数据搞重实时链路延迟上去了怎么定位。这些问题没有标准答案只有踩过坑之后的经验。这篇内容面向三类读者刚转数仓方向、需要一套完整认知框架的同学正在做数据同步与加工、想把链路做扎实的工程师以及准备面试、需要把零散知识点串成体系的人。我会把常用 ETL 工具、ETL 方法以及字节跳动、京东、美团、腾讯这类公司在面试里反复问的点一次性讲透并且每一块都落到能直接抄的配置和写法上。1. ETL在数据仓库里到底承担了什么职责1.1 从ODS到ADS数据是怎么被搬来搬去的先说说分层这件事。绝大多数公司的数仓都逃不开 ODS、DWD、DWS、ADS 这套结构命名可能不同但逻辑是一致的。ODS 层贴源存放几乎不做加工目的是保留一份原始现场将来出了问题能回到源头查DWD 层做清洗和规范化把散落在不同业务库的字段统一口径比如订单表的金额单位统一成分时间统一成标准时间戳DWS 层按主题聚合做宽表和轻度汇总比如用户日粒度行为汇总ADS 层直接面向报表和接口字段可读性优先。ETL 在这四层之间来回穿梭。从 ODS 到 DWD 是一次抽取加转换从 DWD 到 DWS 往往是一次带聚合的加载从 DWS 到 ADS 可能是轻量加工甚至直接映射。你如果只把 ETL 理解成从A库同步到B库那就把这件事看小了——它其实是整条数据链路的血管负责运输也负责在运输过程中完成过滤、转换、校验和补全。我见过不少团队把 ETL 写成一个个孤立的脚本谁需要谁写一个半年之后没人敢动。真正健康的做法是每一段的输入、输出、依赖、校验规则都写清楚脚本本身就是文档。1.2 ETL和ELT差的不是字母顺序这两个词现在被混用得厉害但它们的取舍直接决定架构。ETL 是 Extract-Transform-Load先转换再加载转换发生在数据到达目标仓库之前通常跑在独立的计算引擎上Spark、Flink 都算。ELT 是 Extract-Load-Transform先把原始数据整个灌进目标仓库再用仓库自身的算力做转换典型代表是 Snowflake、BigQuery、Doris、ClickHouse 这类。差别在哪ETL 模式对中间的转换服务器压力大但目标仓库干净ELT 模式依赖仓库的算力弹性省了一套中间层但要求仓库本身足够强。传统 Hadoop 体系里 ETL 是主流因为 HDFS 只负责存计算得靠 MapReduce 或 Spark 干。而现在很多团队往 ELT 靠原因很实在少维护一套东西数据全程可追溯。我的建议是别纠结名词。判断标准就一条你的转换逻辑跑在哪、跑完的数据要不要在中途落盘。如果要落盘做检查点那就是 ETL 思路如果全程在目标库内部完成那就是 ELT 思路。实际项目里两者常常混着用ODS 到 DWD 用 ETLDWD 到 DWS 用 ELT完全没问题。1.3 一条最小可用的ETL链路长什么样抛开所有复杂场景一条能跑起来的链路只需要四样东西一个可靠的数据抽取方式、一个能承载转换的计算引擎、一个存放结果的存储、一个负责任务编排的调度器。举个具体例子。源库是 MySQL 的订单表目标是 Hive 的 DWD 层明细表。抽取用 DataX 配一个 JSON 任务每天凌晨拉全量或按更新时间拉增量转换用 Spark SQL 做字段清洗和关联结果写成 Hive 分区表调度交给 DolphinScheduler配好上游依赖和数据质量检查节点。这四步听着简单但每一步都有讲究。抽取要考虑源库压力不能让同步任务把线上业务拖垮转换要考虑数据倾斜不能让某个 key 拖慢整个任务存储要考虑分区策略不能全表扫调度要考虑失败重试和幂等不能让重跑变成一场灾难。下面几节就把这些点逐个拆开。2. 常用ETL工具盘点什么场景该用什么2.1 批量数据同步DataX、Sqoop、Kettle各自的舒适区批量同步是 ETL 里最基础也最高频的一环工具选错了后面全是补丁。DataX 是阿里开源的离线同步工具单机部署靠多线程和 channel 提高吞吐。它的优势是插件生态完整MySQL、Oracle、HDFS、Hive、ClickHouse 基本都有现成的 reader 和 writer配置文件是 JSON改起来直观。它有个很实用的参数叫speed可以限制字节速率和记录数速率这一点在做线上库同步时是救命的功能——你可以把限速压到每秒几千条让同步任务对源库几乎无感。Sqoop 是老一代 Hadoop 生态的同步工具本质是把同步任务翻译成 MapReduce 作业所以它的并发能力来自 MR 的并行度。它在 HDFS 和关系库之间搬运很成熟但启动开销大小任务跑一次光 JVM 启动就要几十秒而且对非 Hadoop 目标的支持一般。现在新项目里用 Sqoop 的越来越少了但存量系统里还很常见面试也可能问到。Kettle 是图形化的 ETL 工具拖拽式配置转换Transformation和作业Job两级结构非技术人员也能上手。它的弱点是性能调优空间有限处理千万级以上数据时容易成为瓶颈而且图形化配置在版本管理上很别扭——XML 文件 diff 起来一塌糊涂。适合中小规模、逻辑复杂但数据量不大的场景。工具部署形态优势明显短板适用场景DataX单机多线程插件全、限速细、配置直观单机吞吐有上限关系库到数仓的离线同步Sqoop依赖 Hadoop与 HDFS 生态贴合启动开销大、目标受限存量 Hadoop 体系Kettle图形化上手快、逻辑可视化大数据量性能弱、难做版本控制中小规模、复杂转换Flink CDC分布式流式支持全增量一体、Exactly-Once资源占用高、运维复杂度大实时同步、CDC 场景选型的时候我一般问自己三个问题数据量多大、延迟要求多高、源库能不能承受压力。三个答案出来工具基本就定了。2.2 日志与数据库变更采集Flume、Canal和CDC的分工批量同步解决的是已有的数据但日志和数据库变更这两类数据源需要的是采集能力。Flume 主攻日志采集Agent 由 Source、Channel、Sink 三部分组成。生产环境里最常用的是 Taildir Source它能监控多个目录下的文件并且支持断点续传——把读取位置记录在一个 position 文件里进程重启后从上次的位置继续读。Channel 一般选 File Channel 而不是 Memory Channel因为 Memory Channel 一旦 Agent 挂了数据就丢了File Channel 虽然慢一点但能保证不丢。Sink 可以落到 HDFS 或 Kafka。Canal 解决的是 MySQL 增量同步问题原理是伪装成 MySQL 的从库向主库发送 dump 协议请求拿到 binlog 之后解析成结构化事件。它的价值在于不侵入业务代码业务方不用改任何东西DBA 只需要开一下 binlog 权限。需要注意的是 Canal 依赖 binlog 格式必须是 ROW 模式否则拿不到变更前后的完整字段值。Flink CDC 算是新一代方案内置了 Debezium支持全量加增量一体化读取而且能做到 Exactly-Once 语义。它最大的好处是把先全量同步、再切换增量这个容易出错的步骤自动化了不用再手动记录 binlog 位点。代价是资源占用和运维复杂度都上去了一个 Flink 作业要占好几个 TaskManager小团队未必玩得转。选型上我的经验是日志类走高吞吐的采集链路Flume 到 Kafka 是最稳的组合数据库变更如果只是 T1 需求用 DataX 按更新时间拉增量就够了没必要上 CDC只有真正需要秒级同步的场景才值得为 Canal 或 Flink CDC 付出运维成本。2.3 计算与转换引擎Hive、Spark、Flink怎么分工转换这一步跑在什么引擎上直接决定了你能处理的规模和延迟。Hive 是最经典的批处理引擎本质是把 SQL 翻译成 MapReduce 或 Tez、Spark 作业。它的优势是稳定、生态成熟、几乎所有大数据平台都支持缺点是延迟高一个中等复杂度的查询跑几分钟很正常。T1 的离线数仓用 Hive 完全够没必要为了快一点去引入更复杂的东西。Spark 的定位是内存计算比 Hive 快一个数量级是常见现象尤其是需要多次迭代的场景比如机器学习特征加工。Spark SQL 的语法和 Hive 高度兼容迁移成本低。它的核心概念是 RDD、DataFrame、Dataset 三层抽象做 ETL 主要用 DataFrame 和 Spark SQL。有一类 ETL 脚本用 PySpark 写灵活度高但要注意 Python 和 JVM 之间的序列化开销能用 Spark SQL 表达的尽量别用 UDF。Flink 走的是流处理路线把批也当成流的一种特例所以能做到真正的流批一体。实时数仓场景基本是 Flink 的天下Kafka 进、Flink 算、结果写进 Doris 或 ClickHouse。它的状态管理和 Checkpoint 机制是核心竞争力能保证故障恢复后数据不重不丢。代价是开发门槛比 SQL 高需要理解水位线、状态、窗口这些概念。一个常见的分工是ODS 到 DWD 用 Spark 做清洗数据量大、逻辑重DWD 到 DWS 的聚合用 Hive逻辑简单、可容忍延迟实时指标用 Flink 单独走一条链路。三条链路并存不丢人硬凑成一条才丢人。2.4 调度编排Airflow和DolphinScheduler怎么选再好的 ETL 脚本没有调度就是一堆散落的代码。Airflow 是 Python 写的调度平台核心概念是 DAG有向无环图每个任务是一个 Operator依赖关系用代码定义。它的优势是灵活Python 能表达的逻辑它都能表达写自定义 Operator 很方便社区生态也大。缺点是对非技术人员不友好一个 DAG 文件几百行是常态而且它的调度器在高并发下需要额外调优原生对 HA 的支持也是后来才补齐的。DolphinScheduler 是国产的可视化调度系统拖拽定义工作流支持多租户、补数、超时告警、失败重试对国内团队的使用习惯适配得很好。它的补数功能特别实用某个任务失败要重跑历史几天的数据在界面上选个日期范围就行不用手写脚本。缺点是复杂逻辑的表达能力不如 Airflow 代码化灵活。选的时候看团队构成如果团队里 Python 工程能力强的多选 Airflow如果业务和数仓人员占多数、需要可视化操作选 DolphinScheduler。两个都支持分布式部署都够用。提示无论用哪个调度器都要把任务幂等作为强制规范。调度器能帮你重试但重试是否安全取决于你的脚本这个责任在开发者身上。3. ETL方法增量、拉链、幂等这三件事决定成败3.1 全量与增量的取舍以及增量字段怎么挑全量同步简单粗暴每次把源表整个拉一遍。它的好处是永远不会漏数据坏处是数据量大、对源库压力大、耗时随数据量线性增长。几百万行以内的维表用全量完全可以接受上千万行的明细表就必须考虑增量了。增量同步的关键在于找到一个可靠的增量字段。常见的候选有三个自增主键、更新时间戳、binlog 位点。自增主键的问题是物理删除的数据拉不到而且并发插入时可能出现空洞不能作为唯一依据。更新时间戳的问题是有些业务表的 update_time 不更新或者用数据库的 CURRENT_TIMESTAMP而数据库服务器时间有偏差导致边界数据被漏掉。binlog 位点最准确但需要 CDC 工具支持。实际做法通常是组合使用以更新时间戳为主往前冗余一段时间窗口。比如每天同步时条件是update_time 昨天00:00:00 and update_time 今天00:00:00为了避免时间偏差把左边界往前推一小时改成update_time 昨天23:00:00。多拉的一小时数据在写入时用主键覆盖不会产生重复。还有个细节增量字段要建索引。如果源表在 update_time 上没有索引每次同步都会全表扫源库压力极大DBA 迟早找你谈话。3.2 缓慢变化维与拉链表到底该怎么实现缓慢变化维是数仓面试的保留题目。维度变化有几种处理方式类型1是直接覆盖历史不留痕类型2是加一行新记录用起止时间标记有效期也就是拉链表类型3是加一列存上一个值用得少。拉链表是类型2的典型实现核心是三个字段业务主键、开始时间、结束时间。结束时间用一个极大值比如9999-12-31表示当前有效。当一条记录发生变更时做两件事把旧记录的结束时间改成变更日期插入一条新记录开始时间是变更日期结束时间是极大值。用 SQL 实现的话思路是先把增量数据和历史数据做全外连接找出新增和变更两类再对历史中未变更的部分原样保留。伪代码大致是这样-- 1. 找出发生变更的记录历史有效记录与增量记录主键相同但属性不同 -- 2. 历史有效记录的结束时间更新为业务日期 -- 3. 新增记录和变更后的新记录追加进来结束时间为极大值 INSERT OVERWRITE TABLE dwd_user_zip PARTITION (dt ${bizdate}) SELECT user_id, name, phone, start_date, CASE WHEN t1.user_id IS NOT NULL AND t2.user_id IS NOT NULL THEN ${bizdate} ELSE end_date END AS end_date FROM ( -- 历史中未变更的记录 需要关闭的旧记录 SELECT ... FROM dwd_user_zip WHERE dt ${bizdate-1} ) t1 FULL OUTER JOIN ( SELECT ... FROM ods_user WHERE dt ${bizdate} ) t2 ON t1.user_id t2.user_id;坑在哪第一跨分区更新很麻烦所以拉链表通常按天分区每天一个全量快照本质是用空间换逻辑简单。第二属性比对的 NULL 值判断要小心NULL ! NULL在 SQL 里是 true会导致大量假变更得用coalesce兜底。第三首次初始化拉链表时要保证历史全量数据都有正确的起始日期否则回溯查询会出问题。3.3 重跑幂等ETL脚本的生命线任务失败重跑是家常便饭但重跑安不安全取决于你的写入方式。有三种常见的写入模式第一种是追加模式直接 INSERT INTO重跑必然产生重复数据。除非上游有全局唯一约束否则不要用。第二种是覆盖模式INSERT OVERWRITE 按分区覆盖。这是离线数仓最常用的方式只要分区粒度选对了重跑就是安全的。选分区粒度的时候要想清楚按天分区就按天覆盖如果一天内跑了多次后一次会覆盖前一次这正是我们要的。第三种是 upsert 模式按主键 merge适合目标端支持事务的存储比如 Doris、Hudi、Iceberg。这种模式对实时链路更友好但要注意小文件问题频繁 merge 会产生大量小文件。我的规范是离线任务一律按分区覆盖实时任务一律按主键 upsert禁止裸 INSERT INTO。这条规范执行下去能省掉后面一半的数据修复工作。注意分区覆盖有个隐藏陷阱——如果某天的源数据为空覆盖后分区就是空的。这本身没错但下游任务如果是内连接就会把所有数据过滤掉。所以空分区要有专门的告警不能让它悄悄传下去。3.4 数据质量校验该插在链路的哪一环数据质量不是最后做一次统计就完事的应该像单元测试一样嵌在每一环里。我习惯在每个关键节点后面加校验任务校验维度主要有五类完整性主键不为空、行数不为零、唯一性主键不重复、一致性关联表的引用能对上、准确性金额不为负、时间不在未来、及时性分区数据在 SLA 时间内产出。具体做法是写一个通用的校验 SQL 模板把规则配置化。比如这样一条规则校验今天的订单明细行数与昨天相比波动不超过 20%。SELECT CASE WHEN ABS(cnt_today - cnt_yesterday) / cnt_yesterday 0.2 THEN 1 ELSE 0 END AS is_abnormal FROM ( SELECT (SELECT COUNT(*) FROM dwd_order WHERE dt ${bizdate}) AS cnt_today, (SELECT COUNT(*) FROM dwd_order WHERE dt ${bizdate-1}) AS cnt_yesterday ) t;这个阈值怎么定别拍脑袋。先跑一个月的历史数据算出每天的行数波动区间把阈值设成历史最大波动的 1.5 倍左右。设太严会天天误报设太松就失去意义。校验失败要阻断下游而不是发个告警就完事——这一点很多团队做不到结果就是脏数据一路流到报表业务方先发现再回头找数据团队。4. 我在真实项目里踩过的六个ETL坑4.1 小文件从几百个变几十万只用了三个月HDFS 上每个文件对应一个元数据条目NameNode 的内存是有限的。当一个分区下有几十万个小文件时不只是查询变慢整个集群的元数据操作都会受影响。小文件怎么来的最常见的三个来源流式写入频率太高每分钟落一次盘、分区粒度太细、任务的并行度与数据量不匹配。前两个是设计问题第三个是参数问题。解决办法分两层。写入侧控制流式写入改成按大小或时间滚动比如每 128MB 或每 10 分钟生成一个文件Spark 任务输出前用repartition或coalesce把分区数降到合理范围经验值是每个输出文件 128MB 到 256MB。存量治理写一个定时任务扫描小文件过多的分区用INSERT OVERWRITE重新合并一次。注意这个任务本身也要限流不然一次合并几百个分区会把集群打满。# Hive 侧合并小文件的参数 SET hive.merge.mapfiles true; SET hive.merge.mapredfiles true; SET hive.merge.size.per.task 268435456; -- 256MB SET hive.merge.smallfiles.avgsize 134217728; -- 128MB4.2 数据倾斜不是所有热点都能靠加盐解决数据倾斜的表现是一个 Spark 任务 199 个 task 三分钟跑完最后一个跑了两小时。原因通常是某个 key 的记录数远超其他 key比如大促期间某个头部商家的订单量是普通商家的上万倍。排查办法是先看 stage 的 task 耗时分布找出最慢的那个 task 的输入量再反推是哪个 key。定位到 key 之后再选方案。如果热点 key 是 NULL 或空字符串造成的最简单的办法是给 NULL 值加上随机后缀打散到不同分区反正这些数据后续也用不上。如果是真实的业务热点可以用两阶段聚合先给 key 加随机前缀做局部聚合去掉前缀后再做全局聚合。还有一种情况是 join 引起的倾斜。如果其中一张表足够小比如几 MB直接开 map join 广播出去倾斜自然消失SET hive.auto.convert.join true; SET hive.mapjoin.smalltable.filesize 25000000; -- 25MB 以下自动广播我踩过的一个坑是某个 key 的倾斜其实是因为上游数据质量问题同一个用户 ID 被错误地写成了不同的格式带空格和不带空格导致 join 时匹配不到产生了巨量的笛卡尔积。所以遇到倾斜先别急着调参先看看数据本身对不对。4.3 时区、NULL和数据类型漂移这三件小事这三件事单独看不起眼凑一起能让数据对不上账。时区问题的典型场景是源库用的是 UTC 时间数仓用的是东八区同步时不做转换结果每天凌晨那几个小时的订单全跑到前一天去了。解决办法是在 ODS 层就统一转换别留到下游。转换要用带时区的函数比如from_utc_timestamp(ts, Asia/Shanghai)而不是简单加减八小时——夏令时和跨年的时候硬加时间会出错。NULL 值的坑更隐蔽。SQL 里NULL NULL返回的是 NULL 而不是 trueCOUNT(column)会跳过 NULLSUM遇到全 NULL 返回 NULL 而不是 0。做聚合的时候要用COALESCE(SUM(x), 0)。另外 group by 之后 NULL 会归成一组如果不想让它们混在一起得先转成特殊值。数据类型漂移发生在源库改结构的时候。比如某个字段从 int 改成 bigint或者从 varchar(20) 改成 varchar(50)同步任务的 schema 如果不跟着变轻则截断重则报错。我的做法是同步任务里尽量用宽类型string 或 decimal 大精度在 DWD 层再做类型收敛。多花一点存储省下无数排查时间。问题类型典型表现处理位置处理方式时区偏移凌晨数据归属错误ODS 层带时区函数统一转东八区NULL 传播汇总结果为空DWD/聚合层COALESCE 兜底注意 count 语义类型漂移字段截断或任务报错同步配置宽类型落地DWD 层收敛编码不一致中文乱码ODS 层统一 UTF-8源端确认5. 大厂数仓ETL面试题拆解5.1 高频考点数据倾斜的处理套路怎么答面试官问数据倾斜想听的不是加盐两个字而是你有没有完整的排查和决策链路。我一般这么答先确认是哪种倾斜——是 group by 倾斜、join 倾斜还是 count distinct 倾斜三者的解法不同。group by 倾斜用两阶段聚合join 倾斜分大表对小表、大表对大表两种情况小表直接广播大表对大表可以用 map join 加倾斜 key 单独处理。count distinct 倾斜比较特殊可以先 group by 加随机数去重再在外层精确去重或者用 bitmap 做近似计算。然后补一句实战经验加盐不是万能的加完盐之后如果热点数据本身占比极高打散效果也有限这时候更该考虑的是业务上能不能拆分比如把大客户单独走一条链路。5.2 高频考点拉链表的实现和时间维度设计拉链表考的是你对缓慢变化维的理解深度。常见追问有三层拉链表怎么更新、拉链表怎么查询某个历史时点的状态、拉链表怎么处理删除。更新逻辑我在前面讲过核心是全外连接加两段式处理。历史时点查询要注意用start_date 2024-01-01 and end_date 2024-01-01这个条件不能写成between因为结束时间是开区间。删除的处理有两种物理删除的记录可以给一个已删除标记或者在拉链表里关闭这条记录且不新增。我倾向后者语义更干净。面试官还喜欢问拉链表和全量快照的取舍。全量快照每天存一份查询简单但存储成本高拉链表省空间但查询要带时间条件而且更新逻辑复杂。数据量小的维表用全量快照完全合理别为了炫技硬上拉链。5.3 高频考点实时数仓的选型与Lambda、Kappa之争实时数仓的面试题基本围绕两条路线Lambda 架构和 Kappa 架构。Lambda 是批流两条链路并行批处理保证准确性流处理保证实时性最后在服务层做合并。它的好处是稳实时链路出问题还有离线兜底坏处是同一份逻辑写两遍两边的口径很容易不一致维护成本高。Kappa 是只保留流处理一条链路历史数据用流的方式重放。逻辑只有一份口径天然统一但对消息中间件的存储时长和流处理引擎的能力要求很高。现在很多公司走的是折中方案实时链路只处理增量历史数据从离线导入两者在存储层做合并。选型时我关注三个问题业务能接受多长的延迟、实时链路的故障恢复代价有多大、团队有没有维护流任务的能力。三个问题回答完方案基本就出来了。5.4 高频考点数据质量保障和任务治理怎么落地这一类问题考的是工程意识。面试官想知道你有没有从写脚本升级到管体系。我一般分四层讲。规则层定义校验标准覆盖完整性、唯一性、一致性、准确性、及时性五类。执行层把校验嵌到调度流程里关键节点前必须通过校验才放行。响应层定义告警分级P0 直接电话P1 发消息并且要明确谁负责处理。治理层定期做表和任务的盘点下线无人使用的表合并重复的任务控制存储成本。有个容易被忽略的点是血缘。没有血缘出了问题是靠人肉问出来的有了血缘从出问题的报表往上追几分钟就能定位到源头。所以就算团队小也要把血缘采集做起来很多调度器自带这个功能配置一下就能用。6. 把ETL链路落地成可维护工程的清单写到这里我把前面散落的东西收一收给一份我自己在用的落地检查清单。这份清单不是理论是踩过坑之后一条条加上去的。第一所有同步任务必须有明确的限速配置。源库是生产库同步任务再急也不能把它拖垮。限速值一开始设保守一点观察一周再往上调。第二所有离线写入必须按分区覆盖禁止裸 INSERT。这条规则要在代码评审里卡住不能靠自觉。第三所有任务必须能重复执行。判断标准很简单同一个日期连跑三次结果应该完全一样。做不到就说明设计有问题。第四每个关键节点后必须有数据质量校验校验失败要阻断下游。校验规则要配置化新增规则不应该改代码。第五每个任务都要有 owner没有 owner 的任务三个月后自动下线。这条规则看着激进但能有效控制技术债。第六监控要分层。任务层面看成功率和耗时数据层面看行数波动和空值率资源层面看集群水位。三层都覆盖了大部分问题能在业务发现之前暴露出来。第七文档和脚本放在一起。每个目录下写清这个任务的目标、上游依赖、下游影响、常见故障处理方式。半年后接手的人会感谢你。最后说一个我自己体会最深的地方ETL 的价值不在于把数据搬过去而在于让搬过去的数据可信。工具选型、参数调优、架构设计所有这些最终都服务于一件事——业务方看到报表上的数字不需要再问一句这个数准吗。做到这一点链路才算真的做完了。
返回列表