如何使用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
相关产品推荐
相关产品推荐

