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

如何高效保留同支持度项集中的最大项集?

高效筛选同支持度下的最大频繁项集(Spark环境)

需求背景

现有包含ID和Items(字符串列表)列的Spark DataFrame,已通过FP-Growth算法得到频繁项集。需要实现:当多个项集具有相同支持度时,仅保留其中的最大项集(即移除所有是其他项集子集的项,例如["a"]与["a","b"]支持度相同时,仅保留后者)。

原方案按支持度分组后,将同组项集全部加载到内存中进行子集过滤,但面对60个不同项、项集长度可达50的大数据量时,因内存不足崩溃,即使分批处理也无法解决。

优化思路

  1. 避免全量内存加载:放弃将同支持度的所有项集收集到单条记录的方式,改用Spark分布式操作处理每个分组。
  2. 替换Python UDF:使用Spark内置集合函数替代Python自定义函数,减少序列化开销与内存占用。
  3. 逻辑简化:对于同支持度的项集,只需判断是否存在更长的项集包含它——最长项集必然无超集,短项集若被最长项集包含则直接过滤,无需遍历所有组合。

代码实现

1. 生成测试数据并运行FP-Growth

from pyspark.ml.fpm import FPGrowth
from pyspark.sql import functions as F

# 构造测试数据
basic_data = {
    100: [['a', 'b']],
    101: [['a', 'b', 'c']],
    102: [['a', 'b', 'c', 'd']],
    103: [['a', 'b', 'c']],
    104: [['a', 'b']],
    105: [['a', 'b']],
    106: [['c', 'e']],
    107: [['c', 'e']],
    108: [[ 'c', 'e']],
}
data = [(key, value[0]) for key, value in basic_data.items()]
df = spark.createDataFrame(data, schema=['id', 'items'])

# 执行FP-Growth算法
fpGrowth = FPGrowth(itemsCol='items', minSupport=0.3)
model = fpGrowth.fit(df)
frequentItemsets = model.freqItemsets

# 添加项集长度列,用于后续筛选
freq_items = frequentItemsets.withColumn("item_len", F.size(F.col("items")))

2. 分布式筛选最大项集

# 自连接:关联同支持度下长度更长的项集
self_joined = freq_items.alias("a").join(
    freq_items.alias("b"),
    (F.col("a.freq") == F.col("b.freq")) & (F.col("b.item_len") > F.col("a.item_len")),
    how="left"
)

# 用Spark内置函数判断a项集是否是b项集的子集
is_subset_expr = (F.size(F.array_intersect(F.col("a.items"), F.col("b.items"))) == F.size(F.col("a.items")))

# 分组判断每个项集是否存在更长的超集,保留无超集的项集
max_itemsets = self_joined.groupBy("a.freq", "a.items") \
    .agg(F.exists(F.collect_list(is_subset_expr), lambda x: x).alias("has_superset")) \
    .filter(F.col("has_superset") == False) \
    .select("freq", "items")

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

结果验证

运行后得到符合需求的结果:

+----+---------+
|freq|items    |
+----+---------+
|6   |[b, a]   |
|6   |[c]      |
|3   |[c, b, a]|
|3   |[e, c]   |
+----+---------+

性能优势

  • 分布式处理:所有操作基于Spark的分区并行执行,不会将全量项集加载到单节点内存。
  • 内置函数优化:避免Python与JVM之间的序列化开销,计算效率远高于自定义UDF。
  • 时间复杂度降低:无需遍历所有项集组合,仅需判断每个项集是否存在更长的超集,大幅减少计算量。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 03:50:05