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

PySpark 2.4通过join检测ceci_stok值是否同时存在于ceci_l、ceci_p列

问题说明

环境:PySpark 2.4
现有输入DataFrame结构如下:

ceci_p| ceci_l|ceci_stok|
-------+-------+---------+
SFIL401| BPI202|   BPI202|
 BPI202| CDC111|   BPI202|
 LBP347|SFIL402|  SFIL402|
 LBP347|SFIL402|   LBP347|
-------+-------+---------+

需求:通过join操作(支持自连接)筛选ceci_stok列中同时在全表ceci_p列、ceci_l列出现过的取值,最终返回仅包含符合要求的ceci_stok值的新DataFrame。示例中符合要求的值为BPI202。

实现方案

核心逻辑是通过内连接取集合交集,全程用join实现,符合要求,PySpark 2.4版本可直接运行。

步骤1:构造测试数据(已有生产表可跳过)

from pyspark.sql import SparkSession
from pyspark.sql.functions import col

spark = SparkSession.builder.appName("stok_check").getOrCreate()

# 匹配示例的输入数据
source_data = [
    ("SFIL401", "BPI202", "BPI202"),
    ("BPI202", "CDC111", "BPI202"),
    ("LBP347", "SFIL402", "SFIL402"),
    ("LBP347", "SFIL402", "LBP347")
]
df = spark.createDataFrame(source_data, schema=["ceci_p", "ceci_l", "ceci_stok"])

步骤2:核心join逻辑

先分别提取ceci_p、ceci_l、ceci_stok三列的去重值,再通过两次内连接保留同时在p、l列存在的stok值:

# 提取三列的去重值,统一列名用于join
stok_distinct = df.select(col("ceci_stok").alias("target_val")).distinct()
p_distinct = df.select(col("ceci_p").alias("target_val")).distinct()
l_distinct = df.select(col("ceci_l").alias("target_val")).distinct()

# 连续内连接取交集
result = stok_distinct \
    .join(p_distinct, on="target_val", how="inner") \
    .join(l_distinct, on="target_val", how="inner") \
    .select(col("target_val").alias("ceci_stok"))

结果验证

执行result.show()输出如下,符合预期:

+---------+
|ceci_stok|
+---------+
|   BPI202|
+---------+

注意事项

  • 所有临时数据集必须加distinct()去重,否则大表join会产生数据膨胀,严重影响性能
  • 内连接逻辑天然保留两边同时存在的匹配值,不需要额外写过滤条件
  • 不推荐在PySpark 2.4版本中使用IN子查询实现同类需求,该版本对IN子查询的优化不如join稳定,大表场景容易出现执行计划异常

内容的提问来源于stack exchange,提问作者amine jisung

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.03 00:12:53