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

如何在PySpark中对比DataFrame差异并获取差异行全量对应字段

问题根因

你之前拿全字段做subtract时,会比对所有选中列的值,只要任意一列(比如record_id、姓名)两边不一致,哪怕p_user_id同时存在于两个表,也会被判定为差异行,所以行数远超19条。


最优解决方案

核心逻辑

  1. 先单独提取仅存在于单表的p_user_id集合(这一步你已经跑通,得到19条结果)
  2. 用这个ID集合和对应原表做内连接,就能精准拉取这些ID对应的全量字段,不会出现多余行

修正后的可运行代码

已同步修正你原代码里的变量名/视图名笔误:

# ---------- 1. 读取并处理MySQL侧数据 ----------
my_rac = spark.read.parquet("/Users/mysql.parquet")
my_rac.createOrReplaceTempView('my_rac')
# 去重后生成基础表
d_rac = spark.sql('select distinct * from my_rac')
d_rac.createOrReplaceTempView('d_rac')
# 读取需要的字段,统一把p_user_id转成string避免类型匹配问题
rac_p_user_df = spark.sql('''
select 
    cast(p_user_id as string) as p_user_id,
    record_id,
    contact_last_name,
    contact_first_name
from d_rac
''')

# ---------- 2. 读取并处理MSSQL侧CSV数据 ----------
sql_rac = spark.read.csv("/Users/mzn293/Downloads/kavi-20211116.csv")
sql_rac.createOrReplaceTempView('sql_rac')
# 读取对应字段,CSV默认列名规则为_c0、_c1以此类推
d_sql_rac = spark.sql('''
select distinct 
    _c0 as p_user_id,
    _c1 as record_id,
    _c4 as contact_last_name,
    _c5 as contact_first_name
from sql_rac
''')

# ---------- 3. 提取仅在MSSQL侧存在的差异p_user_id ----------
# 仅比对p_user_id字段,得到19个差异ID
diff_p_user_ids = d_sql_rac.select("p_user_id").subtract(rac_p_user_df.select("p_user_id"))

# ---------- 4. 关联获取差异ID对应的全量字段 ----------
# 内连接过滤,只保留19个差异ID对应的全量数据
diff_full_data = d_sql_rac.join(diff_p_user_ids, on="p_user_id", how="inner")
# 输出结果
diff_full_data.show(100, truncate=False)

# 若需要查询仅在MySQL侧存在的ID对应数据,反过来操作即可:
# diff_p_user_ids_mysql = rac_p_user_df.select("p_user_id").subtract(d_sql_rac.select("p_user_id"))
# diff_full_data_mysql = rac_p_user_df.join(diff_p_user_ids_mysql, on="p_user_id", how="inner")

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.25 13:45:03