You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

PySpark按客户分组生成商品购买计数字典实现方法

问题背景

现有一张存储客户购买历史的销售数据表,需要生成按客户维度分组的新DataFrame,新表需包含一列存储该客户所购全部商品的value_counts字典,即记录每个商品对应的购买次数。

当前数据结构

现有DataFrame以CustomerID为索引,包含Description、Counts两个字段,其中Description为逗号分隔的商品名称列表,Counts为当前记录对应的商品总件数,样例数据如下:

Description                                     Counts
CustomerID
3004000304    MAJOR APPLIANCES,HOME OFFICE, OTHER STUFF          3
3004000304    HOME OFFICE, MAJOR APPLIANCES                      2
3004000304    ACCESSORIES, OTHER STUFF                           2
3004002756    MAJOR APPLIANCES, ACCESSORIES                      2
3004002946    HOME OFFICE, HOME OFFICE                           2
3004002946    ACCESSORIES, MAJOR APPLIANCES                      2
3004002946    MAJOR APPLIANCES, OTHER STUFF, ACCESSORIES         3 
期望输出结果

输出DataFrame以CustomerID为索引,仅保留Counts字段,字段值为各客户对应商品购买次数的字典,样例如下:

Counts
CustomerID
3004000304    {'MAJOR APPLIANCES': 2, 'HOME OFFICE': 2, 'ACCESSORIES': 1, 'OTHER STUFF':2}
3004002756    {'MAJOR APPLIANCES': 1, 'ACCESSORIES': 1}
3004002946    {'HOME OFFICE': 2, 'ACCESSORIES': 2, 'MAJOR APPLIANCES': 1,'OTHER STUFF':1}
已尝试方案

由于接触PySpark时间较短、经验不足,先将PySpark DF转换为Pandas DF处理,编写如下lambda函数结合groupby、apply实现聚合,但运行代码未得到预期结果,希望得到可直接使用的PySpark实现方案,可行的Pandas实现方案也可:

f = lambda x: dict(zip(x['Description'], x['Counts']))
df = categories.groupby(level=0).apply(f).to_frame('Counts')
print (df)
实现方案

原有写法的核心问题是没有拆分逗号分隔的商品字段,也没有处理单条记录内多商品的计数平摊、重复商品的次数累加,以下是两种可用的实现:

Pandas 实现

import pandas as pd
from collections import defaultdict

def agg_customer_items(group):
    item_total = defaultdict(int)
    for _, row in group.iterrows():
        # 拆分商品,去除前后空格并过滤空值
        items = [item.strip() for item in row['Description'].split(',') if item.strip()]
        item_cnt_in_row = len(items)
        if item_cnt_in_row == 0:
            continue
        # 计算单条记录中单个商品的分摊计数
        per_item_val = row['Counts'] / item_cnt_in_row
        # 统计单条记录内商品出现频次,累加总计数
        for item, freq in pd.Series(items).value_counts().items():
            item_total[item] += per_item_val * freq
    # 计数转为整数(业务场景下计数均为整数)
    return {k: int(v) for k, v in item_total.items()}

# 分组聚合得到结果
result = categories.groupby(level=0).apply(agg_customer_items).to_frame('Counts')
print(result)

PySpark 实现

from pyspark.sql import functions as F

# 假设源DataFrame名为source_df,若为Pandas转换而来需先将CustomerID从索引转为普通列
# 1. 拆分商品字段,清洗空格
df_step1 = source_df.withColumn(
    "item_arr",
    F.expr("transform(split(Description, ','), x -> trim(x))")
).filter(F.size(F.expr("filter(item_arr, x -> x != '')")) > 0)

# 2. 计算单条记录单商品分摊值,炸开商品列
df_step2 = df_step1.withColumn(
    "per_item_val",
    F.col("Counts") / F.size(F.expr("filter(item_arr, x -> x != '')"))
).withColumn("item", F.explode("item_arr")).filter(F.col("item") != "")

# 3. 按客户+商品维度聚合总计数
df_step3 = df_step2.groupBy("CustomerID", "item").agg(
    F.sum("per_item_val").cast("int").alias("item_cnt")
)

# 4. 按客户维度聚合为字典格式
result = df_step3.groupBy("CustomerID").agg(
    F.map_from_entries(F.collect_list(F.struct("item", "item_cnt"))).alias("Counts")
).orderBy("CustomerID")

result.show(truncate=False)

内容的提问来源于stack exchange,提问作者Abhishek Patil

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.30 21:15:47