Spark执行dropDuplicates()遇ExecutorLostFailure失败,求配置优化
问题分析与配置调整方案
当前配置的核心问题
- Executor资源分配过度挤占节点资源:每个DataNode是8核/24GB,你给每个Executor分配7核/21GB,几乎把节点资源全占了。但节点本身需要运行HDFS DataNode进程、操作系统服务,剩余的1核/3GB资源不足以支撑,会导致节点负载陡增,出现心跳超时、连接重置,最终Executor被判定丢失。
- Shuffle配置未适配超大数据量:
dropDuplicates()本质是全局Shuffle操作(需要按全量字段哈希后聚合去重),20亿条记录的Shuffle数据量极大,默认的200个Shuffle分区会导致每个分区数据量过大,引发内存溢出或任务超时。 - 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: 8gspark.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
相关产品推荐
相关产品推荐

