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

