如何基于Div字段为各Store动态移除PySpark DataFrame中的对应ID?
PySpark实现按门店动态过滤ID集合需求
已知DataFrame结构
- store_df:包含字段
store、ID、Div,存储各门店对应的ID及其所属分区(Div) - final_list:包含字段
Div、ID、Rank、Category,存储不同分区下的ID排名与分类信息
需求说明
针对store_df中的每个门店,根据其所属的Div,移除final_list中属于该Div且存在于该门店对应ID集合中的记录,生成每个门店对应的upd_final_list。
实现步骤
1. 生成门店-Div-ID的排除标记表
从store_df中提取每个门店对应的(Div, ID)组合,作为需要过滤的标记:
# 创建需要排除的(门店, 分区, ID)关联表 exclude_df = store_df.select("store", "Div", "ID")
2. 关联并过滤数据
根据需求分两种场景实现:
场景1:仅过滤全局需排除的记录(不保留门店维度)
直接用left_anti关联,快速筛选出final_list中不在排除列表的记录:
# 左反关联,返回final_list中未被标记为排除的记录 upd_final_list = final_list.join( exclude_df, on=["Div", "ID"], how="left_anti" )
场景2:生成每个门店专属的过滤结果(保留门店维度)
先获取门店与Div的唯一映射,再和final_list交叉关联,最后排除当前门店需移除的ID:
from pyspark.sql.functions import col # 先获取门店-Div的唯一对应关系 store_div_map = store_df.select("store", "Div").distinct() # 交叉关联得到所有门店与final_list的组合,再排除当前门店需移除的ID upd_final_list_per_store = store_div_map.crossJoin(final_list) \ .join( exclude_df, on=["store", "Div", "ID"], how="left_anti" )
代码说明
- left_anti关联:这是PySpark中高效实现“存在性过滤”的方式,直接返回左表中与右表无匹配的记录,完美契合移除指定ID的需求,无需额外判断空值。
- 门店维度结果:通过
distinct()避免门店-Div的重复映射,交叉关联后再过滤,确保每个门店都能得到专属的过滤后列表。
示例验证
假设store_df中有记录:
| store | ID | Div |
|---|---|---|
| 637 | 4000000970 | Pac |
final_list中有记录:
| Div | ID | Rank | Category |
|---|---|---|---|
| Pac | 4000000970 | 1 | A |
| Pac | 4000000971 | 2 | A |
| NY | 4000001000 | 1 | B |
执行场景2的代码后,store 637对应的upd_final_list结果为:
| store | Div | ID | Rank | Category |
|---|---|---|---|---|
| 637 | Pac | 4000000971 | 2 | A |
| 637 | NY | 4000001000 | 1 | B |
内容的提问来源于stack exchange,提问作者Scope
相关产品推荐
相关产品推荐

