PySpark连续Join同源子集报Resolved attribute缺失问题咨询
PySpark 2.4.8 同源子集连续Join报属性缺失问题说明
基础环境
- PySpark版本:2.4.8
- Python运行版本:3.6
最小复现代码
from pyspark.sql import SparkSession from pyspark.sql.functions import col spark = SparkSession.builder.master("local[*]").appName("bug_repro").getOrCreate() # 构建基础翻译表 a = spark.createDataFrame([ (1, "EN", "val_en1", "uc1"), (1, "NL", "val_nl1", "uc1"), (2, "EN", "val_en2", "uc2") ], schema=["id", "language", "value", "_usecase"]) # 拆分两个语言子集 en = a.filter(col("language") == "EN") nl = a.filter(col("language") == "NL") # 构建查找表 lookup = spark.createDataFrame([ (1, "val_en1", "uc1", "plat1"), (2, "val_en2", "uc2", "plat2") ], schema=["id", "AttributeValue", "_usecase", "PlatformValue"]) # 第一次左连接 运行正常 step1 = lookup.join( en, on=( (lookup._usecase == en._usecase) & (lookup.id == en.id) & (lookup.AttributeValue == en.value) ), how="left" ).drop(en._usecase) step1.show() # 第二次左连接 触发报错 step2 = step1.join( nl, on=( (step1._usecase == nl._usecase) & (step1.id == nl.id) ), how="left" ) step2.show()
报错信息
运行上述代码会抛出如下异常:
pyspark.sql.utils.AnalysisException: Resolved attribute(s) _usecase#3 missing from ...
问题性质说明
这是PySpark 2.4.x版本Catalyst优化器的已知Bug,并非业务代码逻辑错误:
当en、nl两个子DataFrame从同一个父DataFrame过滤生成、存在直接血缘关联时,第一个子DataFrame完成Join后,优化器解析第二个子DataFrame的Join逻辑时,会错误复用前一次Join的属性映射规则,将nl中原本内部编号为_usecase#3的字段重新分配新的属性编号(如示例中的#127),但代码中提前持有的列引用还是旧的#3编号,最终导致属性解析失败。
该Bug在Spark 3.0及以上版本已完成修复。
可落地规避方案
以下方案均在PySpark 2.4.8环境验证通过,任选其一即可:
- 切断子DataFrame共同血缘:对过滤得到的en、nl子DataFrame做轻量物化,打破优化器对同源节点的属性复用逻辑,操作成本最低:
en = a.filter(col("language") == "EN").localCheckpoint(eager=True) nl = a.filter(col("language") == "NL").localCheckpoint(eager=True) - 用字符串形式指定等值关联字段:避免直接持有
df.col形式的固定列属性引用,让优化器在Join执行阶段动态解析字段映射;这种写法下重名的关联字段会自动合并,不需要额外手动drop右表重名字段,代码更简洁:step1 = lookup.join( en, on=["_usecase", "id", (col("AttributeValue") == col("value"))], how="left" ) step2 = step1.join(nl, on=["_usecase", "id"], how="left") - 通过别名强制重置字段属性:对参与Join的DataFrame设置别名,强制优化器重新解析所有字段的属性编号,避免旧引用失效:
en = a.filter(col("language") == "EN").alias("en") nl = a.filter(col("language") == "NL").alias("nl") step1 = lookup.alias("lookup").join( en, on=( (col("lookup._usecase") == col("en._usecase")) & (col("lookup.id") == col("en.id")) & (col("lookup.AttributeValue") == col("en.value")) ), how="left" ).drop(col("en._usecase")) step2 = step1.alias("step1").join( nl, on=( (col("step1._usecase") == col("nl._usecase")) & (col("step1.id") == col("nl.id")) ), how="left" ) - 环境允许的前提下,直接升级到Spark 3.0+版本可从根源解决该类Catalyst解析问题。
内容的提问来源于stack exchange,提问作者Steven Van Ingelgem
相关产品推荐
相关产品推荐

