Pandas优化千万级数据集商品关联分析:新增商品名列并提速
高效计算千万级数据集的关联销售商品(含名称列表)
针对千万级订单数据的关联销售计算,核心是避免循环迭代,改用Pandas的矢量化分组、聚合与合并操作,同时通过数据类型优化降低内存占用,大幅提升处理速度。以下是具体实现方案:
核心思路
- 按订单聚合商品信息,获取每个订单内的所有商品ID、名称及金额
- 生成订单内的商品关联对(排除商品自身关联)
- 统计每个商品的关联商品销量/金额,聚合为列表并排序
- (可选)用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
相关产品推荐
相关产品推荐

