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

