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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 11:52:23