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

如何使用Dask实现分类列转多列的大数据宽表转换

解决方案

你可以直接使用Dask内置的pivot_table方法实现该需求,无需依赖多级索引操作,且由于你的category列仅50种取值,属于低基数维度,该方法执行效率极高,完全适配超大数据集的处理场景。

代码实现

针对你提供的合成示例,完整实现代码如下:

import dask.dataframe as dd
import pandas as pd

# 合成测试数据
records = [
    {'idx': 0, 'id': 'a', 'category': 'd', 'value': 1},
    {'idx': 1, 'id': 'a', 'category': 'e', 'value': 2},
    {'idx': 2, 'id': 'a', 'category': 'f', 'value': 3},
    {'idx': 0, 'id': 'b', 'category': 'd', 'value': 4},
    {'idx': 1, 'id': 'c', 'category': 'e', 'value': 5},
    {'idx': 2, 'id': 'c', 'category': 'f', 'value': 6}
]
frame = pd.DataFrame(records)

# 转为Dask DataFrame,实际使用时直接通过dd.read_avro读取源文件即可
ddf = dd.from_pandas(frame, npartitions=2)

# 执行pivot操作
pivoted_ddf = ddf.pivot_table(
    index=['id', 'idx'],
    columns='category',
    values='value',
    aggfunc='first'  # 相同分组下重复值取第一个,满足去重需求
).reset_index()

# 清除列名标识
pivoted_ddf.columns.name = None

# 测试输出结果,实际处理大数据时可直接写入目标存储,无需全量compute到内存
print(pivoted_ddf.compute())

实际场景适配

针对你最开始带timestamp的业务数据集,仅需修改index参数即可:

pivoted_ddf = ddf.pivot_table(
    index=['id', 'timestamp'],
    columns='category',
    values='value',
    aggfunc='first'
).reset_index()
pivoted_ddf.columns.name = None

说明

  • aggfunc可根据实际需求调整,若同一分组下存在多个有效取值,可替换为'max'/'min'/'mean'等聚合函数
  • 由于category仅50种取值,pivot后的列数极少,不会产生内存爆炸问题,Dask会自动按分区并行处理,无需手动分批

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.01 11:36:05