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
相关产品推荐
相关产品推荐

