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

Pandas优化千万级数据集商品关联分析:新增商品名列并提速

高效计算千万级数据集的关联销售商品(含名称列表)

针对千万级订单数据的关联销售计算,核心是避免循环迭代,改用Pandas的矢量化分组、聚合与合并操作,同时通过数据类型优化降低内存占用,大幅提升处理速度。以下是具体实现方案:

核心思路

  1. 按订单聚合商品信息,获取每个订单内的所有商品ID、名称及金额
  2. 生成订单内的商品关联对(排除商品自身关联)
  3. 统计每个商品的关联商品销量/金额,聚合为列表并排序
  4. (可选)用Dask处理超大数据集,实现并行分块计算

Pandas实现代码

假设你的数据集存储在df中,字段为ORDER_CODE(订单号)、ITEM_ID(商品ID)、ITEM_NAME(商品名称)、TOTALPRICE(商品金额/订单总金额):

import pandas as pd

# 第一步:优化数据类型,减少内存占用并加速操作
df['ITEM_ID'] = df['ITEM_ID'].astype('category')
df['ITEM_NAME'] = df['ITEM_NAME'].astype('category')

# 按订单分组,聚合每个订单的商品列表与总金额
# 若TOTALPRICE是商品单价,sum得到订单总金额;若是订单金额,改用first()
order_groups = df.groupby('ORDER_CODE').agg(
    item_ids=('ITEM_ID', list),
    item_names=('ITEM_NAME', list),
    order_total=('TOTALPRICE', sum)
).reset_index()

# 第二步:生成订单内的商品关联对(排除自身关联)
# 将订单聚合的商品列表拆分为单行记录
order_exploded = order_groups.explode(['item_ids', 'item_names']).rename(
    columns={'item_ids': 'item_id', 'item_names': 'item_name'}
)
# 同订单内商品做笛卡尔积,得到所有关联组合
pair_df = order_exploded.merge(order_exploded, on='ORDER_CODE', suffixes=('_x', '_y'))
# 过滤掉商品与自身的关联
pair_df = pair_df[pair_df['item_id_x'] != pair_df['item_id_y']]

# 第三步:统计关联商品的总金额与销量,并按金额排序聚合
agg_df = pair_df.groupby(['item_id_x', 'item_id_y', 'item_name_y']).agg(
    total_amount=('order_total', sum),  # 关联订单的总金额
    sales_count=('ORDER_CODE', 'nunique')  # 关联销量(订单数)
).reset_index()

# 按商品分组,将关联商品按总金额降序排列后聚合为列表
result = agg_df.sort_values(['item_id_x', 'total_amount'], ascending=[True, False])\
               .groupby('item_id_x').agg(
                   related_item_ids=('item_id_y', list),
                   related_item_names=('item_name_y', list),
                   total_amount=('total_amount', sum)
               ).reset_index().rename(columns={'item_id_x': 'item_id'})

关键优化点

  • 避免迭代:全程用Pandas的groupby、agg、merge等矢量化操作替代循环,这些操作底层用C实现,远快于Python迭代
  • 数据类型优化:将ITEM_ID和ITEM_NAME转为category类型,可减少约70%的内存占用,同时加速分组、合并操作
  • 减少中间数据:通过分步聚合过滤,避免生成不必要的中间数据集

超大数据集处理(千万级以上)

如果Pandas单进程内存不足,改用Dask实现并行分块计算,API与Pandas高度兼容:

import dask.dataframe as dd

# 分块读取数据,自动适配内存
ddf = dd.read_csv('your_data_file.csv', dtype={'ITEM_ID': 'category', 'ITEM_NAME': 'category'})

# 后续步骤与Pandas一致,最后调用compute()生成结果
order_groups = ddf.groupby('ORDER_CODE').agg(
    item_ids=('ITEM_ID', list),
    item_names=('ITEM_NAME', list),
    order_total=('TOTALPRICE', sum)
).reset_index()

order_exploded = order_groups.explode(['item_ids', 'item_names']).rename(
    columns={'item_ids': 'item_id', 'item_names': 'item_name'}
)
pair_df = order_exploded.merge(order_exploded, on='ORDER_CODE', suffixes=('_x', '_y'))
pair_df = pair_df[pair_df['item_id_x'] != pair_df['item_id_y']]

agg_df = pair_df.groupby(['item_id_x', 'item_id_y', 'item_name_y']).agg(
    total_amount=('order_total', sum),
    sales_count=('ORDER_CODE', 'nunique')
).reset_index()

result = agg_df.sort_values(['item_id_x', 'total_amount'], ascending=[True, False])\
               .groupby('item_id_x').agg(
                   related_item_ids=('item_id_y', list),
                   related_item_names=('item_name_y', list),
                   total_amount=('total_amount', sum)
               ).reset_index().rename(columns={'item_id_x': 'item_id'})

# 触发计算,得到最终结果DataFrame
result_df = result.compute()

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 21:55:15