PySpark下基于多优先级条件的DataFrame连接实现问询
PySpark 多优先级条件连接DataFrame实现方案
方法一:PySpark DataFrame API 实现
核心思路是逐级匹配、优先级降级,先处理最高优先级的匹配,再对未匹配的行依次使用更低优先级条件,每一步都处理去重逻辑,最后合并所有有效匹配结果。
示例代码
from pyspark.sql import SparkSession from pyspark.sql.functions import col, lit from pyspark.sql.window import Window spark = SparkSession.builder.appName("MultiPriorityJoin").getOrCreate() # 构造示例数据 gr_data = [("A", "X", "1"), ("A", "X", "2"), ("A", "Y", "1"), ("B", "X", "1"), ("C", "Z", "3")] GR_df = spark.createDataFrame(gr_data, ["columnx", "columny", "columnz"]) hk_data = [("A", "X", "1", "val1"), ("A", "X", "1", "val1"), ("A", "X", "2", "val2"), ("A", "Y", "2", "val3"), ("B", "X", "2", "val4"), ("B", "Y", "1", "val5"), ("A", "Z", "1", "val6")] HK_df = spark.createDataFrame(hk_data, ["columnh", "columni", "columj", "value"]) # 1. 最高优先级:三列全匹配,标记优先级1,去重 cond1_matches = GR_df.join( HK_df, (GR_df.columnx == HK_df.columnh) & (GR_df.columny == HK_df.columni) & (GR_df.columnz == HK_df.columj), "inner" ).select(GR_df["*"], HK_df["value"], lit(1).alias("priority")).dropDuplicates() # 筛选GR_df中未被条件1匹配的行 gr_remaining1 = GR_df.join( cond1_matches.select("columnx", "columny", "columnz"), on=["columnx", "columny", "columnz"], "left_anti" ) # 2. 优先级2:前两列匹配,标记优先级2,去重 cond2_matches = gr_remaining1.join( HK_df, (gr_remaining1.columnx == HK_df.columnh) & (gr_remaining1.columny == HK_df.columni), "inner" ).select(gr_remaining1["*"], HK_df["value"], lit(2).alias("priority")).dropDuplicates() # 筛选GR_df中未被条件1、2匹配的行 gr_remaining2 = gr_remaining1.join( cond2_matches.select("columnx", "columny", "columnz"), on=["columnx", "columny", "columnz"], "left_anti" ) # 3. 优先级3:第一+第三列匹配,标记优先级3,去重 cond3_matches = gr_remaining2.join( HK_df, (gr_remaining2.columnx == HK_df.columnh) & (gr_remaining2.columnz == HK_df.columj), "inner" ).select(gr_remaining2["*"], HK_df["value"], lit(3).alias("priority")).dropDuplicates() # 筛选GR_df中未被条件1、2、3匹配的行 gr_remaining3 = gr_remaining2.join( cond3_matches.select("columnx", "columny", "columnz"), on=["columnx", "columny", "columnz"], "left_anti" ) # 4. 最低优先级:仅第一列匹配,标记优先级4,去重 cond4_matches = gr_remaining3.join( HK_df, gr_remaining3.columnx == HK_df.columnh, "inner" ).select(gr_remaining3["*"], HK_df["value"], lit(4).alias("priority")).dropDuplicates() # 合并所有有效匹配结果 final_result = cond1_matches.unionAll(cond2_matches).unionAll(cond3_matches).unionAll(cond4_matches) final_result.show()
方法二:Spark SQL 实现
通过左连接获取所有可能的匹配,用CASE标记优先级,再借助窗口函数筛选每个GR_df行的最高优先级匹配,最后去重并过滤无匹配的行。
示例代码
# 注册临时视图 GR_df.createOrReplaceTempView("GR") HK_df.createOrReplaceTempView("HK") # 编写SQL查询 sql_query = """ WITH all_matches AS ( SELECT g.columnx, g.columny, g.columnz, h.value, CASE WHEN g.columnx = h.columnh AND g.columny = h.columni AND g.columnz = h.columj THEN 1 WHEN g.columnx = h.columnh AND g.columny = h.columni THEN 2 WHEN g.columnx = h.columnh AND g.columnz = h.columj THEN 3 WHEN g.columnx = h.columnh THEN 4 ELSE NULL END AS priority FROM GR g LEFT JOIN HK h ON g.columnx = h.columnh AND (g.columny = h.columni OR g.columnz = h.columj OR g.columny = h.columni AND g.columnz = h.columj) ), ranked_matches AS ( SELECT *, ROW_NUMBER() OVER ( PARTITION BY columnx, columny, columnz ORDER BY priority ASC ) AS rn FROM all_matches WHERE priority IS NOT NULL ) SELECT DISTINCT columnx, columny, columnz, value, priority FROM ranked_matches WHERE rn = 1 """ # 执行查询 final_result_sql = spark.sql(sql_query) final_result_sql.show()
关键说明
- 两种方法都严格遵循优先级规则:高优先级匹配优先保留,低优先级仅处理未匹配的行
- 去重逻辑通过
dropDuplicates()或DISTINCT实现,可根据业务需求调整去重的列(比如需保留特定字段时) - 窗口函数中
ORDER BY priority ASC确保优先级1的匹配被优先选中,若同优先级有多个匹配,可追加其他字段(如value)稳定排序
内容的提问来源于stack exchange,提问作者okm1234
相关产品推荐
相关产品推荐

