Spark 3.2.2多次连接同DataFrame后无法删除列问题咨询
问题描述
我们有一段PySpark代码,将table_a表两次左连接到table_b表,每次连接后需要删除引入的key_hash列。该代码在Spark 3.0.1中运行正常,但升级到Spark 3.2.2后,第一次transform操作成功删除key_hash字段,第二次transform操作后key_hash字段仍保留在output_df中。改用.drop('key_hash')替代.drop(df_a.key_hash)则可正常删除列。
代码示例:
def tr_join_sac_user(self, df_a): def inner(df_b): return ( df_b.join(df_a, on=df_b["sac_key_hash"] == df_a["key_hash"], how="left") .drop(df_a.key_hash) .drop(df_b.sac_key_hash) ) return inner def tr_join_sec_user(self, df_a): def inner(df_b): return ( df_b.join(df_a, on=df_b["sec_key_hash"] == df_a["key_hash"], how="left") .drop(df_a.key_hash) .drop(df_b.sec_key_hash) ) return inner table_a_df = spark.read.format("delta").load("/path/to/table_a") table_b_df = spark.read.format("delta").load("/path/to/table_b") output_df = table_b_df.transform(tr_join_sac_user(table_a_df)) output_df = output_df.transform(tr_join_sec_user(table_a_df))
原因分析
这个差异源于Spark 3.2版本对drop(Column)方法的列匹配逻辑调整:
- 在Spark 3.0.1及更早版本中,
drop(Column)通过列名匹配删除列,只要列名一致就会被删除。 - Spark 3.2+版本中,
drop(Column)改为基于**列的唯一标识(如列的来源、生成上下文)**进行匹配,而非单纯列名。
在你的场景中:
- 第一次join后,使用
df_a.key_hash(原始table_a_df的列实例)删除列,此时join引入的key_hash列和原始列实例标识一致,所以删除成功。 - 第二次join时,
df_b是第一次transform后的DataFrame,再次join同一个table_a_df会生成新的key_hash列实例(尽管列名相同,但内部标识已不同)。此时用原始的df_a.key_hash列实例去匹配,无法识别新生成的列,因此删除失败。 - 而
.drop('key_hash')始终通过列名匹配,不受列实例标识变化的影响,所以两次都能正常删除。
解决方案
推荐以下两种兼容各版本的修改方式:
方式1:使用字符串列名删除(最稳妥)
直接将.drop(df_a.key_hash)替换为.drop("key_hash"),修改后的代码如下:
def tr_join_sac_user(self, df_a): def inner(df_b): return ( df_b.join(df_a, on=df_b["sac_key_hash"] == df_a["key_hash"], how="left") .drop("key_hash") .drop(df_b.sac_key_hash) ) return inner def tr_join_sec_user(self, df_a): def inner(df_b): return ( df_b.join(df_a, on=df_b["sec_key_hash"] == df_a["key_hash"], how="left") .drop("key_hash") .drop(df_b.sec_key_hash) ) return inner
方式2:使用当前DataFrame的列实例删除
如果需要保留Column对象的用法,可以在join后的DataFrame中直接引用列:
def tr_join_sec_user(self, df_a): def inner(df_b): joined_df = df_b.join(df_a, on=df_b["sec_key_hash"] == df_a["key_hash"], how="left") return joined_df.drop(joined_df["key_hash"]).drop(df_b.sec_key_hash) return inner
代码库影响评估
- 影响范围:所有使用
drop(Column)方式删除列、且存在多次复用同一源DataFrame进行join的场景,都可能出现类似的删除失败问题。 - 修复成本:需要遍历代码库,排查所有此类
drop(Column)用法,替换为字符串列名或当前DataFrame的列引用方式。 - 兼容性:替换后的代码在Spark 3.0.1及3.2.2版本中都能正常运行,无版本兼容性问题。
内容的提问来源于stack exchange,提问作者Riaz
相关产品推荐
相关产品推荐

