AWS Spark Glue作业修改DataFrame多列值超时且不分配给executor问题
问题原因分析
- 频繁触发Action导致重复计算:嵌套循环中每次值替换后调用的
show()是触发Spark作业执行的Action操作,配套的distinct()是需要Shuffle的宽依赖操作。每一轮循环都会提交一次完整的Spark作业,且随着循环推进,DataFrame的依赖链(Lineage)越来越长,每一次Action的计算开销都会线性上升,挤占大量集群资源。 - 执行计划过度膨胀:对同一列的多个值替换采用循环调用
withColumn+when的实现方式,每一次withColumn都会给DataFrame的执行计划追加一层节点。当映射规则数量较多时,执行计划会膨胀到Spark Driver端无法正常解析、生成任务的程度,无法将任务下发到Executor执行,出现无日志的永久卡住现象。 - 冗余逻辑放大开销:第一层循环中已经完成了列的字符串类型转换,后续嵌套循环的重复处理逻辑进一步放大了执行计划的复杂度。
优化方案
优化核心是减少中间Action触发、合并转换逻辑压缩执行计划体积,优化后代码如下:
def map_values_in_columns(self, df): for k, v in self.value_mapping.items(): column_name = v['column_name'] values = v['values'] # 一次性构建字符串类型映射规则,单次调用replace完成所有值替换 str_mapping = {str(old_val): str(new_val) for old_val, new_val in values.items()} df = df.withColumn( column_name, F.col(column_name).cast("string").replace(str_mapping) ) # 测试阶段可临时打开验证逻辑,生产环境务必移除 # df.select(column_name).distinct().show(truncate = False) # 所有转换完成后仅触发一次Action df.select(column_name).distinct().show(truncate = False) logger.info("Succesfully mapped values in columns") return df
如果映射规则量级超过1000条,建议将映射规则转为广播变量,通过Join方式完成替换,进一步降低执行计划复杂度。
内容的提问来源于stack exchange,提问作者E. Faslo
相关产品推荐
相关产品推荐

