从 5.5 万到 10 万+: Spark 写 Hive 文件数翻倍 与任务“假卡死”企业案例剖析

发布时间:2026/9/2 7:14:23
从 5.5 万到 10 万+: Spark 写 Hive 文件数翻倍 与任务“假卡死”企业案例剖析 基于 Spark 3.3.4、YARN Cluster、Hive ORC、HDFS Audit 与 Event Log 的跨层诊断。核心结论任务并未真正卡死。计算 Task 完成后Driver 仍在执行 Hive/HDFS Commit将海量 part 文件从 .hive-staging 逐个 rename 到目标分区。输出文件翻倍的直接原因是最终写 Stage 分区数翻倍现有证据排除了 noShowDF 行数增长和 Celeborn 故障但更底层的 input split 变化仍需用 Block/ORC Stripe 与运行时文件状态进一步确认。故障现象数据写完了Application 为什么还不结束某 Spark Application 通过 spark-submit 直接提交到 YARN业务 SQL 的最后一步是 INSERT OVERWRITE Hive 静态分区。Spark UI 中所有计算 Task 已经完成目标分区中也能够看到 part 文件但 Application 长时间保持 RUNNINGDriver 日志每 10 秒只出现一次 ApplicationHeartbeater。Driver 尾部日志与此同时异常批次输出文件由常态约 5.5 万个上涨到 10 万以上另一个异常样本甚至写出了 117,147 个文件。文件数上涨与 Application 延迟结束同时出现首先需要判断二者是否存在因果关系。图 1 正常样本写出约 55,984 个文件Job Commit 约 10.7 分钟图 2 异常样本写出 117,147 个文件Job Commit 增长到 28.7 分钟任务代码为什么 spark.sql.shuffle.partitions8000 没控制住文件数任务先读取上一日全量表通过 BloomFilter 将历史数据拆成 noShowDF 与 showDF再将 showDF 与当天增量合并并去重最后与 noShowDF 做一次 Union 后写入 Hive。核心数据链路如下。核心逻辑精简val noShowDF history_df.filter(row !rnidBF.mightContainString(row.getAs[String](rnid)) || !unidBF.mightContainString(row.getAs[String](unid))) val showDF history_df.filter(row rnidBF.mightContainString(row.getAs[String](rnid)) unidBF.mightContainString(row.getAs[String](unid))) val combinedDF showDF.union(todayDetailDF).dropDuplicates() noShowDF.union(combinedDF).createTempView(result_view) spark.sql(INSERT OVERWRITE TABLE ... SELECT ... FROM result_view)关键机制spark.sql.shuffle.partitions8000 只控制发生 Shuffle 的 Exchange例如 dropDuplicates()。Filter 和普通 Union 都不会主动把上游分区重排成 8000 个。这段代码实际上包含两条物理路径路径 Ahistory_df → noShowDF → final Union。该路径没有 Shuffle历史表的输入分区会一路保留到最终写 Stage。路径 BshowDF today_detail_df → dropDuplicates()。该路径发生 Shuffle分区数受 sql.shuffle.partitions 与 AQE 影响。最终写文件数并不是简单等于 8000而更接近“路径 A 的分区数 路径 B 的最终 Shuffle 分区数”。只要 noShowDF 上游扫描分区上涨最终文件数就会同步上涨。两组 Event Log 对比记录数几乎没变Task 与物理读取量却大幅增加选取 application_1712849500823_0001 作为正常样本、application_1712849500823_0002 作为异常样本对 SQL Execution、Stage 与 Environment 数据进行对比。数据给出的信号非常一致最终业务输出行数和输出字节几乎没有变化但扫描输入字节与 Task 数量明显上涨。也就是说问题发生在“相同数量记录如何被物理读取和切分”而不是业务结果集规模突然翻倍。Stage 12 DAG102,866 个 Task 是怎样形成的Stage 12 同时扫描历史全量表 A 和当天增量表 B 两个 HadoopRDD 经过 MapPartitionsRDD 和 UnionRDD 后进入聚合与 Shuffle Write。UnionRDD 的语义是拼接父 RDD 分区而不是合并父分区。图 3 正常样本 Stage 12约 11.0 TiB 输入Locality 汇总约 4.9 万个 Task图 4 异常样本 Stage 1218.1 TiB 输入102,866 个 Task异常截图中的 Locality Level Summary 可以直接相加Task 数核验Any 1,459 Node local 87,326 Rack local 14,081 102,866 Tasks因此 102,866 不是 UI 展示错误而是 Stage 12 实际生成的输入 Task 数。需要特别注意Stage 12 顶部的 Input Size 是两张表的合计值不能直接拿 18.1 TiB 与某一个历史分区的 HDFS du 结果进行一对一比较。业务字段变化是否导致 noShowDF 增长一度存在这样的候选假设rnid 或 unid 的业务分布发生变化更多历史记录没有命中 BloomFilter从而进入 noShowDF最终导致文件数增加。Event Log 中的 SQL 节点指标可以直接验证这个假设。已排除异常样本的 noShowDF 行数不仅没有增加反而下降约 0.27%。因此“更多业务记录进入 noShowDF 导致文件数翻倍”与事实不符。不过“业务字段内容影响物理大小”仍不能完全排除。例如 pkg、rnid、unid 出现超长值、高熵随机值或基数骤增可能降低 ORC 压缩率、增加物理读取字节和 Memory Spill。它与“noShowDF 行数增长”是两个不同问题必须分开验证。HDFS 现状核验文件数和分区总大小并未翻倍对异常任务读取的历史分区 day20260809 与正常任务读取的 day20260811 执行 hdfs dfs -count 与 hdfs dfs -du当前结果均为约 2.9K 个文件、11.0 TiB 逻辑数据、22.0 TiB 含副本空间。当前 HDFS 统计day20260809 files≈2.9K logical_size11.0T replicated_size22.0T day20260811 files≈2.9K logical_size11.0T replicated_size22.0T这可以排除“当前异常历史分区文件数或总大小直接翻倍”但还不能排除以下情况两个分区的单文件 Block Size 不同例如一个为 128 MiB、另一个为 256 MiB。ORC Stripe 数量、Stripe 大小或 Hive ORC InputFormat 的 split 结果不同。分区在 Application 运行后被重新生成、合并或压缩当前文件状态已不同于运行时。Stage 12 的另一输入——当天 version_full 分区——物理大小或字段压缩率发生变化。结论边界目前已确认的是 input split/Task 数量发生了变化Block Size、ORC Stripe 或具体业务字段是哪一个底层触发因素仍需补充验证不能提前写成已确认根因。为什么最后一个文件 03:15 已生成Application 到 04:55 才结束目标分区中可观察到的最晚 part 文件时间约为 03:15但 HDFS Audit 显示直到 04:55:49Hive 仍在把 .hive-staging 下的 part 文件 rename 到正式分区04:55:51 才删除 staging 目录。图 5 HDFS Audit04:55:49 仍在逐文件 rename04:55:51 删除 .hive-staging跨层时间线03:15 目标目录已出现最后一批 part 文件 03:17—04:55 Driver 仍存活只输出 ApplicationHeartbeater 04:55:49 rename staging/ext-10000/part-* → day20260817/part-* 04:55:51 delete .hive-staging_hive_...Spark SQL 的 INSERT OVERWRITE 只有在 Hive 文件提交和 staging 清理完成后才会返回。代码中的 spark.stop() 位于 spark.sql(full_pkg_data_sql) 之后因此在 Commit 完成前根本不会执行。Driver 持续发送心跳说明 JVM 和 ApplicationMaster 仍然存活并不代表 Spark 还在计算数据。图 6 文件数翻倍与 Application 延迟结束的完整链路Celeborn 是否是根因从现有证据看Celeborn 不是本次 Application 延迟结束的直接原因。Stage 12 的 Shuffle Write 和后续读取已经完成最终写 Task 也已经结束卡住的时间窗口中日志没有 Celeborn push/fetch/revive/fallback 异常而 HDFS Audit 明确显示 Hive 正在执行 rename 和 staging 清理。验证底层触发因素建议按这个顺序执行分别对比两张输入表的异常日与正常日分区避免把 Stage 12 总 Input Size 错归到单张表。history_pathA表的hdfs存储路径 version_pathB表的hdfs存储路径 for p in \ ${history_path}/day20260809 \ ${history_path}/day20260811 \ ${version_path}/day20260810 \ ${version_path}/day20260812 do echo ${p} hdfs dfs -count -q -h ${p} hdfs dfs -du -s -h ${p} done检查历史分区中每个文件记录的 HDFS Block Size而不是只看当前集群 blocksize 默认值。for day in 20260809 20260811 do echo day${day} block size hdfs dfs -stat %o \ ${history_path}/day${day}/part-* | sort -n | uniq -c done估算每个分区的 Block 数并与 Spark Stage 的输入 Task 数交叉验证。hdfs dfs -stat %b %o ${history_path}/day${day}/part-* | awk { files; bytes $1; blocks int(($1 $2 - 1) / $2); } END { printf files%d size%.2f TiB estimated_blocks%d\n, files, bytes/1024/1024/1024/1024, blocks; }如果文件数、大小和 Block Size 都一致检查 ORC Stripe 以及文件修改时间确认是否在任务之后发生过重写。hdfs dfs -ls -t ${history_path}/day20260809 | head hdfs dfs -ls -t ${history_path}/day20260811 | head # 从每个分区选择一个典型 ORC 文件使用 orc-tools meta # 或 hive --orcfiledump 查看 Stripe 数量、Stripe 大小与压缩方式如果某个 version_full 分区物理大小明显异常再验证字段长度、极端值和基数。SELECT day, COUNT(*) AS rows, AVG(LENGTH(rnid)) AS rnid_avg_len, MAX(LENGTH(rnid)) AS rnid_max_len, percentile_approx(LENGTH(rnid), array(0.5,0.9,0.99,0.999)) AS rnid_pct, AVG(LENGTH(unid)) AS unid_avg_len, MAX(LENGTH(unid)) AS unid_max_len, percentile_approx(LENGTH(unid), array(0.5,0.9,0.99,0.999)) AS unid_pct, AVG(LENGTH(pkg)) AS pkg_avg_len, MAX(LENGTH(pkg)) AS pkg_max_len, percentile_approx(LENGTH(pkg), array(0.5,0.9,0.99,0.999)) AS pkg_pct, COUNT(DISTINCT pkg) AS pkg_distinct FROM dm_mid_master.dwd_rnid_pkg_it_version_full WHERE day IN (20260810,20260812) GROUP BY day;短期修复在最终 Union 之后控制输出分区最直接、风险相对可控的措施是在 noShowDF.union(combinedDF) 之后、写 Hive 之前执行 coalesce。这样输出文件数不再直接跟随历史表 input split 数量。val targetPartitions 60000 val resultDF noShowDF .unionByName(combinedDF) .coalesce(targetPartitions) resultDF.createOrReplaceTempView(result_view) spark.sql(fullPkgDataSql)18 TiB 数据写成 60,000 个文件平均约 315 MiB/文件与正常批次约 5.5 万个文件的规模接近。coalesce 不产生全量 Shuffle成本通常明显低于 repartition但它可能让少数输出 Task 偏大需要观察最大 Task 耗时与输出文件分布。**不建议直接使用 repartition(60000)**repartition 会对约 18 TiB 全量数据再做一次 Shuffle。除非 coalesce 后出现严重倾斜否则不应为了控制文件数引入如此昂贵的网络与磁盘开销。中长期优化建议统一上游数据生产参数HDFS Block Size、ORC Stripe Size、压缩编码和写入并发防止不同日期物理布局漂移。评估sql.hive.convertMetastoreOrctrue。在兼容性验证通过后使用 Spark Native ORC使 spark.sql.files.maxPartitionBytes 等文件扫描参数更可控。将输出文件数作为任务级指标记录 history/show/noShow/today/combined/final 的 getNumPartitions、输入字节、输出文件数与 Commit Time。为 Hive Commit 阶段增加可观测性统计 staging 文件数、rename QPS、NameNode RPC 延迟与 Audit 日志中 rename/delete 的持续时间。使用 try/finally 包裹 SparkSession 生命周期保证异常时也执行 stop()但要明确它不能缩短正常的 Hive Commit。val spark SparkSession.builder() .appName(sPkg2VertexStep1$day) .enableHiveSupport() .getOrCreate() try { compute(spark, day, p1day) } finally { spark.stop() }在生产改造前验证fileoutputcommitter.algorithm.version2 对当前 Hive/Spark/CDH 版本的兼容性。它可能减少部分提交开销但无法替代减少文件数。建议新增的运行时诊断日志为避免下一次只能依赖 History Server 反推可在提交 INSERT 前打印各关键 DataFrame 的物理分区数。getNumPartitions 主要读取物理计划分区信息不等同于触发完整数据计算。def logPartitions(name: String, df: DataFrame): Unit { println(s[partition-check] $name${df.rdd.getNumPartitions}) } logPartitions(history, historyDF) logPartitions(noShow, noShowDF) logPartitions(show, showDF) logPartitions(today, todayDetailDF) logPartitions(combined, combinedDF) val finalDF noShowDF.unionByName(combinedDF) logPartitions(final-before-coalesce, finalDF) logPartitions(final-after-coalesce, finalDF.coalesce(60000))根因分层与最终定性结语看到“所有 Task 已完成、文件也已经出现”不能直接把持续心跳理解成 Spark 或 Celeborn 卡死。SQL 写 Hive 在计算之后还可能存在昂贵的文件提交阶段。有效的排查方式是把 SQL Plan、Stage Task、HDFS 文件布局、Hive staging 与 NameNode Audit 串成一条时间线分别回答 Task 为什么变多、文件为什么变多以及 Driver 最终在等待什么。最终建议最终 Union 后增加 coalesce(60000)短期控制文件数和 Commit 压力同时核验四个输入分区的 Block、ORC Stripe 与 mtime长期修复 input split 漂移。