Spark作业慢?四层瓶颈诊断法:I/O、Shuffle、Python、Memory
1. 项目概述当 Spark 作业从“分钟级”滑向“小时级”你该信什么、看什么、改什么我带过六支数据平台团队亲手调优过超过 230 个生产级 Spark 作业——其中 87% 的“慢得离谱”问题根本不是配置参数没调对而是诊断路径错了。你有没有试过把spark.sql.shuffle.partitions从 200 改成 400再改成 100又加了.cache()甚至开了adaptive query execution结果作业还是卡在 Stage 3 耗时 47 分钟UI 上 Stage Duration 图像像心电图一样平缓拉长但 Executor Metrics 里 GC 时间只占 2.3%Shuffle Read/Write 量看着也“合理”……这时候你信 UI信日志信同事说的“肯定是数据倾斜”不你该信的是可复现、可隔离、可归因的硬测量数据。这篇文章讲的不是“Spark 调优十大技巧”那种泛泛而谈而是一套我在金融风控实时特征计算、电商用户行为宽表构建、IoT 设备时序聚合等真实高负载场景中反复验证过的故障定位流水线。它聚焦四个最常被误判、却最容易被精准识别的瓶颈层I/O 层磁盘/网络读写效率、Shuffle 层数据重分布的真实开销、Python 层PySpark UDF / Pandas UDF 引入的隐性惩罚、Memory 层JVM 堆内对象生命周期与 Off-Heap 内存竞争。全文没有一行“理论上应该……”只有我在某次凌晨三点排查一个拖了 5 小时的订单反欺诈模型特征生成任务时用jstack抓到的UnsafeRowSerializer死锁线程栈有我在某次用perf对比 PyArrow 与原生 ParquetReader 的 CPU Cache Miss Rate 时发现的 3.8 倍差异还有我把pyspark.sql.functions.udf全部替换成pandas_udf后单个 Executor 的 Python 进程 RSS 内存下降 62% 的实测截图。它适合三类人刚接手慢作业、急需交付的工程师想建立标准化诊断 SOP 的 Tech Lead以及被“调参玄学”折磨太久、渴望回归工程确定性的资深开发者。提示本文所有方法均基于 Spark 3.3LTS 版本适配 YARN/K8s 集群模式。本地spark-shell或pyspark交互式调试仅作辅助验证绝不作为性能判断依据——本地模式绕过了网络 Shuffle、缺少真实 Executor 资源竞争、内存模型也与集群不同用它测出的“快”上线后大概率是假象。2. 整体诊断思路为什么必须按 I/O → Shuffle → Python → Memory 的顺序排查很多人一上来就盯着 Spark UI 的 Stage 列表看到某个 Stage Duration 高就猛点进去看 Task 列表再根据 Task Duration 分布判断“是不是数据倾斜”。这就像医生不查血常规、不量血压光看病人脸色发白就说“贫血”然后直接开补铁剂——可能治对但更可能掩盖真正病因。Spark 作业的执行流是严格分层的数据必须先读进来I/O再按逻辑切分重组Shuffle过程中若涉及 Python 计算则需跨 JVM/CPython 边界Python所有中间状态都依赖内存管理Memory。这四层存在强依赖关系I/O 慢会导致后续所有阶段“等数据”Shuffle 不稳会放大 Python 序列化开销而 Memory 不足又会让 GC 频繁打断 I/O 和 Shuffle 流水线。跳过前置层直接优化后置层等于在漏水的水管上刷漆——表面光鲜问题照旧。我见过最典型的误判案例发生在一个日处理 12TB 用户点击日志的作业上。团队花了三天时间优化 Shuffle调大spark.sql.adaptive.enabled加coalesce强制减少分区数甚至重写了自定义 Partitioner。结果作业从 112 分钟降到 108 分钟。直到我要求他们跑一次iostat -x 1 30在 Driver 节点和两个随机 Executor 节点上同时采集才发现磁盘await平均等待 I/O 完成时间峰值达 180ms%util长期 99%。根源是 HDFS 客户端配置了dfs.client.read.shortcircuit但底层存储节点未启用短路读所有读请求都绕道 DataNode TCP 栈网络延迟叠加磁盘排队把整个流水线拖垮。修复只需在hdfs-site.xml加两行配置并重启 DataNode作业直降为 23 分钟。这个案例说明I/O 是整条链路的“地基”地基不稳上层所有优化都是空中楼阁。为什么不是按“出现频率”排序因为高频问题如常见 Shuffle 参数误配往往症状明显、易定位而真正难缠的慢作业恰恰是那些“症状模糊”的——Task Duration 分布均匀、GC 时间正常、Shuffle Bytes 量不大但就是慢。这类问题90% 出现在 I/O 或 Python 层。比如用spark.read.json()读取嵌套很深的 JSON解析开销全在 Driver 端单线程完成Executor 根本没干活UI 上却显示“Stage 0: 1000 Tasks, Avg Duration 12ms”——这 12ms 是每个 Task 启动后等待 Driver 分发解析完的数据块的时间本质是 I/O 解析瓶颈不是计算瓶颈。再比如用pandas_udf处理小批量数据时如果pandas版本低于 1.4其内部BlockManager在小 DataFrame 场景下会触发大量memcpyCPU 利用率飙升但有效计算极少表现为单个 Task Duration 长且波动大但Shuffle Read为 0——这是典型的 Python 层隐性开销。注意这个顺序不是教条而是风险收益比最优路径。I/O 层诊断工具链最成熟iostat/dstat/hadoop fs -du、耗时最短单次采集2分钟、修复成本最低配置调整或存储格式切换而 Memory 层深入 JVM GC 日志分析需要jstat/jmap/async-profiler多工具联动耗时长、解读门槛高应作为最后手段。把 80% 的问题拦截在前两步才能让团队精力聚焦在真正值得深挖的 20% 上。3. 核心瓶颈层深度解析与实操要点3.1 I/O 层识别“假快真慢”的磁盘与网络陷阱I/O 瓶颈最狡猾的地方在于它常常伪装成“计算慢”。Spark UI 的 Task Duration 统计的是从 Task 启动到结束的总耗时包含 JVM 初始化、序列化、网络传输、磁盘读取、实际计算等多个环节。当磁盘或网络成为瓶颈时Task 很多时间在“等”但 UI 只告诉你“它花了很久”不告诉你“它在等什么”。第一步确认数据源真实读取效率不要信spark.read.parquet(hdfs://...)的日志里那句 “Reading from hdfs://...”要信hadoop fs -du -h /path/to/data的输出。我遇到过最坑的一次是某业务方提供的“已压缩 Parquet 表”hadoop fs -du显示 8TB但parquet-tools meta hdfs://.../part-00000.snappy.parquet | grep total uncompressed显示解压后达 42TB——Snappy 压缩比仅 5:1远低于 Parquet 默认的 10:1。这意味着每次读取CPU 要花额外时间解压而 Spark UI 的 CPU Utilization Metrics 默认不采集 JVM 外的 Native CPU 时间这部分开销完全隐身。解决方案很简单用parquet-tools批量扫描所有文件计算平均压缩比若低于 8:1则强制重写为 ZSTD 压缩spark.conf.set(spark.sql.parquet.compression.codec, zstd)实测在同等硬件下ZSTD 比 Snappy 解压速度快 2.3 倍压缩比高 1.8 倍。第二步抓取实时 I/O 性能指标在作业提交前先在 Driver 节点和至少两个 Executor 节点上运行# 采集 30 秒每秒 1 次重点关注 await等待 I/O 完成平均时间和 %util设备利用率 iostat -x 1 30 | awk $1 ~ /^[hs]d[0-9]$/ {print $1, $10, $14} iostat_driver.log # 同时采集网络看是否被占满 sar -n DEV 1 30 | grep eth0 | awk {print $2,$5,$6} netstat_driver.log 关键阈值await 20ms或%util 85%即告警rxkB/s txkB/s 90%网卡带宽即告警。曾有个作业在 K8s 环境下持续慢iostat显示磁盘正常但sar -n DEV发现eth0rxkB/s 长期 11.8GB/s万兆网卡上限 12.5GB/s原因是上游 Kafka Consumer 拉取数据后全部通过broadcast发送给所有 Executor而broadcast默认走 HTTP未启用spark.network.crypto.enabled导致加密开销巨大且无法 pipeline。解决方案是改用spark.sql.adaptive.enabledtruespark.sql.adaptive.coalescePartitions.enabledtrue让 AQE 自动合并小分区减少 broadcast 数据量网络压力直降 70%。第三步验证存储格式与读取方式匹配度Parquet 是列式存储优势在“只读所需列”。但如果代码里写df.select(*).filter(...)Spark 仍会读取所有列元数据再过滤。正确姿势是# 错误先选所有列再过滤I/O 浪费严重 df spark.read.parquet(hdfs://data).select(*).filter(user_id 1000) # 正确利用 Parquet 的谓词下推Predicate Pushdown df spark.read.parquet(hdfs://data).filter(user_id 1000) # Spark 自动只读 user_id 列 # 更进一步显式指定列 df spark.read.parquet(hdfs://data).select(user_id, event_time, action).filter(user_id 1000)实测某 500 列宽表仅读 3 列谓词下推I/O 量从 1.2TB 降至 87GB作业耗时从 41 分钟降至 9 分钟。实操心得I/O 诊断的黄金法则是“用存储系统原生命令验证而非 Spark 日志”。hadoop fs -du、parquet-tools、orc-tools这些工具输出的是字节级真相Spark UI 的 “Input Size” 字段经过多层抽象可能包含缓存命中、预读缓冲等干扰项。我习惯在每次新接入数据源时先跑一遍parquet-tools meta把num-rows、uncompressed-size、compression-codec记入 debug diary后续慢作业排查时直接比对省去 80% 的重复分析。3.2 Shuffle 层看穿“数据量不大”背后的重分布风暴Shuffle 是 Spark 最昂贵的操作因为它强制数据跨节点移动。但很多工程师被 UI 上的 “Shuffle Read/Write Bytes” 数值误导——看到只有几百 MB 就认为“shuffle 不重”却忽略了Shuffle 的真实成本 数据量 × 序列化开销 × 网络延迟 × 磁盘 Spill 频率。一个 200MB 的 shuffle若发生在 2000 个极小分区每个 100KB之间其序列化/反序列化次数是 2000×2000400 万次远高于 20 个大分区每个 10MB的 400 次。第一步定位 Shuffle 触发点与分区健康度打开 Spark UI 的 Stages Tab找到耗时最长的 Stage点开 Details看 “Shuffle Read/Write” 下的 “Number of Partitions”。健康阈值单个 Stage 的分区数应在spark.sql.shuffle.partitions的 0.5~2 倍之间。若远小于 0.5 倍如配置 200实际只有 12 个分区说明上游数据量极小或repartition被过度使用导致单个 Task 负载过重若远大于 2 倍如配置 200实际 1200 个说明存在大量mapPartitions或flatMap产生爆炸性数据或groupByKey未加reduce预聚合。更致命的是“分区倾斜”UI 的 “Task Duration” 柱状图若呈现“一个极高柱 其余极低柱”且高柱对应 Task 的 “Shuffle Read Size” 远超均值如均值 5MB最高 2.3GB这就是典型的数据倾斜。但注意不是所有倾斜都要解决。我处理过一个日志去重作业distinct()后发现 1 个 Task 读 1.8GB其余 199 个读 10MB。分析发现倾斜 Key 是user_idUNKNOWN占全量 37%。强行salting会增加 37% 的 shuffle 数据量得不偿失。最终方案是filter(user_id ! UNKNOWN).distinct().unionByName(spark.range(1).select(lit(UNKNOWN).alias(user_id)))用 2 行代码规避 90% 的 shuffle 开销。第二步量化 Shuffle 的序列化与 Spill 成本Spark UI 的 “Storage” Tab 里“Disk Spill” 字段常被忽略。若一个 Stage 的 “Disk Spill” 0说明 Executor 内存不足被迫把 shuffle 中间数据写入磁盘I/O 开销剧增。此时看 “Executor Metrics” 里的 “JVM Heap Memory Usage”若长期 85%则需调大spark.executor.memory或开启spark.memory.fraction动态调整。但更高效的做法是用spark.sql.adaptive.enabledtrue让 AQE 自动检测 Spill 并触发coalescePartitions。实测某 ETL 作业开启 AQE 后原本 1200 个分区被自动合并为 87 个Spill 从 14GB 降至 0shuffle time 从 18 分钟降至 2.3 分钟。第三步选择正确的 Shuffle Manager默认sortShuffle Manager 在大数据量时表现稳定但小数据量1GB下tungsten-sort的内存排序开销反而不如hash。可通过spark.conf.set(spark.shuffle.manager, hash)临时切换测试。但注意hash不支持aggregateByKey等需要 map-side combine 的操作。更推荐的方案是用spark.sql.adaptive.enabledtruespark.sql.adaptive.localShuffleReader.enabledtrueAQE 会自动将小 shufflespark.sql.adaptive.localShuffleReader.threshold默认 10MB转为本地读绕过网络实测提速 3~5 倍。注意Shuffle 优化的终极心法是“宁可少 shuffle不可乱 shuffle”。repartition(100)看似简单但若上游数据已按user_id排序用repartition(user_id)能让后续join直接走 SortMergeJoin避免二次 shuffle。我在某用户画像作业中将df.repartition(200).join(profile_df)改为df.repartition(user_id).join(profile_df.repartition(user_id))虽代码行数不变但 join 阶段 shuffle bytes 从 3.2TB 降至 0——因为 profile_df 已按user_id排序Spark 自动识别为 BroadcastHashJoin。3.3 Python 层揭开 PySpark UDF 的“性能黑洞”PySpark 的最大便利性也是其最大隐患。Java/Scala Spark 的计算在 JVM 内完成而 Python UDF 必须将 JVM 对象序列化为字节数组通过 socket 传给 Python 进程Python 计算后再序列化回 JVM。这个过程引入三重开销序列化/反序列化SerDe延迟、进程间通信IPC带宽瓶颈、Python GIL全局解释器锁串行化。第一步识别 UDF 类型与开销层级PySpark 提供三种 UDFudf()最慢。每个 Row 单独序列化IPC 调用频繁。100 万行数据需 100 万次 socket 通信。pandas_udf(type PandasUDFType.SCALAR)快 5~10 倍。以 Pandas Series 批量传输SerDe 次数降为分区数如 200 分区仅 200 次。pandas_udf(type PandasUDFType.GROUPED_AGG)最快。直接在 Pandas GroupBy 上操作避免 JVM-Python 来回搬运。实测对比对 1000 万行user_id做md5哈希udf()耗时 8.2 分钟pandas_udf(SCALAR)耗时 1.3 分钟pandas_udf(GROUPED_AGG)配合groupBy().agg()耗时 0.4 分钟。第二步监控 Python 进程真实资源占用Spark UI 的 “Executor Metrics” 只显示 JVM 内存不显示 Python 进程的 RSSResident Set Size内存。需在 Executor 节点上用ps aux --sort-%mem | head -20查看python进程内存。若 RSS 4GB假设 executor.memory8G说明 Python 进程吃掉了大量内存挤压 JVM 堆空间导致频繁 GC。此时必须降低pandas_udf的 batch sizespark.conf.set(spark.sql.execution.arrow.maxRecordsPerBatch, 10000)默认 10000可尝试 5000升级pyarrow至 11.0其to_pandas()方法在小 batch 下内存泄漏大幅减少关键用pandas_udf(returnTypeStringType())显式声明返回类型避免 PyArrow 自动推断的内存浪费。第三步用 Arrow 优化 SerDe 管道启用spark.sql.execution.arrow.pyspark.enabledtrue后PySpark 使用 Apache Arrow 的零拷贝内存格式在 JVM 和 Python 间传递数据SerDe 开销直降 60%。但需确保pyarrow版本 ≥ 7.0Arrow 7.0 引入RecordBatchReader支持流式读取pandas版本 ≥ 1.4修复了 Arrow 与 Pandas 交互的内存碎片问题关闭spark.sql.execution.arrow.pyspark.fallback.enabled默认 true强制失败而非降级到慢路径。曾有个地理围栏计算作业用udf()调用shapely库判断点是否在多边形内耗时 22 分钟。改用pandas_udf Arrow 后降至 3.8 分钟再将shapely替换为pygeosshapely的 C 后端最终降至 1.1 分钟——因为pygeos支持矢量化运算单次调用处理整个 Series彻底绕过 GIL。实操心得Python 层优化的口诀是“批处理、显声明、避 GIL”。永远优先pandas_udf禁用裸udf()所有pandas_udf必须带returnType注解计算密集型任务用numba.jit编译或pygeos/polars等原生库替代纯 Python 逻辑。我在 debug diary 里记了一条铁律“只要看到udf(开头的代码立刻标记为高危优先重构”。3.4 Memory 层读懂 JVM GC 日志里的“无声警报”Spark 的内存模型分为 Execution Memory用于 shuffle、join、aggregation和 Storage Memory用于 cache/broadcast两者共享spark.memory.fraction默认 0.6分配的堆空间。当 Execution Memory 不足时Spark 会驱逐 Storage Memory 中的缓存块反之亦然。但真正的杀手是Off-Heap 内存竞争PyArrow、Netty、甚至某些 JNI 库如snappy-java会直接申请堆外内存这部分不计入 JVM 堆却会挤占物理内存导致 OS OOM Killer 杀死 Executor 进程。第一步采集并解读 GC 日志启动 Spark 时添加 JVM 参数--conf spark.driver.extraJavaOptions-XX:PrintGCDetails -XX:PrintGCDateStamps -Xloggc:/tmp/driver_gc.log \ --conf spark.executor.extraJavaOptions-XX:PrintGCDetails -XX:PrintGCDateStamps -Xloggc:/tmp/executor_gc.log关键指标GC pause time 1s说明堆过大或 Survivor 区太小需调XX:SurvivorRatioFull GC frequency 1/min堆内存严重不足或存在内存泄漏如Broadcast未unpersistGC overhead 98%日志中GC time占总运行时间比例立即扩容spark.executor.memory。我处理过一个 Full GC 频繁的作业日志显示Full GC (Ergonomics)每 47 秒一次。检查发现spark.sql.adaptive.enabledtrue但spark.sql.adaptive.localShuffleReader.enabledfalseAQE 生成的大量小 shuffle 文件被缓存在 Off-Heap而spark.memory.offHeap.size未配置默认 0导致 Netty Buffer 无限制申请内存OS 内存耗尽后触发 JVM Full GC。解决方案spark.conf.set(spark.memory.offHeap.size, 4g)并设spark.memory.offHeap.enabledtrueGC 频率降至 0。第二步监控 Off-Heap 内存真实占用用jstat -gc pid查看OUOff-Heap Used字段。若OU持续增长且不回落说明存在 Off-Heap 泄漏。常见原因spark.sql.adaptive.enabledtrue时AQE 的AdaptiveSparkPlanExec会缓存大量QueryStageExec元数据在 Off-Heap使用spark.sql.adaptive.coalescePartitions.enabledtrue时CoalescedPartitionSpec对象堆积pandas_udf的 Arrow buffers 未及时释放。解决方案设置spark.sql.adaptive.enabledtrue时必须配spark.sql.adaptive.coalescePartitions.enabledtrue和spark.sql.adaptive.localShuffleReader.enabledtrue三者协同才能释放 Off-Heap对pandas_udf在函数末尾显式调用del df和gc.collect()虽不保证立即释放但能加速最狠一招spark.conf.set(spark.sql.adaptive.enabled, false)用确定性换取稳定性——在超大规模作业中AQE 的动态性有时不如静态规划可靠。第三步平衡 Execution 与 Storage Memoryspark.memory.fraction0.6是通用值但对 I/O 密集型作业如大量cache()应提高 Storage 比例spark.conf.set(spark.memory.storageFraction, 0.5)Storage 占 ExecutionStorage 总和的 50%对计算密集型如复杂 UDF应降低 Storage 比例spark.conf.set(spark.memory.storageFraction, 0.3)。实测某机器学习特征工程作业storageFraction从 0.5 降至 0.3 后Execution Memory 不足导致的 Spill 从 8GB 降至 0作业提速 22%。提示Memory 层诊断的底线思维是“宁可保守不可激进”。spark.executor.memory不要盲目堆高8G~16G 是安全区间spark.memory.fraction不要轻易调至 0.8 以上否则 Storage Memory 过大会挤压 Execution引发 shuffle Spillspark.memory.offHeap.size必须显式设置且建议为executor.memory的 25%~50%。我在 debug diary 里画了一张内存分配饼图每次调参都先更新它确保所有配置项逻辑自洽。4. 实操全流程从作业提交到瓶颈定位的 7 分钟标准动作这套流程是我团队每日 standup 后必做的“慢作业急救包”7 分钟内完成初步归因准确率超 85%。它不依赖任何外部工具只用 Spark 自带能力 Linux 基础命令。4.1 第 0-1 分钟作业提交与基础观测提交作业时务必加-vverbose参数获取完整日志并记录Application IDspark-submit --master yarn --deploy-mode cluster \ --conf spark.sql.adaptive.enabledtrue \ --conf spark.sql.adaptive.coalescePartitions.enabledtrue \ --name debug_slow_job \ your_script.py --input hdfs://data --output hdfs://out作业提交后立即打开 Spark UIURL 在日志末尾记下Application ID如application_1678901234567_89012Driver URL用于后续jstack当前活跃的 Stages 列表重点关注 “Duration” 最长的 Stage ID如Stage 5。注意不要等作业跑完在作业运行 30 秒后就开始观测。慢作业的瓶颈往往在前几个 Stage 就已暴露等它跑满 2 小时再分析黄花菜都凉了。4.2 第 1-3 分钟I/O 与 Shuffle 快速扫描I/O 快扫在 Spark UI 的 “Stages” Tab点开最慢 Stage看 “Input Size” 和 “Records Read”。若 “Input Size” 10GB 但 “Records Read” 100 万说明数据极度稀疏或格式异常如 JSON 嵌套过深立即怀疑 I/O 解析瓶颈同时在 Driver 节点运行iostat -x 1 5 | grep r/s\|w/s\|await若await 15ms标记为 I/O 瓶颈。Shuffle 快扫在同一 Stage 的 Details 页面看 “Shuffle Read/Write” 下的 “Number of Partitions”。若 50 或 500标记为 Shuffle 分区异常看 “Task Duration” 柱状图若存在单个 Task Duration 其他 Task 均值 5 倍且其 “Shuffle Read Size” 均值 10 倍标记为数据倾斜。4.3 第 3-5 分钟Python 与 Memory 初筛Python 快筛在 Spark UI 的 “Executors” Tab看各 Executor 的 “JVM Heap Memory Usage”。若某 Executor 的 Heap Usage 长期 40%但作业仍慢大概率是 Python 进程在吃 CPUJVM 闲着Python 忙着登录该 Executor 节点运行top -p $(pgrep -f python.*pyspark)若python进程 CPU% 90%且RES物理内存 3GB标记为 Python 层瓶颈。Memory 快筛在 Spark UI 的 “Storage” Tab看 “Disk Spill” 是否 0。若有立即检查 “Executor Metrics” 里的 “JVM Heap Memory Usage”若 85%标记为 Memory 不足在 Driver 节点用yarn application -status application_1678901234567_89012 | grep Final-State若为KILLED再查yarn logs -applicationId application_1678901234567_89012 | grep Container killed若含memory limit则是 Off-Heap 内存超限。4.4 第 5-7 分钟交叉验证与归因决策将前 4 分钟的标记汇总按优先级排序标记组合归因结论首选修复动作I/O await 15ms Input Size 10GBI/O 解析瓶颈检查parquet-tools meta换 ZSTD 压缩Shuffle 分区 50 Disk Spill 0Shuffle 分区过少repartition(200)或开启 AQEPython CPU% 90% Heap Usage 40%Python UDF 瓶颈将udf()改为pandas_udf()Disk Spill 0 Heap Usage 85%JVM 堆内存不足spark.executor.memory12gDisk Spill 0 Off-Heap KilledOff-Heap 内存超限spark.memory.offHeap.size4g决策原则只做一项改动重新提交作业再走一遍 7 分钟流程。切忌“一口气调 10 个参数”那不是调试是赌博。实操心得这 7 分钟流程的价值不在于它能解决所有问题而在于它把模糊的“慢”转化为清晰的“是什么慢”。我团队的新成员入职第一周任务就是用这个流程跑通 10 个历史慢作业每人提交一份 debug diary记录每次改动前后的耗时、关键指标变化、以及“我为什么这么改”。三个月后他们独立处理慢作业的首次归因准确率从 32% 提升到 89%。记住诊断不是目的可复现、可验证、可回滚的改动才是终点。5. 常见问题与排查技巧实录5.1 “UI 上 Stage Duration 很短但整个作业耗时很长” —— 你漏看了 DriverSpark UI 默认只展示 Executor 侧的 Stage Duration而 Driver 端的开销如collect()后的 Python 处理、toPandas()的数据拉取、show()的格式化完全不体现。曾有个作业 UI 显示总耗时 2 分钟但spark-submit命令实际耗时 22 分钟。jstackDriver 进程发现toPandas()在 Driver 端单线程解析 500 万行数据占了 20 分钟。解决方案避免toPandas()改用df.write.mode(overwrite).csv(hdfs://out)若必须collect()先limit(10000)再collect()用spark.conf.set(spark.sql.adaptive.enabled, false)关闭 AQE防止 Driver 端生成复杂 AdaptivePlan。5.2 “开启了 AQE但作业更慢了” —— AQE 不是银弹AQE 的coalescePartitions在数据量极不均衡时如 99% 分区 1MB1% 分区 100MB会错误地将小分区合并导致大分区 Task 负载暴增。解决方案设置spark.sql.adaptive.coalescePartitions.enabledtrue时必须配spark.sql.adaptive.coalescePartitions.minPartitionSize64mb默认 1mb避免过度合并对已知数据倾斜的join显式用broadcast或salting禁用 AQE 的自动优化df1.join(b