Spark DataFrame转RDD耗时极长:是懒执行机制还是真实故障?
核心结论
你猜测的懒执行特性导致耗时统计错位是完全正确的:.rdd本身是转换操作,不会触发执行,你在UI上看到的归到这一行的耗时,本质是这一行之前所有的懒加载操作(Solr数据读取、SQL查询、列新增)加上后续map、MongoSpark.save动作的总耗时,UI的统计只是把整个作业的耗时挂载到了第一个触发shuffle/action关联的转换节点上,和DataFrame转RDD本身无关。
常见诱因及排查方案
spark-solr分区配置不合理导致调度开销过大
spark-solr默认会根据Solr collection的分片数、splits_per_shard参数生成大量分区,哪怕最终计算结果只有90行,成百上千个空/小分区的调度、资源申请、任务分发都会在集群环境下产生极高的 overhead。本地standalone模式因为调度链路短,所以开销可以忽略。
排查方法:在resultDataFrame上执行println(result.rdd.getNumPartitions),如果分区数超过10基本就属于配置不合理。
解决方案:读Solr时调低splits_per_shard参数,或者在调用.rdd之前执行coalesce(1)合并分区(90行数据单分区完全足够)。Solr读取逻辑未下压过滤条件,拉取数据量远超预期
你单独操作Solr快不代表Spark侧的查询逻辑是最优的:如果你的SQL过滤、group by逻辑没有被spark-solr识别并下压到Solr服务端执行,就会出现Spark拉取Solr全表数据,再在内存中做计算的情况,哪怕最终结果只有90行,拉取全量数据的耗时也会极高,且全部算到.rdd关联的作业耗时里。
排查方法:在resultDataFrame上单独执行result.count(),看这个操作的耗时,如果单独count就需要数分钟,即可确认是前置读取和计算逻辑的问题。
解决方案:检查spark-solr的谓词下压配置是否开启,将SQL中的过滤条件尽量放到Solr的读取参数中,避免Spark侧拉取无用数据。算子内的序列化初始化和依赖冲突问题
你将org.json4s.DefaultFormats的初始化放在了map算子内部,等于每处理一行数据/每个分区都会初始化一次序列化规则,同时如果你的作业依赖包和集群自带的json4s版本存在冲突,会在每个任务执行时产生大量类加载开销,集群环境下这个问题会被放大。
解决方案:将implicit val formats的定义提到map算子外部,最好放到类的静态成员位置,避免重复初始化;打包时将json4s相关依赖做shade/重定位,和集群自带的版本隔离。Mongo写入配置不合理
虽然你测试Mongo本身响应快,但如果mongo-spark-connector的配置不合理,比如默认开启了高并发写入、批量提交阈值设置过大、重试次数过多,也会导致小数据量写入耗时极高。
排查方法:把写入逻辑替换成map后的RDD执行count()动作,看耗时是否明显下降,如果是则说明问题出在写入环节。
解决方案:调整mongo写入配置,调低spark.mongodb.output.batchSize、spark.mongodb.output.maxBatchSize参数,小数据量下可以直接设置为100。
内容的提问来源于stack exchange,提问作者uylmz

