1. 从"报表数字对不上"说起:ETL在数据仓库里的真实位置
那段经历我印象很深:业务方早上九点发来消息,说昨天的成交日报比财务系统少了三十多万,两边都坚称自己没改口径。我们几个人把整条链路从头捋了一遍,最后问题出在最不起眼的地方——抽取环节用的时间戳字段是订单创建时间,而业务改口径之后需要的是支付完成时间,中间那批"跨天支付"的订单就集体漏掉了。这就是ETL最真实的样子:它不像算法模型那样出彩,但只要有一个环节想当然,整条数据链路的可信度就归零。
数据仓库(Data Warehouse)这套体系从上世纪九十年代被系统化提出,核心目标一直没变:把分散在各业务系统里、格式各异、口径不一的原始数据,加工成能支撑分析与决策的一致数据集合。而ETL(Extract-Transform-Load,抽取-转换-加载)就是连接这两端的那条流水线。很多人以为ETL就是写几条SQL把数据搬过去,实际上它同时承担着口径治理、质量兜底、性能保障和可追溯性这四件事,任何一件没做好,报表都会在某个不经意的时刻打你的脸。
这篇文章适合三类人看:刚接手数据仓库建设、被安排去"先把数据同步过来"的工程师;写了半年ETL脚本、开始觉得维护成本压不住的开发;以及需要评估ETL工具和方案的技术负责人。我会把每个环节的取舍逻辑讲透,而不只是贴一段能跑的代码——因为能跑的脚本满地都是,敢让它跑三年的脚本很少。
1.1 一个典型对账现场的复盘
回到开头那个漏单问题。表面看是字段选错,往深挖其实是三个层次的问题叠在一起。第一层是需求理解:业务说的"昨日成交",在订单表里有创建时间、支付时间、完成时间、更新时间四个候选字段,选哪个取决于报表要回答什么问题。第二层是口径固化:这个选择应该被写进数据字典和ETL配置,而不是散落在某个人的SQL注释里。第三层是校验缺失:如果当时有一条"上游订单数 vs 下游汇总订单数"的日校验规则,这个问题在第一次跑批时就该被拦住,而不是等业务方发现。
我后来给自己定了一条规矩:任何ETL任务上线前,必须能回答"如果源端少传了一条,我怎么知道"。回答不上来的,就不算完工。这条规矩拦下的问题比任何代码审查都多。
1.2 ETL和ELT不是同一件事,别混着用
近几年ELT(先加载再转换)因为云数仓的算力弹性变得流行,很多人就把两者当成一回事,这是个危险的简化。ETL的价值在于"数据进仓之前就已经干净且合规",适合强监管、口径复杂、目标存储算力有限的场景;ELT的价值在于"原始数据先落湖落仓,用目标端的算力做转换",适合数据量大、转换逻辑频繁迭代、目标端是云数仓的场景。
判断逻辑其实很朴素:如果你的转换逻辑需要访问多个源系统做关联、且目标端算力紧张,那转换必须前置,走ETL;如果你的转换大量依赖目标端的窗口函数和列存扫描能力,且原始数据有留存价值,那走ELT更划算。我见过最糟的做法是两边都不彻底——在抽取阶段做了一半清洗,又在数仓里重复做一遍,结果两套口径打架,出了问题谁都说不清是哪一层改的。
1.3 数据仓库为什么必须有ETL这一层
有人会问,业务库直接连BI工具查不就行了,为什么要多这一层?答案藏在三个现实约束里。业务库的 schema 是为交易设计的,第三范式、频繁join、字段命名充满缩写,直接查会拖垮线上;业务库的数据质量是按交易需求保证的,空值、脏数据、逻辑删除的记录它不在乎,但分析在乎;业务库只保留当前状态,历史变化它不存,而分析最需要的恰恰是"当时是什么样"。
ETL这一层的本质,是把"面向交易"的数据翻译成"面向分析"的数据。这个翻译过程必须显式、可版本化、可重跑,否则口径就变成了口口相传的玄学。我在做第一个数仓项目时图省事,把一部分清洗逻辑写在了BI工具的查询里,半年后换了BI工具,那部分逻辑直接蒸发,只能凭记忆重建——那次之后我再也没让转换逻辑离开数仓。
2. 抽取环节:全量与增量之间的取舍逻辑
抽取看着最简单,实际是最容易埋雷的地方。核心问题只有一个:每次跑批,我该从源端拿多少数据?拿多了浪费资源和时间,拿少了数据就缺。这个平衡点怎么找,取决于源表的数据特征和你能接受的数据延迟。
2.1 全量抽取的实现与它的隐性成本
全量抽取说白了就是每天把源表整张捞一遍。实现简单到不能再简单:select * from source_table,然后整体覆盖目标表。小表用这个方式完全没问题,几十万行、每天跑一次,几分钟就结束了。
问题出在它会被"顺手复制"到不该用的地方。一张两千万行的订单表,全量抽取每天要拉几十GB网络流量,源库要承受一次全表扫描的IO压力,目标端要做一次全量覆盖写。这三件事叠在一起,轻则跑批时间从半小时涨到三小时,重则把源库的连接数占满,影响线上交易。我见过一次事故,凌晨跑批的ETL任务和早上七点的对账程序撞在一起抢源库连接,导致对账失败,最后是靠给ETL加独立的只读副本才解决的。
全量抽取还有一个更隐蔽的问题:它掩盖了数据丢失。因为每天都是覆盖写,如果某天源端因为网络抖动少传了十万行,目标表当天就"少了"这十万行,第二天全量拉回来又恢复正常,中间那一天的报表就错了,而且没有任何痕迹。所以用全量抽取时,一定要配一条行数波动告警,比如当日行数低于七日均值的90%就触发。
2.2 基于时间戳的增量抽取:字段选错的代价
增量抽取最常见的做法是找一个"更新时间"字段,每次只拉update_time > 上次跑批时间的记录。这个方案对源表有一个硬要求:该字段必须在每次数据变更时被可靠更新。这个要求听起来基础,但踩坑率极高。
我总结过几类典型翻车。第一类是有物理删除的业务:一条记录被删掉了,update_time根本没机会更新,增量抽取永远感知不到它的消失,下游还留着这条幽灵数据。第二类是批量导入绕过ORM:有些中台系统用批量insert加载数据,不会触发应用层的自动更新逻辑,update_time保持原值,这批数据全部漏掉。第三类是时间精度问题:如果字段精度是秒,而跑批间隔也是秒级,同一秒内的边界数据可能被漏掉或重复。
我的做法是,用时间戳增量时一律把水位线往前回退一小段,比如回退五分钟:
-- 每次跑批取水位线前5分钟开始,避免边界丢失 select * from source_table where update_time >= date_sub('${last_watermark}', interval 5 minute) and update_time < '${current_batch_time}'配合下游的去重逻辑,重复比丢失好处理得多。这个"宁可重复不可丢失"的原则,是我做ETL这些年最值钱的一条经验。
2.3 CDC与日志解析:什么时候值得上
当业务要求秒级延迟,或者源表根本没有可靠的update_time时,就得上CDC(Change Data Capture)了。它的思路是直接读数据库的变更日志,把每一条insert/update/delete都捕获下来,以事件流的形式推到下游。好处很明显:没有增量字段的依赖,删除也能捕获,延迟可以做到秒级甚至更低。
代价同样明显。第一是运维复杂度:日志解析组件需要独立部署、监控、处理日志清理导致的位点失效。第二是schema变更的连锁反应:源表加了一列,下游消费程序可能直接崩掉。第三是历史数据的初始快照问题:CDC通常只捕获开启之后的变更,那开启之前的历史数据怎么补?必须单独做一次全量快照,再和增量流做合并,这个合并逻辑本身就容易出错。
我的判断标准是这样:如果业务能接受T+1的延迟,就别上CDC,时间戳增量加回退窗口足够用;如果业务明确要求准实时,并且团队有人力长期维护这套链路,再考虑CDC。介于两者之间的场景,用"微批"(比如每十五分钟跑一次增量)往往是性价比最高的选择。
2.4 抽取环节的水位线必须持久化
最后说一个看起来很小、出问题很致命的细节:水位线存哪儿。新手最容易犯的错是把水位线存在脚本变量或者临时文件里,任务重启后从默认值开始,要么全量重拉,要么从很久以前重跑。
正确做法是把水位线落到一张专门的元数据表里,跑批成功后再更新,更新和业务写入放在同一个事务里(能放就放),保证"数据成功写入"和"水位线推进"要么都成功要么都失败。这样即使任务在中途崩了,下次跑批依然从上次成功的位置继续。
-- 水位线表结构建议 create table etl_watermark ( job_name varchar(128) primary key, last_value varchar(64), last_run_time datetime, status varchar(16) -- running / success / failed );把status加上去,任务启动时先检查上一次是不是running状态,如果是,说明上次异常中断了,需要人工确认或者自动回退重跑。这个小小的状态字段,帮我避免过好几次数据重复写入。
3. 转换环节:数据仓库分层如何约束ETL脚本的写法
转换是ETL里最考验功力的部分,因为它既是技术活也是业务活。数据仓库分层不是为了让架构图好看,每一层都有明确的职责边界,这个边界直接决定了你的转换脚本该写什么、不该写什么。
3.1 ODS、DWD、DWS、ADS四层各自承担什么
ODS层(操作数据存储)几乎不做转换,它的唯一职责是"原样落地",保留源系统的字段和结构,方便回溯。这一层我坚持不加任何业务逻辑,连字段改名都尽量少做,因为一旦ODS掺了业务逻辑,出问题时就失去了"原始参照物"。
DWD层(数据仓库明细层)开始做清洗和规范化:字段命名统一、枚举值翻译、脏数据过滤、逻辑删除的记录剔除。这一层保留最细粒度,一条业务事件一行,不聚合。缓慢变化维的历史追踪通常也在这一层或专门的维度表里完成。
DWS层(数据服务层)做轻度聚合,按主题域组织,比如"用户主题"下按天汇总用户的订单数、金额、活跃天数。这一层的聚合粒度要谨慎设计,粒度太细下游用不上,太粗又满足不了灵活分析。
ADS层(应用数据层)直接面向具体报表和接口,可以为了某个看板做专门的宽表,甚至可以允许一定冗余。这一层变化最频繁,所以逻辑要尽量薄,能引用DWS就别重新算。
我踩过的一个坑是在DWD层就做了聚合,后来业务要看明细,只能从ODS重算一遍,等于白做。分层的铁律是:越底层越细,越上层越合。
3.2 缓慢变化维在ETL里怎么落地
维度建模里最绕的就是缓慢变化维(SCD)。简单说,一个用户的等级从普通升到了VIP,历史订单应该算哪个等级?这不是技术问题,是业务口径问题,必须先定清楚再写代码。
Type 1(直接覆盖)最简单,只保留最新值,适合"填错了要修正"的场景。Type 2(拉链表)保留历史,用start_date、end_date、is_current三个字段标记每个版本的生效区间,适合"要按当时状态分析"的场景。Type 3(加列)只保留有限的几个历史值,用得较少。
拉链表的ETL实现有个固定套路:新数据进来后,先找出发生变化的维度(用业务主键比对最新版本),把旧记录end_date置为变更时间、is_current置为0,再插入新版本。这个过程一定要用一次性的merge操作完成,不能先删后插,否则中间状态会被下游读到。
-- 拉链表更新的核心逻辑(以Hive/Spark SQL为例) -- 1) 关闭旧版本 update dim_user set end_date = '${batch_date}', is_current = 0 where user_id in (select user_id from stg_user_change) and is_current = 1; -- 2) 插入新版本 insert into dim_user select user_id, level, '${batch_date}' as start_date, '9999-12-31' as end_date, 1 as is_current from stg_user_change;实际生产里我会把这两步合并成一个merge into,或者用全量重算的方式保证一致性,因为update在部分计算引擎上性能很差。
3.3 清洗、口径统一与指标原子设计
转换层最容易被忽视的工作是"消歧"。同一个概念在不同源系统里可能有不同的叫法:A系统管它叫status,取值1/2/3;B系统叫state,取值A/B/C。分析要用的时候,必须映射到统一的一套枚举值上,这个映射表要单独维护,不能硬编码在SQL里。
指标的定义更要命。什么叫"活跃用户"?登录算不算?只看不算不算?下单算不算?这些必须在指标字典里写死,然后所有ETL任务都引用同一个定义。我们现在的做法是把原子指标注册到一张元数据表里,每个指标对应一段可复用的SQL片段,脚本通过引用指标ID来拼装,避免同一个指标在十个地方有十种写法。
提示:指标定义一旦被两个以上报表使用,就必须进指标字典;只在一个报表里出现的临时口径,可以先用注释标记,但要在需求评审时明确它的临时属性。
3.4 转换逻辑该放SQL还是放代码
这是个老争论。我的答案是按复杂度分:纯字段映射、过滤、简单聚合,放SQL,可读性和可维护性最好,还能让熟悉业务的分析师参与维护;涉及多层循环、复杂状态机、需要调用外部服务的,放代码(Java/Python/Scala);介于两者之间的,优先拆分——把复杂逻辑拆成几个SQL步骤,而不是硬塞进一个巨大的UDF。
一个反面案例:我曾经维护过一个两千行的SQL,里面嵌了七个嵌套子查询和三个自定义函数,改一个字段要找半小时。后来我把它拆成了五个中间表,每步不超过两百行,加上中间表的落地,虽然多花了存储,但排查问题时能快速定位到具体哪一步错了。宁可多几张中间表,也别写那种没人敢碰的巨型SQL。
4. 加载环节:幂等、分区与重跑机制的设计
加载看起来只是"把结果写出去",但数据仓库的加载和普通写文件完全是两码事,因为它必须能应对重跑、补数、并发这些真实需求。
4.1 覆盖写、追加写、Upsert的适用边界
覆盖写(overwrite)适合全量重算的任务,比如每天的维度快照、按天分区的汇总表。它的好处是天然幂等,不管跑几次结果都一样。追加写(append)适合不可变的事件流,比如日志、埋点,写进去就不再修改。Upsert适合需要更新已有记录的场景,比如维表、状态表。
选择的关键是问自己:这条记录将来会不会被修改?不会,就追加;会,就upsert;如果整个分区都是重新算的,就覆盖。新手最容易犯的错是给应该upsert的表用了追加写,结果同一用户出现多条记录,下游聚合时金额翻倍——这类bug排查起来特别耗时,因为数据看起来"都合理"。
4.2 分区设计与动态分区的实操细节
分区是数据仓库性能的生命线,但分区字段选错,会比不分还惨。分区字段要满足两个条件:查询时经常作为过滤条件,且值的基数不能太高。按天分区是最常见的,按小时分区适合日志类数据,按业务ID分区基本是灾难——几百万个分区,元数据服务直接被拖垮。
动态分区写入时有两个参数必须注意。以Hive为例,hive.exec.dynamic.partition.mode要设为nonstrict,否则必须指定至少一个静态分区;hive.exec.max.dynamic.partitions要调大,否则分区数超过默认值(通常1000)直接报错。Spark里对应的是spark.sql.sources.partitionOverwriteMode,用dynamic时只覆盖涉及的分区,用static会覆盖整表,这个参数设错会误删历史数据,我们内部把它列为"上线前必查项"。
-- Spark动态分区覆盖,只影响本次写入涉及的分区 set spark.sql.sources.partitionOverwriteMode=dynamic; insert overwrite table dws_order_di partition(dt) select ..., dt from dwd_order_detail where dt = '${batch_date}';4.3 幂等性为什么是ETL的生命线
幂等的意思是同一个任务跑一次和跑十次,结果完全一样。它的重要性在于:跑批失败是常态,重跑是必需能力。如果任务不幂等,重跑一次数据就多一份,运维根本不敢重跑,只能人工修数据,这就掉进了无底洞。
实现幂等的三个层次:最彻底的是"先删后写",写之前把目标分区清掉,比如insert overwrite就自带这个语义;次优是"唯一键去重",写入时用row_number()按唯一键取最新一条;最弱的是"什么都不做,靠上游保证不重复",这个基本等于没保证。
我见过最典型的非幂等场景是"累加型"任务:每天把增量加到一张汇总表上。这种写法一旦重跑,昨天的增量就被加了两遍。正确做法是改成按天分区重算全量,或者引入一个批次号,写入前检查该批次是否已处理过。
4.4 补数与重跑的完整链路
数据出问题是必然的,关键是补数要能跑通。我会在任务设计时就预留补数能力:所有任务支持传入日期参数,而不是硬编码当天;DAG依赖支持按日期回溯重跑;补数时自动触发下游重算。
补数最容易忽视的是下游联动。上游补了一天,下游所有依赖它的汇总表都得重算,否则数据不一致。所以血缘关系必须清晰,最好能根据血缘自动生成重跑计划。我曾经手工补过数,漏了一个下游宽表,导致报表连续一周都是错的,从那以后我坚持补数必须走自动化流程。
5. Spark ETL脚本从能跑到能维护的工程化改造
用Spark写ETL,能跑起来和能维护之间隔着一条河。我经手过不少别人留下的脚本,问题往往高度雷同,这里把改造思路完整说一遍。
5.1 一段典型脚本的问题清单
先看一段几乎人人都写过的脚本:
from pyspark.sql import SparkSession spark = SparkSession.builder.appName("order_etl").getOrCreate() df = spark.sql("select * from ods.order where dt='2024-01-01'") df.createOrReplaceTempView("order") result = spark.sql(""" select user_id, count(1) as cnt, sum(amount) as amt from order where status = 1 group by user_id """) result.write.mode("overwrite").saveAsTable("dws.user_order")这段代码能跑,但问题一堆:日期硬编码、没有资源参数、没有分区写入、异常没有处理、表名写死无法复用、没有质量校验。要把它变成生产脚本,需要逐项改造。
5.2 参数化与配置分离
第一件事是把日期、表名、并行度这些全部外置。Spark任务通过--conf或脚本参数传入,代码里只引用变量。更进一步,把任务的配置抽成一份独立文件(json/yaml),代码读配置执行,这样改逻辑不用改代码,一个脚本能跑多个相似任务。
spark-submit \ --master yarn \ --deploy-mode cluster \ --num-executors 20 \ --executor-memory 8g \ --executor-cores 4 \ --conf spark.sql.shuffle.partitions=800 \ order_etl.py --biz_date 2024-01-01参数化的好处在补数时体现得最明显:改一个日期就能重跑任意一天,不用去改代码再提交。
5.3 数据倾斜与分区调优的实际手法
数据倾斜是Spark ETL最头疼的问题,表现为任务卡在最后几个task上,其他都跑完了。根因是某个key的数据量远超其他key,比如某个大客户贡献了80%的订单。
处理手法按场景选:如果是聚合类倾斜,可以先给key加随机前缀打散,聚合后再合并;如果是join类倾斜,小表广播(broadcast)是首选,大表和大表join则可以考虑拆分热key单独处理;如果是null值导致的倾斜,先在join前过滤掉null。
# 热key打散示例 from pyspark.sql.functions import rand, concat, lit df_salted = df.withColumn("salt", (rand() * 10).cast("int")) \ .withColumn("key_salted", concat("user_id", lit("_"), "salt"))spark.sql.shuffle.partitions这个参数也常被忽视。默认200,对于大任务明显不够,导致单分区数据量过大;对于小任务又太多,产生大量小任务调度开销。经验值是总数据量除以每个分区128MB左右,再取个整数。
5.4 小文件与Shuffle的治理
小文件问题在数仓里特别普遍,因为每天往分区写一次数据,时间长了分区里全是几KB的小文件,元数据压力和读取效率都会恶化。治理手段有两个层面:写入时控制分区数,用repartition或coalesce把输出文件数控制在合理范围;写入后用定时任务合并小文件。
Shuffle是最贵的操作,能避免就避免。常见的优化是把多次join合并成一次,把groupBy换成reduceByKey(RDD场景),尽量用宽依赖操作之前先过滤缩小数据量。
# 写入前控制文件数,避免产生大量小文件 result.repartition(10).write.mode("overwrite") \ .partitionBy("dt").saveAsTable("dws.user_order")这里的repartition(10)是相对分区数据的文件数,不是总文件数,实际用的时候要结合单分区数据量估算。
6. 调度、监控与质量校验:上线之后才真正开始
ETL任务上线不是终点,而是运维的开始。一个没有调度、没有监控、没有质量校验的ETL链路,迟早会以你意想不到的方式崩掉。
6.1 DAG依赖与任务编排的常见错误
调度工具的核心是DAG(有向无环图),把任务按依赖关系串起来。听起来简单,但配错依赖的情况太常见了。最典型的是"依赖到了任务级别,但没依赖到数据级别":A任务今天跑的是昨天的分区,B任务依赖A,但B读的也是昨天的分区,结果B在A还没写完最新分区时就启动了,读到的是旧数据。
正确做法是让依赖精确到分区,调度系统要能表达"B任务的分区2024-01-01依赖A任务的分区2024-01-01完成"。另外要避免循环依赖和隐式依赖,所有依赖都要显式声明。我见过最隐蔽的问题是有人靠"任务执行顺序"来保证依赖,没有配置显式依赖关系,某天调度系统并发调度的顺序变了,数据就错了。
6.2 监控告警该怎么设计
监控分三个维度:任务是否成功、跑了多久、产出的数据对不对。任务失败告警是最基础的,但只做这个远远不够。要加运行时长告警,比如平时跑30分钟的任务突然跑了两小时,可能是数据量异常或资源不足,这时候即使没失败也要预警。
产出数据量告警最有价值。按分区记录每日行数,和七日均值对比,波动超过阈值就告警。这个规则帮我抓到过很多问题:上游少传数据、去重逻辑误删、字段变更导致过滤条件失效。
# 简单的行数波动校验 import statistics history = [get_row_count(d) for d in last_7_days] avg = statistics.mean(history) today = get_row_count(today_partition) if today < avg * 0.7 or today > avg * 1.5: send_alert(f"行数异常: 今日{today}, 均值{avg}")6.3 数据质量校验的六类规则
我把校验规则归成六类,覆盖绝大多数场景:完整性(关键字段不为空)、唯一性(主键不重复)、有效性(枚举值在允许范围内)、一致性(上下游行数或金额对齐)、及时性(数据按时产出)、准确性(抽样比对源系统)。
校验的时机也有讲究:源数据落地后做一次,核心转换后做一次,出仓前再做一次。三次校验的成本不高,但能拦住大部分错误数据。落地校验失败的策略要谨慎选择——是阻断整条链路,还是告警但放行?我的建议是核心指标相关的一定要阻断,非核心的可以告警放行,避免因为一个边缘校验规则把整条链路卡死。
6.4 血缘关系与影响分析
血缘是数仓的神经系统。它回答两个问题:这个字段是从哪来的,以及改这个字段会影响谁。前者用于排查问题,后者用于评估变更风险。
血缘的维护成本不低,手工维护几乎不可能长久。常见做法是在ETL脚本里显式声明输入输出表,由调度系统或元数据平台自动解析SQL生成血缘。现在不少平台能解析SQL自动推导字段级血缘,准确率已经够用。
有了血缘,补数、变更、下线这些操作才有依据。我们内部有个硬规定:任何表下线前,必须确认下游影响面为零,或者通知到所有下游负责人。这条规定执行起来麻烦,但比"某天下线一张表导致十个报表集体变空"要好得多。
7. 几个让我加班到凌晨的ETL故障与排查链路
理论讲完了,接下来说几个真实故障。我把排查过程完整写出来,因为排查思路比结论更有价值。
7.1 时区导致的日期错位
现象:日报里每天凌晨0点到8点的数据,总是归到前一天。排查链路是这样的:先确认源库时间和数仓服务器时间是否一致,发现源库用的是UTC,数仓用的是本地时间,两者差8小时;再确认抽取逻辑用的是源库的create_time,这个是UTC时间,直接拿来按本地时间分区,自然就错位了。
修复方式是在抽取时统一做时区转换,把源端时间转成业务约定的时区再入库。更稳妥的做法是在ODS层就存两份时间——原始时间和转换后的业务时间,DWD之后统一用业务时间。这个问题的教训是:时间字段进仓的那一刻就要想清楚时区语义,越晚处理越麻烦。
7.2 上游字段类型变更引发的静默失败
现象:某天开始某张表的金额字段全部变成0,但任务没报错。排查发现上游把金额字段从decimal(10,2)改成了varchar,而我们的转换逻辑有一句cast(amount as decimal(10,2)),遇到带千分位的字符串(比如"1,234.56")转换失败,Spark在这类场景下可能返回null,然后被下游的coalesce逻辑兜底成了0。
修复分两步:短期加上"转换失败计数"的校验,一旦出现非零就告警;长期建立schema变更的感知机制,源端表结构变化时通知ETL负责人。这个坑的可怕之处在于它不报错,数据静默地错,全靠业务方发现。
7.3 重复数据引发的指标翻倍
现象:某天的销售额比前一天翻了一倍,但订单量正常。排查逻辑是:先查订单量,正常,说明不是数据重复导致的整体翻倍;再查销售额的明细分布,发现个别用户金额异常大;最后定位到join环节,维表里同一个用户有多条有效记录(拉链表更新时没关闭旧版本),导致订单记录被放大了。
根因是拉链表的更新逻辑在并发情况下出了竞态,两条记录同时被标记为"当前有效"。修复方式是在拉链表上加唯一约束或跑完后做一次去重校验,确保每个业务主键只有一条is_current=1的记录。这个校验后来成了我们的标配。
7.4 空值语义带来的意外结果
null在SQL里既不是0也不是空字符串,它的行为经常违反直觉。null = null返回的是null而不是true,count(null)不计数,sum(null)返回null,not in (1, null)结果是空集——最后这个坑我踩过,一条where status not in (select status from blacklist)因为子查询里有null,直接返回了空结果,报表突然全空。
规范做法是:所有可能为null的字段在入库时就明确处理,该填默认值的填默认值,该保留null的用is null明确判断,永远不要依赖=或in去和null比较。数据字典里也要标注哪些字段允许null,以及null代表什么含义——是"未知"还是"不适用",这两者的处理方式完全不同。
8. ETL工具选型:从传统商业工具到现代数据栈
最后聊聊工具。市面上的ETL工具从老牌商业套件到开源框架再到云原生服务,选择很多,但选型的逻辑其实不复杂。
8.1 商业ETL工具与自研脚本的对比
商业ETL工具(比如那些带图形化界面的数据处理平台)的优势在于可视化编排、内置的调度和监控、对非技术人员的友好度,以及厂商支持。劣势也很明显:license成本高、复杂逻辑表达受限、性能调优空间小、出问题时要看厂商脸色。
自研脚本(SQL、Spark、Python)的优势是灵活、可控、成本低,任何逻辑都能实现,性能调优空间大。劣势是运维成本高,团队得有相应的工程能力,监控调度都得自己搭,人员流动时知识传承困难。
我的经验是,中小团队、数据量不大、转换逻辑相对标准的场景,商业工具能显著降低人力成本;数据量大、逻辑复杂、团队有工程能力时,自研脚本更划算。最不划算的是拿商业工具做它不擅长的复杂转换,最后写出来一堆难以维护的图形化流程。
8.2 云原生与开源方案的取舍
现在很多团队在开源方案(Spark、Flink、Airflow、DolphinScheduler等)和云原生托管服务之间纠结。开源方案的优势是可控、无厂商锁定、社区活跃、成本相对透明;云托管服务的优势是免运维、弹性伸缩、和周边服务集成好。
判断的核心是团队规模和成本结构。如果团队有专职的数据平台工程师,开源方案能带来更好的长期收益;如果团队小、人力紧张,托管服务把运维成本转移出去,往往更划算。要警惕的是隐性成本:云服务的数据传输费、存储费在数据量大起来之后可能远超预期,选型前一定要算清楚量级。
8.3 一张选型判断清单
我把选型时真正需要确认的问题列成清单,逐条过一遍,基本能得出结论:
| 判断维度 | 关键问题 | 倾向 |
|---|---|---|
| 数据量级 | 日增数据是否超过百GB | 超则优先分布式方案 |
| 延迟要求 | 能否接受T+1 | 能则批处理足够,否则考虑流式 |
| 团队能力 | 是否有专职数据工程人员 | 无则倾向托管/商业工具 |
| 逻辑复杂度 | 转换逻辑是否频繁变化 | 否则脚本化,否则可视化 |
| 成本结构 | 三年总成本(含人力运维) | 都要算,别只看license |
| 合规要求 | 数据是否要求本地化 | 是则自建或私有部署 |
这张表我用过很多次,它能逼着团队把隐性需求想清楚。很多选型失败不是因为工具不好,而是因为用错了场景——拿批处理工具硬做实时,拿可视化工具硬做复杂状态机,最后都会付出代价。
选型定下来之后,还有一件事常被忽略:迁移成本。从旧方案迁到新方案的代价,往往比新方案本身的价值还大。如果现有方案不是完全不能用,先优化而不是先替换,通常是更明智的选择。我在一次迁移上吃过亏,花了三个月把三十个任务迁到新框架,结果新框架在某类场景下性能还不如旧的,最后又迁回来一部分——那三个月成了我职业生涯里最不想回顾的阶段之一。