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

如何基于另一DataFrame的两列标记PySpark DataFrame子集?

解决方案

要实现基于表B的规则标记表A的子集,核心是判断同一code分组下是否同时存在表B中的cod_1和cod_2,再填充对应rule值。以下是满足性能要求的PySpark和SQL实现:

PySpark 实现

from pyspark.sql import functions as F

# 1. 预计算每个code对应的所有servico集合,用于快速匹配判断
code_servicos = table_a.groupBy("code").agg(F.collect_set("servico").alias("servico_set"))

# 2. 将集合关联回原表A,再与表B做条件关联
table_a_with_set = table_a.join(code_servicos, on="code", how="inner")
joined_df = table_a_with_set.join(
    # 若表B数据量小,可添加F.broadcast()广播表B,减少shuffle
    F.broadcast(table_b),
    (table_a_with_set["servico"] == table_b["cod_2"]) & F.array_contains(table_a_with_set["servico_set"], table_b["cod_1"]),
    how="left"
)

# 3. 生成结果,优先取表B的rule值,无匹配则保留原表A的rule
result = joined_df.select(
    "code",
    "servico",
    F.coalesce(table_b["rule"], table_a_with_set["rule"]).alias("rule")
)

# 查看结果
result.show()

SQL 实现

先创建临时视图(适配Databricks环境):

CREATE OR REPLACE TEMP VIEW table_a AS
SELECT * FROM VALUES
    ('123', 'A1',''), 
    ('123', 'E1',''),
    ('123', 'A3',''),
    ('123', 'B1',''),
    ('123', 'B2',''),
    ('123', 'B3',''),
    ('321', 'C1',''),
    ('321', 'C2',''),
    ('321', 'C3',''),
    ('321', 'C4',''),
    ('321', 'D1','')
AS t(code, servico, rule);

CREATE OR REPLACE TEMP VIEW table_b AS
SELECT * FROM VALUES
    ('E1', 'A1','aaa'), 
    ('E2', 'A2','bbb'),
    ('F1', 'A3','ccc'),
    ('F2', 'B1','ddd'),
    ('F3', 'B2','eee'),
    ('G1', 'B3','fff'),
    ('G2', 'C1','ggg'),
    ('G3', 'C2','hhh'),
    ('G4', 'C3','iii'),
    ('H1', 'C4','jjj'),
    ('H2', 'D1','kkk')
AS t(cod_1, cod_2, rule);

CREATE OR REPLACE TEMP VIEW code_servicos AS
SELECT code, collect_set(servico) AS servico_set
FROM table_a
GROUP BY code;

执行查询获取结果:

SELECT
    a.code,
    a.servico,
    COALESCE(b.rule, a.rule) AS rule
FROM table_a a
JOIN code_servicos cs ON a.code = cs.code
LEFT JOIN table_b b 
    ON a.servico = b.cod_2 
    AND array_contains(cs.servico_set, b.cod_1);

性能优化说明

  • 减少Shuffle:通过collect_set预计算每个code的servico集合,仅需一次分组shuffle,避免多次关联带来的性能损耗。
  • 广播小表:若表B数据量较小,使用F.broadcast()(PySpark)或Databricks自动广播优化,将表B广播到所有节点,避免大表shuffle。
  • 适配Databricks:上述逻辑完全兼容Databricks的Spark环境,利用其分布式执行引擎最大化性能。

内容的提问来源于stack exchange,提问作者Leonardo Kanashiro Felizardo

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.12 03:01:22