PySpark合并DataFrame并新增来源标识列的问题求助
问题描述
现有两个PySpark DataFrame:
df1数据:
mail_id e_no aaa@ven.com 111 bbb@ven.com 222 ccc@ven.com 333
df2数据:
email emp_no aaa@ven.com 111 222 333
需求:通过left join保留df1的所有记录,用df2填充emp_no,同时新增src列标识emp_no的来源,期望结果:
email emp_no src aaa@ven.com 111 tbl1 bbb@ven.com 222 tbl2 ccc@ven.com 333 tbl2
原代码:
final_df = df1.join(df2, df1['e_no'] == df2['emp_no'], 'left') final_df = final_df.withColumn('src', F.when('mail_id'.isNull() | mail_id =='', 'tbl2') .when('mail_id' != '', 'tbl1') .otherwise(F.lit('')))
错误分析
- 列引用错误:
'mail_id'.isNull()是对字符串调用isNull(),不是对DataFrame列的操作,正确写法应为F.col('mail_id').isNull();同理mail_id ==''也未正确引用列,需改为F.col('mail_id') == ''。 - 逻辑判断错误:
mail_id是df1的字段,left join后df1的所有记录都会保留,因此mail_id永远不会为空或空字符串,导致第二个when条件始终成立,所有src都会被设为tbl1,完全不符合需求。 - 未明确emp_no填充规则:原代码未指定最终
emp_no的取值逻辑,无法实现“用df2填充emp_no”的要求。
修改后的代码
from pyspark.sql import functions as F # 执行left join,保留df1所有记录 final_df = df1.join(df2, df1['e_no'] == df2['emp_no'], 'left') # 生成最终email列:直接取df1的mail_id final_df = final_df.withColumn('email', F.col('mail_id')) # 生成最终emp_no列:优先取df2的emp_no,为空则用df1的e_no(此处两者值一致,也可直接取df1的e_no) final_df = final_df.withColumn('emp_no', F.coalesce(F.col('emp_no'), F.col('e_no'))) # 新增src列:当df2的email与df1的mail_id匹配(即df2.email非空)时,来源为tbl1,否则为tbl2 final_df = final_df.withColumn( 'src', F.when(F.col('df2.email').isNotNull() & (F.col('df2.email') == F.col('mail_id')), 'tbl1') .otherwise('tbl2') ) # 选择目标列并输出 final_df.select('email', 'emp_no', 'src').show()
代码说明
- 关联逻辑:通过
e_no和emp_no关联,确保df1的每条记录都能匹配到df2中对应的emp_no行。 - email列处理:直接复用df1的
mail_id作为最终的email列,匹配期望结果格式。 - emp_no列处理:用
coalesce函数优先取df2的emp_no,为空则 fallback 到df1的e_no,保证数据完整性。 - src列判断:通过检查df2的
email是否非空且与df1的mail_id一致,来区分emp_no的来源,完全符合需求。
内容的提问来源于stack exchange,提问作者user175025
相关产品推荐
相关产品推荐

