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

Spark Scala单元测试挂起问题排查求助

Spark单元测试挂起问题排查与疑问

问题背景

生产环境运行正常的Spark作业,在单元测试时出现挂起。测试采用Mockito模拟主类的方式执行。

测试代码示例

mainInput = spy(
  new Input(
    // some initialization
  ))
val spiedApp = spy(MainApp) // MainApp是主类
doReturn(args).when(spiedApp).getCliArgs // getCliArgs是类中的方法,其他方法省略
....
spiedApp.main(mainInput)

线程Dump信息

挂起后通过线程dump得到关键锁信息:

- locked <0x00000005c128d510> (a java.lang.ref.ReferenceQueue$Lock)
- locked <0x00000005c084d350> (a java.lang.ref.ReferenceQueue$Lock)
- locked <0x00000005c084dd38> (a io.netty.channel.nio.SelectedSelectionKeySet)
- locked <0x00000005c084dd50> (a java.util.Collections$UnmodifiableSet)
.....
- locked <0x00000005c1e47998> (a org.apache.spark.sql.execution.datasources.v2.BatchScanExec)
- locked <0x00000005c19e97d8> (a org.apache.spark.sql.execution.QueryExecution)

挂起位置与日志

挂起发生在以下代码行:

spark.sql(s"select xx from xxx").dropDuplicates.collect

挂起时控制台日志

[scalatest] 23/12/30 09:27:53 INFO V2ScanRelationPushDown: 
[scalatest] Output: _corrupt_record#30

有无cache()的日志对比

添加cache()时的日志:

[scalatest] 23/12/30 11:02:25 INFO V2ScanRelationPushDown: 
[scalatest] Output: _corrupt_record#30, type#31, metadata#32, timestamp#33, sequence#34L, data#35
[scalatest]          
[scalatest] 23/12/30 11:02:25 INFO V2ScanRelationPushDown: 
[scalatest] Output: _corrupt_record#30, type#31, metadata#32, timestamp#33, sequence#34L, data#35

移除cache()后(挂起时)的日志:

[scalatest] 23/12/30 11:02:25 INFO V2ScanRelationPushDown: 
[scalatest] Output: _corrupt_record#30, type#31, metadata#32, timestamp#33, sequence#34L, data#35
[scalatest]          
[scalatest] 23/12/30 11:02:25 INFO V2ScanRelationPushDown: 
[scalatest] Output: 

已排查结论

  • 排除已知UDF导致Spark挂起的问题
  • 多次dump进程输出完全一致,无重复计算迹象
  • 当前测试仅启动单个Java进程,怀疑是Spark内部问题,询问是否可在测试中分离Driver和Executor进程

代码变更影响

本次问题由移除cache()调用引发:

  • 理论上移除cache后,isEmpty和collect操作会从头触发执行,对应两次V2ScanRelationPushDown日志
  • 因测试数据量过大必须移除cache,且计算成本较低,本应仅重复读取测试文件即可,但实际出现挂起
  • 无cache时第二次V2ScanRelationPushDown日志为空,未触发自定义的FilePartitionReaderFactory实现

自定义数据读取逻辑

作业使用自定义下推过滤,修改了以下类的逻辑,在文件载入Spark内存前先执行grep处理:

  • FilePartitionReaderFactory
  • TextBasedFileScan
  • FileScanBuilder

数据读取的核心代码:

spark.read.format(format).options(options.get).schema(schema.get).load(prefixesToRead: _*)

核心疑问

  1. spark.read与V2ScanRelationPushDown的关系是什么?
  2. 为什么移除cache()后,第二次执行未触发自定义的FilePartitionReaderFactory实现?
  3. 无cache时第二次V2ScanRelationPushDown的空输出是否是导致挂起的直接原因?

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.03 04:07:09