Azure Databricks执行pandas_udf时报Checkpoint block不存在错误如何排查
问题根因判断
该问题90%以上和内存压力强相关,核心逻辑如下:
你用到的applyInPandas + pandas UDF跑fbprophet的场景,单分组计算时会把对应分组的所有数据加载到executor内存中,而Spark默认的local checkpoint是将中间RDD块存储在产生该块的executor本地存储中。当分组最多、数据量最大的国家运行时,executor负载过高触发OOM,被Azure Databricks集群自动回收,对应存储的checkpoint块就会丢失,就会抛出你遇到的报错。之前数据量小时executor内存足以承载计算负载,不会被回收,所以运行正常,新增数据后内存阈值被突破就变成必现问题。
可行排查路径
1. 验证内存关联猜想
- 打开Azure Databricks集群的Spark UI,找到对应失败任务的Executor标签页,查看异常节点的
memoryUsed、diskSpilled指标,确认是否有executor出现内存使用率超过90%、大量磁盘溢写的记录 - 查看集群的日志,搜索
OutOfMemoryError、ExecutorLostFailure关键词,确认报错时间点是否有executor因为OOM被kill的记录
2. 短期快速修复方案
- 按照报错提示的建议,把默认的local checkpoint换成持久化到分布式存储的checkpoint:先在代码中配置checkpoint目录到ADLS等Azure分布式存储路径
spark.sparkContext.setCheckpointDir("你的分布式存储路径"),然后在applyInPandas之后调用result = result.checkpoint()再执行toPandas(),分布式checkpoint不会因为单executor挂掉丢失数据 - 临时调大executor的内存配置,比如把
spark.executor.memory调整到20G,spark.executor.memoryOverhead调整到4G,降低executor OOM概率 - 减少单executor同时处理的分组数,设置
spark.sql.shuffle.partitions为集群总核心数的2倍(你当前3个worker共12核,可设为24或32),避免单个executor同时加载太多大分组的数据导致内存溢出
3. 长期优化方案
- 对分组进行分片处理,不要让数据量最大的国家一次性跑所有产品分组,比如把该国家的key再拆成2-3批串行执行,降低单次计算的内存压力
- 优化fbprophet的预测逻辑,比如减少不必要的参数计算、截断过长的历史时序数据,降低单分组的内存占用
- 如果后续数据还会持续增长,可以考虑升级worker节点规格,或者增加worker节点数量
内容的提问来源于stack exchange,提问作者Ali Waheed
相关产品推荐
相关产品推荐

