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

PySpark中如何比较两个DataFrame数组列并过滤匹配行

解决PySpark数组列匹配筛选问题

你原来的代码存在两个核心问题:

  • 错误使用F.lit(a):a是DataFrame的数组列对象,lit仅支持传入单个常量值,无法直接处理列;且循环中错误引用了整个main_df的列,而非当前迭代的行数据。
  • 采用rdd.toLocalIterator循环处理完全违背Spark分布式计算的设计思路,属于效率极低的反模式。

以下是两种高效的Spark原生处理方案:

方案一:数组交集判断法

通过笛卡尔积关联两个DataFrame,利用array_intersect计算数组交集,过滤交集长度大于0的行:

from pyspark.sql import functions as F

# 关联两个DataFrame
cross_df = main_df.crossJoin(df)
# 筛选数组存在共同元素的行
final_df = cross_df.filter(F.size(F.array_intersect(F.col("refer_array_col"), F.col("array_column"))) > 0)

# 查看结果
final_df.show(truncate=False)

方案二:Explode数组关联法(更适用于大数据量)

将数组列拆分为单行元素,通过元素匹配完成关联,最后去重得到唯一匹配行:

from pyspark.sql import functions as F

# 拆分main_df的数组列
main_exploded = main_df.withColumn("refer_element", F.explode(F.col("refer_array_col")))
# 拆分df的数组列
df_exploded = df.withColumn("array_element", F.explode(F.col("array_column")))

# 通过元素匹配关联,去重保留唯一行
final_df = main_exploded.join(df_exploded, main_exploded.refer_element == df_exploded.array_element) \
    .drop("refer_element", "array_element") \
    .dropDuplicates(["No", "ID"])

# 查看结果
final_df.show(truncate=False)

输出结果说明

两种方案都会得到以下匹配行:

  • main_df中No=1(数组含YYY)匹配df中ID=1A、ID=3C(value-1)的行
  • main_df中No=2(数组含XXX、YYY)匹配df中ID=1A、ID=3C(value-1)的行
  • main_df中No=3、No=4无匹配行

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.06 02:05:18