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)或资源分配不匹配,具体包括:
- Upsert阶段Shuffle压力过载:Hudi Upsert需要基于RecordKey做全局查找,触发大量Shuffle操作。你设置的
hoodie.upsert.shuffle.parallelism=500远超Executor可承载的并发任务数,导致单个Executor内存耗尽。 - 分区数据倾斜:如果
rpt_partition_id存在数据倾斜(某分区数据量远大于其他分区),处理该分区的Executor会因负载过高被杀死。 - 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
相关产品推荐
相关产品推荐

