You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Spark执行dropDuplicates()遇ExecutorLostFailure失败,求配置优化

问题分析与配置调整方案

当前配置的核心问题

  1. Executor资源分配过度挤占节点资源:每个DataNode是8核/24GB,你给每个Executor分配7核/21GB,几乎把节点资源全占了。但节点本身需要运行HDFS DataNode进程、操作系统服务,剩余的1核/3GB资源不足以支撑,会导致节点负载陡增,出现心跳超时、连接重置,最终Executor被判定丢失。
  2. Shuffle配置未适配超大数据量:dropDuplicates()本质是全局Shuffle操作(需要按全量字段哈希后聚合去重),20亿条记录的Shuffle数据量极大,默认的200个Shuffle分区会导致每个分区数据量过大,引发内存溢出或任务超时。
  3. Driver资源未配置:未指定Driver的内存和核心数,处理20亿数据的元数据和任务调度时,Driver容易成为瓶颈。

具体调整方案

1. 调整Executor基础资源配置

每个DataNode预留至少2核/4GB给系统和HDFS进程,剩余资源拆分给多个Executor:

  • spark.executor.cores: 3(每个Executor用3核,每个DataNode可跑2个Executor,3个DataNode共6个)
  • spark.executor.memory: 8g(每个Executor分配8GB,预留8GB给节点)
  • spark.executor.instances: 6(3个DataNode×2个Executor/节点)

2. 优化Shuffle相关配置

针对超大数据量的去重Shuffle,调整以下参数:

  • spark.sql.shuffle.partitions: 4000(大幅增加Shuffle分区数,让每个分区数据量控制在合理范围,避免OOM)
  • spark.sql.adaptive.enabled: true(开启自适应执行,Spark会根据实际数据量动态调整分区数和任务并行度)
  • spark.shuffle.spill.compress: true(开启Shuffle溢写数据压缩,减少磁盘IO开销)
  • spark.network.timeout: 300s(延长网络超时时间,避免临时负载高导致心跳判定失败)
  • spark.executor.heartbeatInterval: 60s(调整Executor心跳间隔,降低误判概率)

3. 配置Driver资源

给Driver分配足够资源处理元数据:

  • spark.driver.memory: 8g
  • spark.driver.cores: 2

4. 数据读取与去重优化

  • 读取Parquet时开启向量化读取,提升读取效率:spark.sql.parquet.enableVectorizedReader=true
  • 如果不是全字段去重,在dropDuplicates()中指定具体字段(比如dropDuplicates(["id", "create_time"])),减少Shuffle时的Key大小,降低内存开销
  • 先过滤无效数据(比如空值、不符合业务规则的记录),减少后续处理的数据量

修改后的配置示例

from pyspark.conf import SparkConf
from pyspark.sql import SparkSession

configuration = {
    "spark.sql.debug.maxToStringFields": 100,
    # Executor配置
    "spark.executor.cores": 3,
    "spark.executor.memory": "8g",
    "spark.executor.instances": 6,
    # Driver配置
    "spark.driver.memory": "8g",
    "spark.driver.cores": 2,
    # Shuffle优化
    "spark.sql.shuffle.partitions": 4000,
    "spark.sql.adaptive.enabled": "true",
    "spark.shuffle.spill.compress": "true",
    "spark.network.timeout": "300s",
    "spark.executor.heartbeatInterval": "60s",
    # Parquet读取优化
    "spark.sql.parquet.enableVectorizedReader": "true"
}
config = SparkConf().setAll(configuration.items())
spark = (
    SparkSession.builder.master("spark://192.168.1.1:7077")
    .config(conf=config)
    .appName("my_app")
    .getOrCreate()
)

# 先过滤无效数据(示例,根据实际业务调整)
df = spark.read.parquet(f"{directory}/*.parquet").filter("id is not null")
# 指定去重字段(如果不是全字段)
df = df.dropDuplicates(["id", "create_time"])
df.count()

额外建议

  • 如果集群支持,可以开启外部Shuffle服务(spark.shuffle.service.enabled=true),避免Executor丢失时Shuffle数据丢失,减少任务重试开销
  • 监控集群节点的CPU、内存、磁盘IO情况,观察调整后的负载变化,再进一步微调配置

内容的提问来源于stack exchange,提问作者CopyOfA

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.15 11:50:22