PySpark:如何基于另一DataFrame的类别均价过滤删除低价记录?
解决方法:按类别匹配均价并过滤记录
嘿,我懂你遇到的问题了——你之前写的代码之所以没达到预期效果,核心问题出在collect()[0]["avg(Price)"]这部分:它直接提取了均价表的第一条记录的均价,完全没有实现“每个类别对应各自均价”的匹配逻辑,自然没法正确过滤所有类别的数据。
要解决这个问题,我们需要通过关联(Join)操作把两个DataFrame按Category字段绑定,让每条商品记录都带上对应类别的均价,之后再做过滤就简单了。下面是具体的实现步骤(从代码风格判断你用的是PySpark):
步骤1:关联两个DataFrame
先把商品表data和类别均价表avgCategoryPrice按Category字段做内关联,这样每条商品记录都会新增对应类别的均价列:
# 关联两个表,保留所有类别匹配的记录 joined_df = data.join(avgCategoryPrice, on="Category", how="inner")
步骤2:过滤符合条件的记录
现在每条记录都有Price和对应类别的AVG,直接过滤出价格大于等于均价的记录即可,最后可以按需去掉均价列:
# 过滤价格 >= 对应类别均价的记录 dataGreaterAvge = joined_df.filter(joined_df.Price >= joined_df.AVG) # 如果不需要保留AVG列,选择原表的字段即可 dataGreaterAvge = dataGreaterAvge.select("Category", "Name", "Price")
可选优化:小表广播(如果均价表数据量很小)
如果你的类别均价表avgCategoryPrice数据量不大,可以用广播关联来提升性能,避免大表的shuffle操作:
from pyspark.sql.functions import broadcast joined_df = data.join(broadcast(avgCategoryPrice), on="Category", how="inner") dataGreaterAvge = joined_df.filter(joined_df.Price >= joined_df.AVG).select("Category", "Name", "Price")
这样处理后,就能精准地过滤掉所有价格低于对应类别均价的记录,再也不会只取第一个类别的均价啦~
内容的提问来源于stack exchange,提问作者Maciej Wawrzyniak
相关产品推荐
相关产品推荐

