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

Hudi Upsert任务因ExecutorDeadException失败问题排查求助

问题描述

本人刚接触Hudi,对其内部机制了解有限。目前尝试在Spark上以Upsert模式运行一个处理约50GB基准数据的Hudi任务,代码如下:

from pyspark.sql import SparkSession 
spark = SparkSession \
    .builder \
    .appName("hudi_test") \
    .enableHiveSupport().getOrCreate()  

tableName = "hudi_test12"
basePath = "/tmp/rahul/hudi_test12" 

df = spark.read.parquet("/user/data/input")
df = df.repartition(500)
#df.show() 
hudi_options = {
    'hoodie.table.name': tableName,
    'hoodie.datasource.write.recordkey.field': 'row_id',
    'hoodie.datasource.write.partitionpath.field': 'rpt_partition_id',
    'hoodie.datasource.write.table.name': tableName,
    'hoodie.datasource.write.operation': 'upsert',
    'hoodie.datasource.write.precombine.field': 'created',
    'hoodie.upsert.shuffle.parallelism': 500,
    'hoodie.insert.shuffle.parallelism': 500
} 
df.write.format("hudi"). \
    options(**hudi_options). \
    mode("overwrite"). \
    save(basePath)

任务在SparkUpsertCommitActionExecutor阶段生成约3000个任务,完成1000个后其余任务开始失败,最终任务报错。Executor端日志多次打印如下错误后被终止:

23/01/24 12:37:55 ERROR client.TransportResponseHandler: Still have 1 requests outstanding when connection from mnplld-shddn02.india.airtel.itm/10.240.8.108:42916 is closed
23/01/24 12:37:55 INFO shuffle.RetryingBlockTransferor: Retrying fetch (1/3) for 1 outstanding blocks after 5000 ms
23/01/24 12:38:01 INFO client.TransportClientFactory: Found inactive connection to mnplld-shddn02.india.airtel.itm/10.240.8.108:42916, creating a new one.
23/01/24 12:38:01 ERROR shuffle.RetryingBlockTransferor: Exception while beginning fetch of 1 outstanding blocks (after 1 retries)
org.apache.spark.ExecutorDeadException: The relative remote executor(Id: 19), which maintains the block data to fetch is dead.
    at org.apache.spark.network.netty.NettyBlockTransferService$$anon$2.createAndStart(NettyBlockTransferService.scala:136)
    at org.apache.spark.network.shuffle.RetryingBlockTransferor.transferAllOutstanding(RetryingBlockTransferor.java:154)
    at org.apache.spark.network.shuffle.RetryingBlockTransferor.lambda$initiateRetry$0(RetryingBlockTransferor.java:184)
    at java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:515)
    at java.base/java.util.concurrent.FutureTask.run(FutureTask.java:264)
    at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1128)
    at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:628)
    at io.netty.util.concurrent.FastThreadLocalRunnable.run(FastThreadLocalRunnable.java:30)
    at java.base/java.lang.Thread.run(Thread.java:834)
23/01/24 12:38:01 ERROR storage.ShuffleBlockFetcherIterator: Failed to get block(s) from mnplld-shddn02.india.airtel.itm:42916
org.apache.spark.ExecutorDeadException: The relative remote executor(Id: 19), which maintains the block data to fetch is dead.
    at org.apache.spark.network.netty.NettyBlockTransferService$$anon$2.createAndStart(NettyBlockTransferService.scala:136)
    at org.apache.spark.network.shuffle.RetryingBlockTransferor.transferAllOutstanding(RetryingBlockTransferor.java:154)
    at org.apache.spark.network.shuffle.RetryingBlockTransferor.lambda$initiateRetry$0(RetryingBlockTransferor.java:184)
    at java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:515)
    at java.base/java.util.concurrent.FutureTask.run(FutureTask.java:264)
    at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1128)
    at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:628)
    at io.netty.util.concurrent.FastThreadLocalRunnable.run(FastThreadLocalRunnable.java:30)
    at java.base/java.lang.Thread.run(Thread.java:834)

Spark提交配置如下:

spark3-submit --master yarn --jars avro-1.10.0.jar,hudi-spark3.2-bundle_2.12-0.13.0-SNAPSHOT.jar --deploy-mode cluster --conf 'spark.serializer=org.apache.spark.serializer.KryoSerializer'   --conf 'spark.sql.catalog.spark_catalog=org.apache.spark.sql.hudi.catalog.HoodieCatalog'   --conf 'spark.sql.extensions=org.apache.spark.sql.hudi.HoodieSparkSessionExtension' --driver-memory 5g  --executor-memory 22g --num-executors 20 --executor-cores 12 --queue ocp test.py

我认为Executor和核心数量足以处理该数据,请问失败原因是什么?该如何解决?


问题分析与解决

核心原因

日志中的ExecutorDeadException明确说明Executor进程被YARN杀死,核心诱因是内存溢出(OOM)或资源分配不匹配,具体包括:

  1. Upsert阶段Shuffle压力过载:Hudi Upsert需要基于RecordKey做全局查找,触发大量Shuffle操作。你设置的hoodie.upsert.shuffle.parallelism=500远超Executor可承载的并发任务数,导致单个Executor内存耗尽。
  2. 分区数据倾斜:如果rpt_partition_id存在数据倾斜(某分区数据量远大于其他分区),处理该分区的Executor会因负载过高被杀死。
  3. Hudi索引内存开销过大:默认的Bloom索引在RecordKey基数极高时,会占用大量Executor堆内存,引发OOM。

解决方法

1. 调整Executor资源配置

  • 优化CPU与内存比例:当前executor-cores=12搭配executor-memory=22g,单核心内存仅约1.8g,建议将executor-cores降至8,保持executor-memory=22g,减少单个Executor的并发任务数,降低内存压力。
  • 提升Shuffle内存占比:通过spark.shuffle.memoryFraction=0.3或spark.memory.fraction=0.8增大Shuffle可用内存,避免Shuffle过程中OOM。

2. 优化Hudi Upsert参数

  • 合理设置并行度:根据总核心数(如20*8=160)将hoodie.upsert.shuffle.parallelism调整为200-300,避免过多并行度导致Shuffle开销激增。
  • 切换高效索引:改用GLOBAL_BLOOM索引并开启分区过滤hoodie.bloom.index.filter.by.partition=true,缩小索引查找范围;若数据量极大,可考虑HBase索引降低内存占用。
  • 移除不必要的Hive支持:如果无需与Hive交互,删除.enableHiveSupport(),减少额外资源消耗。

3. 处理数据倾斜

  • 打散倾斜分区:对rpt_partition_id加盐打散数据,比如:
    from pyspark.sql.functions import rand
    df = df.repartition(500, "rpt_partition_id", rand())
    
  • 开启动态分区:设置hoodie.upsert.dynamic.partitioning=true,让Hudi自动均衡分区数据分布。

4. 其他优化

  • 替换写入模式:Upsert模式下用mode("append")替代mode("overwrite"),避免删除原有表的额外开销,Hudi会自动处理更新与插入。
  • 查看YARN容器日志:通过yarn logs -applicationId <你的任务ID>查看被杀死的Executor日志,确认是否存在OutOfMemoryError,这是定位内存问题最直接的依据。

内容的提问来源于stack exchange,提问作者rk.the1

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 06:20:53