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

PySpark能否使用join替代低性能filter实现高效数据过滤?

业务需求与问题描述

现有存储全量产品、竞品每日观测数据的Spark DataFrame,业务规则要求仅保留近1个月观测次数≥14次的竞品对应的近2个月全量观测数据。
原有实现逻辑为:先通过临时DataFrame筛选出符合规则的竞品与product_id集合,再基于该集合过滤原始DataFrame。实测发现filter(col in list)语法存在严重性能瓶颈,因此需要寻找更低开销的替代实现方案。此前曾考虑使用inner join实现过滤,但担忧该方式会丢失近2个月表中的有效观测记录,需要确认是否存在和inner join开销接近、可仅保留存在于符合条件集合内的product_id对应数据的实现方式。
过滤逻辑参考示意图:
数据过滤逻辑示意图

现有实现代码

待过滤集合生成逻辑

以下代码用于生成后续过滤用的竞品集合,可正常运行但执行效率偏低,最终生成的列表长度为12.4万:

# Adding period column
df_spark_frequent = df_spark_raw.filter((f.col('date') >= dt_1m)
                               & (~f.col('store').isin('somestore', 'secondstore', ))) \
                      .withColumn('period', f.when(f.col('date') > dt_1m, f.lit('1m'))) \
                      .withColumn('xcount', f.count(f.col('period')).over(Window.partitionBy('store', 'product_id', 'period'))) 

df_spark_frequent = df_spark_frequent.filter((f.col('period') == '1m') & (f.col('xcount') > 13)) \
                                     .withColumn('concat', concat(f.col('store'), f.col('product_id')))

# Creating a list of all product ids and stores with >14 observations last month 
concat_l = df_spark_frequent.select('concat').distinct().rdd.flatMap(lambda x: x).collect() # This list has poor performance

数据过滤逻辑

使用上述生成的列表执行过滤时,会出现性能问题甚至任务运行中断,对应代码如下:

df_spark_2m = df_spark_raw.filter((f.col('date') >= dt_2m)
                               & (f.col('sales_price') > 0)
                               & (f.col('store').endswith('dk')) 
                               & (~f.col('store').isin('somestore', 'secondstore', ))) \
                       .withColumn('concat', concat(f.col('store'), f.col('product_id'))).filter(f.col('concat').isin(df_spark_frequent.concat)) # this filter has a poor performance
核心性能瓶颈说明
  • 通过collect()方法将分布式存储的符合条件的竞品数据拉取到Driver端生成本地列表,会产生大量跨节点数据传输开销
  • 12.4万长度的大列表用于isin过滤时,Spark需要将该列表广播到所有Executor节点,广播开销极高,易触发OOM导致任务中断
  • 直接使用inner join的顾虑点:join操作若未去重关联键,可能导致数据膨胀,若关联逻辑有误可能丢失近2个月的有效观测记录

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.27 14:06:26