如何在PySpark中对比DataFrame差异并获取差异行全量对应字段
问题根因
你之前拿全字段做subtract时,会比对所有选中列的值,只要任意一列(比如record_id、姓名)两边不一致,哪怕p_user_id同时存在于两个表,也会被判定为差异行,所以行数远超19条。
最优解决方案
核心逻辑
- 先单独提取仅存在于单表的
p_user_id集合(这一步你已经跑通,得到19条结果) - 用这个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
相关产品推荐
相关产品推荐

