PySpark批量对比DataFrame:按区域编码校验商品数量一致性
批量区域商品编码数量校验方案
核心思路
以Table2(在售清单)为基准,先分别统计两张表中各区域的商品编码数量,再关联对比生成校验标记;如果只需要指定区域的结果,最后过滤即可。优先用分布式关联查询替代循环,避免多次触发计算。
代码实现(PySpark)
1. 统计各区域商品编码数量
假设需要校验的是唯一商品编码的数量(若要统计总记录数,替换countDistinct为count即可):
from pyspark.sql import functions as F # 统计Table1(购买记录)各区域的唯一商品编码数 df1_cnt = df.groupBy("area_code")\ .agg(F.countDistinct("item_code").alias("buy_item_cnt")) # 统计Table2(在售清单)各区域的唯一商品编码数 df2_cnt = df2.groupBy("area_code")\ .agg(F.countDistinct("item_code").alias("sell_item_cnt"))
2. 关联对比生成校验标记
# 关联两个统计结果,左连接保留Table2所有区域 compare_result = df2_cnt.join(df1_cnt, on="area_code", how="left")\ .fillna(0, subset=["buy_item_cnt"]) # 无购买记录的区域数量填0 # 生成校验标记:数量一致为1,不一致为0 compare_result = compare_result.withColumn( "check_flag", F.when(F.col("sell_item_cnt") == F.col("buy_item_cnt"), 1).otherwise(0) ) # 过滤指定区域列表 area_code_sample_list = [12345, 10000] final_result = compare_result.filter(F.col("area_code").isin(area_code_sample_list))
3. 循环实现方案(不推荐,效率低)
如果必须用循环处理指定区域:
from pyspark.sql import Row area_code_sample_list = [12345, 10000] result_rows = [] for area_code in area_code_sample_list: # 获取Table1该区域的商品编码数,无记录则为0 cnt1 = df.filter(df.area_code == area_code)\ .select(F.countDistinct("item_code")).first()[0] or 0 # 获取Table2该区域的商品编码数 cnt2 = df2.filter(df2.area_code == area_code)\ .select(F.countDistinct("item_code")).first()[0] or 0 # 生成校验标记 flag = 1 if cnt1 == cnt2 else 0 result_rows.append(Row(area_code=area_code, sell_item_cnt=cnt2, buy_item_cnt=cnt1, check_flag=flag)) # 转为DataFrame final_result = spark.createDataFrame(result_rows)
注意事项
- 分布式场景下优先用关联查询:循环会为每个区域单独触发计算,性能远低于一次性统计关联
fillna(0)处理边界情况:当某区域无购买记录时,购买数量视为0,确保对比逻辑完整- 字段名需匹配:上述代码中的
item_code需替换为你实际的商品编码字段名
内容的提问来源于stack exchange,提问作者Hackerds
相关产品推荐
相关产品推荐

