PySpark中基于匹配字段数量更新flag的实现求助
PySpark 解决方案
方法一:DataFrame API 实现
代码实现
from pyspark.sql import functions as F # 替换为你实际需要对比的非键字段列表 non_key_cols = ["col1", "col2", "col3", "col4", "col5"] # 1. 过滤table1中flg为'new'的记录,与table2按accountid内连接(仅保留匹配成功的记录) t1_new = table1.filter(F.col("flg") == "new") joined_df = t1_new.join(table2, on="accountid", how="inner") # 2. 计算非键字段的匹配总数 match_count_expr = sum( F.when(F.col(f"table1.{col}") == F.col(f"table2.{col}"), 1).otherwise(0) for col in non_key_cols ).alias("match_count") # 3. 根据匹配数更新flg字段 processed_new_df = joined_df.select( "accountid", F.when(F.col("match_count") >= 3, "ready") # 不匹配数>2 → 匹配数 ≤ 总非键数-3 .when(F.col("match_count") <= len(non_key_cols) - 3, "wip") .alias("flg"), *[F.col(f"table1.{col}") for col in non_key_cols] ) # 4. 合并table1中flg不为'new'的原记录 final_df = processed_new_df.unionByName(table1.filter(F.col("flg") != "new")) final_df.show()
补充说明
如果需要保留flg='new'但未匹配到table2的记录,只需将连接方式改为left,并在when逻辑最后添加.otherwise("new")即可。
方法二:Spark SQL 实现
代码实现
# 将DataFrame注册为临时视图 table1.createOrReplaceTempView("table1") table2.createOrReplaceTempView("table2") # 替换为你实际需要对比的非键字段列表 non_key_cols = ["col1", "col2", "col3", "col4", "col5"] # 生成匹配计数的SQL表达式 match_count_sql = " + ".join( f"CASE WHEN t1.{col} = t2.{col} THEN 1 ELSE 0 END" for col in non_key_cols ) # 执行SQL逻辑 sql_query = f""" WITH processed_new AS ( SELECT t1.accountid, CASE WHEN ({match_count_sql}) >= 3 THEN 'ready' WHEN ({match_count_sql}) <= {len(non_key_cols)-3} THEN 'wip' END AS flg, {', '.join([f't1.{col}' for col in non_key_cols])} FROM table1 t1 JOIN table2 t2 ON t1.accountid = t2.accountid WHERE t1.flg = 'new' ) SELECT * FROM processed_new UNION ALL SELECT * FROM table1 WHERE flg != 'new' """ final_df = spark.sql(sql_query) final_df.show()
内容的提问来源于stack exchange,提问作者Matthew
相关产品推荐
相关产品推荐

