You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何从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

原代码问题分析

你提供的代码存在多处语法和逻辑错误:

  1. join方法参数顺序错误:PySpark中join的正确格式是df.join(other, on_condition, how),不能直接在join后罗列列选择;
  2. 列名引用错误:未给列名加引号(如serial_no、s_no需写成'serial_no'、's_no');
  3. 关联逻辑混乱:原条件无法覆盖两种补全场景;
  4. 列名不存在:people['serial_no.people']是无效列名,实际应为people['s_no'];
  5. 语法错误:括号不匹配、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()

代码说明

  1. 使用coalesce优先保留原DataFrame的非空值,仅当原字段为空时取关联后的补全值;
  2. 分步版本通过两次join分别处理两种补全场景,逻辑更易懂;
  3. 简化版本通过或条件关联两种匹配场景,用when精准对应补全规则,最后distinct避免重复数据;
  4. 用F.lit("NA")处理无匹配值的填充需求。

内容的提问来源于stack exchange,提问作者user175025

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.04 11:07:30