一、我改了三个习惯,都是被半夜的电话逼出来的
2019 年,我在杭州一家做跨境的小公司。凌晨 2 点 47 分,值班同事打电话过来,说 DWD 层那张订单明细表当天只有 37 万行,前一天是 2100 万行,BI 上所有 GMV 全变成 0 了。
我爬起来打开 Airflow(那会儿还是 1.10.9,调度器是单节点,重启一次要 8 分钟,而且重启期间 DAG 全都卡在那里不动),翻了半天日志才发现:任务状态是绿色的,成功的。原因是我前一天把 join 的顺序调了一下,右边那张维表当天上游还没出数,分区是空的。Spark 不会因为 join 到一张空表就报错,它老老实实写出一个空分区,然后下游一路空到底。
后来我养了三个习惯,到现在还在用。第一条,任务成功不等于数据正确。核心表必须挂行数断言,波动超过阈值直接 fail,别让它往下游跑。我一般这么写,放 DWD 任务后面,或者放进 dbt test 里:
select dt, count(1) as cnt
from dwd_order_di
where dt = '${ds}'
group by dt
having cnt < 0.6 * (
select avg(cnt) from (
select dt, count(1) cnt from dwd_order_di
where dt between date_sub('${ds}', 7) and date_sub('${ds}', 1)
group by dt
) t
);
第二条,任何 join 之前先采样看两边 key 的重合率,低于 70% 我就先停下来,八成是口径变了或者上游换了枚举值。第三条,核心 DAG 全加 sla=timedelta(hours=2) 和 on_failure_callback,推到钉钉群。别指望人盯,人盯一定会漏,尤其是连着加班三天以后。
这三条听起来都像废话,但我后来换过两家公司,见过至少五个团队到现在还是靠「报表数不对了业务跑来骂」这个机制发现问题的。
二、200 万个小文件,和那些没人愿意承认的债
小文件这事儿说大不大,但它特别能耗人。当年我们 HDFS 上一个跑了三年的库,光 dt 分区下面就有 200 多万个文件,平均 8MB 一个,最小的 2KB。NameNode 那边每个文件大概占 150 字节元数据,200 万个文件就是 300MB 堆内存,听着不多,但你架不住它天天涨,而且 hdfs fsck 一跑就是四十多分钟,运维看你的眼神都不对了。
问题出在哪?出在 Spark 的 spark.sql.shuffle.partitions,默认 200。一张三千万行的表写 200 个文件,每个文件 15MB 上下,其实还行。但如果你按 dt 跑补数,90 天就是 18000 个文件;要是再叠一层小时分区,24 × 90 = 2160 个分区,乘上去就崩了。当年我就是那个给全表按小时分区的人,现在想起来还脸红。
治理办法我试过几种,从土到洋:最简单的,写完加一步合并任务,用 distribute by cast(rand() * 20 as int) 把输出打散成 20 个文件,代价是多跑一遍;稍微优雅点,Spark 3 里用 spark.sql.files.maxRecordsPerFile 限制单文件行数,配合 repartition 控制文件个数;最省心的是上 Iceberg,写完直接调 rewrite_data_files,我们当时把 90 天的文件从 180 万降到 12 万,同一个查询的平均耗时从 40 秒降到 6 秒左右。要注意的是 Iceberg 的 rewrite 本身也吃资源,最好放低峰期跑,别跟 ETL 抢队列。
顺带说一句,小文件不是技术问题,是排期问题。业务每天都催新需求,合并这种活儿永远排在最后。所以我后来学乖了,在 DAG 末尾加一个 check_file_count 的校验节点,文件数超过阈值就告警,逼着自己去处理。不告警就永远不会有人管。
三、数据倾斜:AQE 能救你 70%,剩下 30% 还得自己动手
数据倾斜我印象最深的一次,是一个 user_id = -1 的脏数据。上游有个埋点没做登录态判断,所有未登录用户都打成了 -1,占了整张表 60% 的量。那个 Spark 任务 200 个 task 里,199 个 3 分钟跑完,剩下 1 个跑了 4 小时 12 分,整个作业卡在最后那个 task 上不动。
排查顺序我一般是这样的:先看 Spark UI 的 Stage 页面,把 task 的 duration 排序,如果最大和最小差 10 倍以上基本就是倾斜了;然后去 SQL 里找 group by 或 join 的 key,采样看 top 10 的 key 占比,超过 30% 就要警惕;最后确认这个 key 是脏数据还是真实业务分布——这两种处理方式完全不一样,前者是清洗,后者是改算法。
具体解法分三层。第一层开 AQE,也就是 spark.sql.adaptive.enabled=true,Spark 3.2 之后这个是默认开的,连同 spark.sql.adaptive.skewJoin.enabled。我的经验是 AQE 大概能解决 70% 的倾斜,尤其是那种分区大小不均的 shuffle。第二层手动加盐,给热 key 拼一个随机后缀打散,聚合完再去掉后缀二次聚合,这个代码写起来丑,但确实有用。第三层是广播,小表小于 spark.sql.autoBroadcastJoinThreshold(默认 10MB)的时候会自动广播,但很多同学不知道这个值可以调,调到 50MB 有时候能直接绕开倾斜。
不过我还想说个反直觉的事:不是所有倾斜都值得修。一个 T+1 的离线任务,最慢 task 跑 40 分钟、整体 2 小时能出数,那就别动它。你花两天优化到 1 小时 50 分,业务根本感知不到,但这两天你本来可以去修那个一直没人管的空分区问题。优先级这东西,得看业务卡在哪。
四、要不要上 Flink?我们把 60 个「实时需求」数了一遍
这个话题我被问过至少三遍,每次都是以「我们业务要实时」开头。
后来我做了一件事:把各部门提的「实时」需求整理成一张表,一共 60 条,然后逐条去问业务方,你希望数据多久延迟。结果挺有意思——真正要求 5 分钟以内的只有 4 条,两条是风控,两条是 C 端库存扣减。剩下 56 条,问到最后业务自己说「其实半小时也行」「小时级够了」「早上 9 点前能看到就行」。
我们最后没上 Flink。方案是 T+1 批加 15 分钟微批,用 Spark Structured Streaming 读 Kafka,trigger(processingTime='15 minutes'),输出到 Iceberg 表。延迟从 24 小时压到 15 分钟,集群成本只涨了 20% 左右(20 台 c5.4xlarge 的 EMR,按月付大概 4.7 万,加微批任务之后变成 5.6 万)。那 4 条真需要秒级的,单独用 Flink 做,2 秒端到端,这部分机器贵,但业务体量小,撑得住。
我不是说 Flink 不好,恰恰相反,它处理状态、checkpoint、watermark 乱序那套东西是很扎实的。我想说的是另一件事:流计算的真实成本从来不是机器,是人。状态后端怎么选、checkpoint 间隔设多少(我们试过 30 秒和 3 分钟,30 秒对小文件特别不友好)、背压了怎么定位、晚到数据怎么处理——这些没有半年养不出手感。一个三个人、还要同时扛离线报表的团队,硬上流计算,大概率是半年后大家都不敢动那个作业。
顺便说个对比:Iceberg、Hudi、Delta 这三个,我实际用过前两个。选型的时候别只看功能表,去看你们最常用的查询引擎版本,Spark 3.1 和 3.5 对同一张 Iceberg 表的表现能差出好几倍。这个坑我们踩过,后来老老实实先把 Spark 升了。
五、真要说什么经验,我觉得是「让别人能复现你的数字」
前面讲的都是技术,最后说点虚的,但我觉得这是我干了这么多年最重要的一条。
数据工程师真正的工作量,我觉得大概只有 30% 在写 SQL 和调 Spark 参数,剩下 70% 花在「为什么财务部的 GMV 和运营部的 GMV 差 12 万」这种事情上。我见过一张宽表字段涨到 380 多个,光是有赞、退款、取消这三状态的口径就有四种写法,谁也不知道哪个对。
所以后来我强制做两件事:一是每个指标在 dbt 或者别的元数据平台里写清楚定义,尤其是过滤条件;二是每张核心表必须能被人从上游原始表一步步复现出来,中间不掺手工改数。这两件事做起来特别枯燥,也出不了什么业绩,但它是唯一能让半年后的自己少加几次班的东西。
我现在也不太相信「数据中台」这种词了。我信的是:一个新人入职三天内,能不能看着文档把昨天的订单数捞出来,并且跟报表对上。能做到,这个团队的数据工程就算及格了。