PySpark join时基于参考DataFrame补全缺失值计算order_share的问题
解决方案
核心修改点
- 原代码仅通过
customer_id单字段内连接订单统计结果,只会保留存在订单记录的地址,丢失了无订单的地址数据 - 调整关联逻辑,以全量地址表
df_address作为主表,关联订单统计结果时使用customer_id、address_id双字段左连接,再对空值填充0即可实现需求
完整可运行代码
from pyspark.sql import functions as f from pyspark.sql.types import FloatType # 原有统计逻辑不变 # 统计每个用户的总地址数 df1 = df_address.groupBy('customer_id').count().select('customer_id',f.col('count').alias('address_counts')) # 统计df_orders中每个地址的订单数 df2 = df_orders.groupBy('customer_id','address_id').count().select('customer_id','address_id',f.col('count').alias('order_count')) # 给全量地址表关联对应用户的总地址数 df_address_with_total = df_address.join(df1, on='customer_id', how='inner') # 左关联订单统计结果,保留所有地址 df = df_address_with_total.join(df2, on=['customer_id', 'address_id'], how='left') # 无订单地址的order_count空值填充为0 df = df.fillna(0, subset=['order_count']) # 用原生算子计算占比,无需自定义UDF,原生算子性能远高于UDF df = df.withColumn("order_share", (f.col('order_count') / f.col('address_counts') * 100).cast(FloatType())) # 查看结果 df.orderBy('customer_id', 'address_id').show()
运行结果说明
104、105、202、302这些订单表中不存在的地址都会出现在结果中,对应order_count为0,order_share也为0,完全符合预期。
内容的提问来源于stack exchange,提问作者Jitesh Malipeddi
相关产品推荐
相关产品推荐

