PySpark无关联条件合并两个DataFrame列及merge报错解决
报错原因
AttributeError: 'DataFrame' object has no attribute 'merge' 报错是因为merge()是Pandas DataFrame的API,PySpark的DataFrame对象没有提供该方法,直接调用必然抛出属性错误。
无关联条件的横向列拼接不能直接套用Pandas的写法,PySpark中DataFrame是分布式数据集,没有内置的按位置直接拼列的逻辑,需要通过临时行号锚定行位置后再关联,且必须保证两个DataFrame行数完全一致、行顺序符合预期,否则会出现数据错位。
实现代码
推荐方案:行号关联(无笛卡尔积风险)
引入行号作为临时关联键,内连接后删除临时字段即可得到包含x1、x2、x3、x4的目标结果:
from pyspark.sql.functions import monotonically_increasing_id, row_number from pyspark.sql.window import Window # 为df1生成连续行序号 df1_rn = df1.withColumn("tmp_id", monotonically_increasing_id()) df1_rn = df1_rn.withColumn( "rn", row_number().over(Window.orderBy("tmp_id")) ).drop("tmp_id") # 为df2生成连续行序号 df2_rn = df2.withColumn("tmp_id", monotonically_increasing_id()) df2_rn = df2_rn.withColumn( "rn", row_number().over(Window.orderBy("tmp_id")) ).drop("tmp_id") # 按行号关联后删除临时列 new_df = df1_rn.join(df2_rn, on="rn", how="inner").drop("rn")
简化方案:crossJoin过滤(仅适合行数完全确定的场景)
如果能100%确认两个DataFrame行数完全一致,可以用交叉连接后过滤行号相等的记录,写法更简洁,但如果行数不匹配会直接生成笛卡尔积脏数据,不建议在生产环境使用:
new_df = df1_rn.crossJoin(df2_rn).filter("rn = rn").drop("rn")
注意:禁止不加行号直接调用
crossJoin做横向拼接,否则会生成两表行数乘积的超大数据集,直接引发任务异常。
内容的提问来源于stack exchange,提问作者Kiran_t
相关产品推荐
相关产品推荐

