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: _*)
核心疑问
spark.read与V2ScanRelationPushDown的关系是什么?- 为什么移除
cache()后,第二次执行未触发自定义的FilePartitionReaderFactory实现? - 无cache时第二次V2ScanRelationPushDown的空输出是否是导致挂起的直接原因?
内容的提问来源于stack exchange,提问作者user2988877
相关产品推荐
相关产品推荐

