如何从people DataFrame补全persons中缺失的serial_no与mail字段?
DataFrame缺失字段补全的代码修正
现有DataFrames
persons DataFrame
name age serial_no mail John 25 100483 john@abc.com Sam 49 448900 sam@abc.com Will 63 will@abc.com Robert 20 299011 Hill 78 hill@abc.com
people DataFrame
name s_no e_mail John 100483 john@abc.com Sam 448900 sam@abc.com Will 229809 will@abc.com Robert 299011 Hill 567233 hill@abc.com
补全需求
- 若persons的
serial_no缺失,通过mail字段关联people获取s_no作为补全值; - 若persons的
mail缺失,通过serial_no字段关联people获取e_mail作为补全值; - 若无匹配值,填充为"NA"。
期望结果
name age serial_no mail John 25 100483 john@abc.com Sam 49 448900 sam@abc.com Will 63 229809 will@abc.com Robert 20 299011 NA Hill 78 567233 hill@abc.com
原代码问题分析
你提供的代码存在多处语法和逻辑错误:
join方法参数顺序错误:PySpark中join的正确格式是df.join(other, on_condition, how),不能直接在join后罗列列选择;- 列名引用错误:未给列名加引号(如
serial_no、s_no需写成'serial_no'、's_no'); - 关联逻辑混乱:原条件无法覆盖两种补全场景;
- 列名不存在:
people['serial_no.people']是无效列名,实际应为people['s_no']; - 语法错误:括号不匹配、
alias方法的括号使用错误。
修正后的代码(PySpark)
分步处理版本(逻辑清晰)
from pyspark.sql import functions as F from pyspark.sql.functions import coalesce # 第一步:通过mail关联,补全persons中缺失的serial_no df_mail_join = persons.join( people, persons['mail'] == people['e_mail'], how='left' ).withColumn( 'serial_no_filled', coalesce(persons['serial_no'], people['s_no']) ).select( persons['name'], persons['age'], 'serial_no_filled', persons['mail'] ).withColumnRenamed('serial_no_filled', 'serial_no') # 第二步:通过serial_no关联,补全mail字段 final_df = df_mail_join.join( people, df_mail_join['serial_no'] == people['s_no'], how='left' ).withColumn( 'mail_filled', coalesce(df_mail_join['mail'], people['e_mail'], F.lit("NA")) ).select( df_mail_join['name'], df_mail_join['age'], df_mail_join['serial_no'], 'mail_filled' ).withColumnRenamed('mail_filled', 'mail') # 查看结果 final_df.show()
简化版本(一次关联处理)
from pyspark.sql import functions as F final_df = persons.join( people, (persons['mail'] == people['e_mail']) | (persons['serial_no'] == people['s_no']), how='left' ).withColumn( 'serial_no', coalesce(persons['serial_no'], F.when(persons['mail'] == people['e_mail'], people['s_no'])) ).withColumn( 'mail', coalesce(persons['mail'], F.when(persons['serial_no'] == people['s_no'], people['e_mail']), F.lit("NA")) ).select( persons['name'], persons['age'], 'serial_no', 'mail' ).distinct() # 去重避免多匹配产生重复行 final_df.show()
代码说明
- 使用
coalesce优先保留原DataFrame的非空值,仅当原字段为空时取关联后的补全值; - 分步版本通过两次join分别处理两种补全场景,逻辑更易懂;
- 简化版本通过或条件关联两种匹配场景,用
when精准对应补全规则,最后distinct避免重复数据; - 用
F.lit("NA")处理无匹配值的填充需求。
内容的提问来源于stack exchange,提问作者user175025
相关产品推荐
相关产品推荐

