如何为从Parquet导入的Dask DataFrame新增列?
问题描述
有一个存储为parquet格式的大数据集,共40个文件,索引为datetime类型。读取DataFrame后想要新增一列,伪代码如下:
import dask.dataframe as dd import dask.delayed from dask.dataframe import read_parquet import random from numpy import array from datetime import datetime # 补充原代码缺失的导入 CUSTOM_DATE_PARSER = lambda x: datetime.strptime(x, "%Y-%m-%d %H:%M:%S") IMPORTED_DF = read_parquet("./data/regression_df.parquet", parse_dates=['time'], date_parser=CUSTOM_DATE_PARSER, index_col='time') new_column = array([random.randint(1,30) for _ in range(len(IMPORTED_DF))])
尝试用assign方法添加时出现报错:
ValueError: Not all divisions are known, can't align partitions. Please use `set_index` to set the index.
尝试重置再设置索引的临时方案仍报错:
df = df.reset_index().set_index('time') # 'time'是索引列名 df.assign(label=dd.from_array(new_column))
对应的错误信息:
TypeError: '<' not supported between instances of 'int' and 'Timestamp'
考虑将数组转为Series后用dd.from_pandas,但不知道chunksize或npartitions的值(原数据集导出前做过NaN删除,row_chunk_size=10000已不准确),询问如何正确为该Dask DataFrame新增列。
解决方案
方法1:用map_partitions生成随机列(推荐)
直接在每个分区上生成对应长度的随机数,无需提前生成全量数组,既节省内存,又自动匹配原DataFrame的分区结构,避免对齐问题:
import dask.dataframe as dd import random from datetime import datetime CUSTOM_DATE_PARSER = lambda x: datetime.strptime(x, "%Y-%m-%d %H:%M:%S") IMPORTED_DF = read_parquet("./data/regression_df.parquet", parse_dates=['time'], date_parser=CUSTOM_DATE_PARSER, index_col='time') def add_random_partition(partition): # 为当前分区生成对应行数的随机整数 partition['new_col'] = [random.randint(1,30) for _ in range(len(partition))] return partition # 对每个分区应用函数,生成带新列的DataFrame result_df = IMPORTED_DF.map_partitions(add_random_partition)
方法2:基于原DataFrame的分区数创建Dask数组
如果已经提前生成了new_column数组,可直接匹配原DataFrame的npartitions参数创建Dask数组,再通过重置索引实现位置对齐:
# 假设已生成new_column数组 new_dask_col = dd.from_array(new_column, npartitions=IMPORTED_DF.npartitions) # 重置索引后按位置新增列,再恢复原索引 temp_df = IMPORTED_DF.reset_index() temp_df = temp_df.assign(new_col=new_dask_col) result_df = temp_df.set_index('time')
方法3:构造匹配索引的Dask Series(仅适用于索引分区已知的情况)
如果原DataFrame的divisions不为None(即索引分区信息明确),可以构造和原索引完全匹配的Dask Series,直接用assign添加:
import pandas as pd # 创建与原DataFrame索引一致的Pandas Series new_series = pd.Series(new_column, index=IMPORTED_DF.index.compute()) # 转为Dask Series,分区数与原DataFrame一致 new_dask_series = dd.from_pandas(new_series, npartitions=IMPORTED_DF.npartitions) # 直接新增列 result_df = IMPORTED_DF.assign(new_col=new_dask_series)
内容的提问来源于stack exchange,提问作者wildcat89
相关产品推荐
相关产品推荐

