如何用Dask高效生成按日期值命名目录的分区Parquet文件?
Dask生成指定分区结构的最佳实践
核心思路
先持久化清洗后的数据,避免重复执行读取和清洗逻辑;再通过分组自定义写入,完全控制目录结构和字段保留。
步骤1:预处理并持久化清洗数据
直接基于原始CSV做后续分组操作会重复触发清洗流程,先把清洗后的数据写入临时存储,后续操作复用这份数据:
import dask.dataframe as dd # 读取原始数据并执行清洗 ddf = dd.read_csv('big_data.csv').map_partitions(clean_data) # 将清洗后的数据写入临时Parquet(可根据数据量调整分区数) ddf.to_parquet('tmp_cleaned', overwrite=True) # 加载已清洗的临时数据 ddf_cleaned = dd.read_parquet('tmp_cleaned')
步骤2:格式化日期并自定义分组写入
把日期转换为YYYYMMDD格式,然后按日期分组,每个分组写入对应目录:
from dask.diagnostics import ProgressBar import os # 将datetime类型的date列转为YYYYMMDD字符串(如果原date是字符串,可直接替换分隔符) ddf_cleaned['date_dir'] = ddf_cleaned['date'].dt.strftime('%Y%m%d') # 定义单组写入逻辑 def write_single_date(group): date_dir = group['date_dir'].iloc[0] target_path = f'test/{date_dir}' os.makedirs(target_path, exist_ok=True) # 写入Parquet,保留所有字段(包括date) group.to_parquet(f'{target_path}/part.0.parquet', index=False) return None # 按日期分组执行写入,指定meta参数让Dask识别返回类型 write_tasks = ddf_cleaned.groupby('date_dir').apply(write_single_date, meta=object) # 启动计算并显示进度 with ProgressBar(): write_tasks.compute()
步骤3:清理临时文件(可选)
临时存储不再需要时,可删除释放空间:
import shutil shutil.rmtree('tmp_cleaned')
方案优势
- 无重复计算:清洗逻辑仅执行一次,后续分组写入基于已处理的数据。
- 完全符合目录要求:生成纯
YYYYMMDD命名的文件夹,而非Hive风格的date=xxx。 - 保留date字段:写入时不会自动移除分区字段,无需后续补全操作。
备选:调整Hive分区格式(若可接受前缀)
如果能接受文件夹带date=前缀,可通过自定义分区文件名回调简化操作,同时调整日期格式:
ddf_cleaned.to_parquet( 'test', partition_on='date', # 自定义分区文件夹和文件名格式 partition_filename_cb=lambda part_str: part_str.split('=')[1].replace('-', '') + '/part.0.parquet', write_index=False )
此方案会生成date=20211201格式的文件夹,无需额外分组操作,但不符合纯日期目录的要求。
内容的提问来源于stack exchange,提问作者Nezo
相关产品推荐
相关产品推荐

