如何基于array_contains合并DataFrame?求快速统计年度词频方法
问题解答
1. 能否基于filtered包含word的条件合并两个DataFrame?
可以实现这种匹配合并,你写的df2.join(df1, F.array_contains(col("filtered"), col("word")), "inner")语法是可行的,但要注意两个细节:
- 这种join本质是笛卡尔积式的匹配,每个
word会和所有包含它的filtered行建立关联,数据量较大时会产生大量中间数据,可能引发性能问题。 - 若两个DataFrame存在同名列(比如都有
id之类的),要提前给列加别名区分,比如df1.withColumnRenamed("filtered", "df1_filtered"),避免列名冲突报错。
另外,你提到DF2的word是唯一关键词,所以不用担心同一word重复匹配的问题。
2. 更快的年度词频统计方法
完全有更高效的方案,不需要先join再做交叉表,直接基于DF1处理即可,步骤如下:
- 用
explode函数把filtered数组拆分成单行记录,同时保留年度列(假设DF1有year列); - 按
year和拆分后的word分组,直接统计次数。
示例代码:
from pyspark.sql import functions as F # 展开filtered数组,关联年度列 word_year_df = df1.select("year", F.explode(F.col("filtered")).alias("word")) # 分组统计年度词频,按年份和词频降序排序 year_word_freq = word_year_df.groupBy("year", "word") \ .count() \ .orderBy("year", F.desc("count")) year_word_freq.show()
这种方法绕开了大表join的性能损耗,explode是轻量的行转列操作,直接基于原数据分组统计,数据量越大,效率优势越明显。如果需要限定只统计DF2中的关键词,在explode后加过滤即可:
# 提取DF2的唯一关键词列表 target_words = [row.word for row in df2.collect()] # 过滤后再统计 word_year_df.filter(F.col("word").isin(target_words)).groupBy("year", "word").count().show()
内容的提问来源于stack exchange,提问作者Anđela Todorović
相关产品推荐
相关产品推荐

