基于DASK分组读取日数据并行计算的代码优化问询
Dask代码评估与优化建议(分组滚动相关性场景)
问题背景
我有100个对应100天数据的文件,需要按每10个文件一组(比如第0-9天、10-19天等)分组,每组内执行滚动相关性计算,最终将所有组的结果聚合为一个DataFrame。作为Dask新手,不确定当前代码是否充分利用了Dask特性,求评估和优化建议。
生成数据代码(无需加速)
import os import datetime import dask df = dask.datasets.timeseries(start='2022-01-01', end='2022-04-11', freq='1d').drop(columns=['name','id']) if not os.path.exists('data'): os.mkdir('data') def name(i): return str(datetime.date(2022, 1, 1) + i * datetime.timedelta(days=1)) df.to_csv('data/*.csv', name_function=name)
原数据处理代码
import pandas as pd from datetime import date, timedelta N_FILES = 100 GROUP_SIZE = 10 CORR_WINDOW = 5 start_date = date(2022, 1, 1) # calculate correlation of two series def rolling_corr(s1: pd.Series, s2: pd.Series, window: int) -> pd.Series: return s1.rolling(window=window).corr(other=s2) # delayed, read one file @dask.delayed def read_file(filename: str) -> pd.DataFrame: return pd.read_csv(filename) # delayed, read several files in parallel and compute correlation @dask.delayed def read_and_cal_corr(start_date: date, ndays: int) -> pd.DataFrame: df_ls = [] cur_date = start_date for _ in range(ndays): df = read_file(f'./data/{str(cur_date)}.csv') df_ls.append(df) cur_date += timedelta(days=1) df_ls = dask.compute(df_ls)[0] df = pd.concat(df_ls) df['corr'] = rolling_corr(df['x'], df['y'], CORR_WINDOW) return df # main, call read_and_cal_corr() on each group and aggregate results df_ls = [] cur_date = start_date for _ in range(N_FILES // GROUP_SIZE): df = read_and_cal_corr(cur_date, GROUP_SIZE) cur_date += timedelta(days=GROUP_SIZE) df_ls.append(df) df_ls = dask.compute(df_ls)[0] result = pd.concat(df_ls) print(result)
原代码问题评估
- 并行性未充分利用:在
read_and_cal_corr内部调用dask.compute(df_ls)[0],会提前触发该组文件读取任务的同步计算,打断Dask的全局任务调度,失去自动并行读取的优势。 - 手动管理文件与日期冗余:通过循环拼接日期生成文件名,代码繁琐且易出错,没有利用Dask内置的批量文件读取能力。
- 内存压力风险:最终用
pd.concat合并所有组结果,会将所有数据加载到本地内存,数据量变大时容易出现内存溢出。 - 任务粒度不合理:将整个组的读取和计算打包成单个延迟任务,Dask无法对组内的单个文件任务进行细粒度调度优化。
优化建议
- 改用
dask.dataframe.read_csv批量读取文件,自动实现并行读取,替代手动延迟任务。 - 避免在延迟任务内部调用
dask.compute,让Dask统一管理整个任务图的执行。 - 利用Dask的
groupby+apply实现分组内的滚动计算,自动调度各组并行执行。 - 保留结果为Dask DataFrame,可选择直接保存为分区文件,避免一次性加载所有数据到内存。
- 确保时间序列的顺序正确性,避免滚动计算出错。
优化后的代码
import dask.dataframe as dd import pandas as pd from datetime import date, timedelta N_FILES = 100 GROUP_SIZE = 10 CORR_WINDOW = 5 start_date = date(2022, 1, 1) # 生成所有文件路径 file_paths = [] cur_date = start_date for _ in range(N_FILES): file_paths.append(f'./data/{str(cur_date)}.csv') cur_date += timedelta(days=1) # 批量读取文件为Dask DataFrame,自动解析时间列 ddf = dd.read_csv(file_paths, parse_dates=['timestamp']) # 生成分组键:按每10天一组划分 def get_group_id(date_val): days_elapsed = (date_val.date() - start_date).days return days_elapsed // GROUP_SIZE ddf['group_id'] = ddf['timestamp'].apply(get_group_id, meta=('group_id', int)) # 定义组内滚动相关性计算函数 def calculate_group_corr(df): # 确保数据按时间排序,滚动窗口计算依赖顺序 df_sorted = df.sort_values('timestamp') df_sorted['corr'] = df_sorted['x'].rolling(window=CORR_WINDOW).corr(other=df_sorted['y']) return df_sorted # 按组应用计算,指定元数据保证类型正确 result_ddf = ddf.groupby('group_id').apply( calculate_group_corr, meta=ddf._meta.append(pd.Series([], name='corr', dtype=float)) ) # 可选:计算结果到本地DataFrame,或直接保存为分区文件 # result_ddf.to_csv('result/*.csv') # 直接保存,无需加载到内存 result = result_ddf.compute() print(result)
优化点说明
- 并行读取自动化:
dd.read_csv会自动将文件读取任务拆分为多个并行任务,由Dask调度执行。 - 全局任务调度:整个任务图由Dask统一管理,避免手动调用
compute打断并行流程。 - 内存友好:结果保留为Dask DataFrame,可直接保存为分区文件,无需一次性加载所有数据到内存,适合大数据场景。
- 逻辑更简洁:通过
groupby+apply实现分组计算,代码结构清晰,减少手动循环的冗余。
内容的提问来源于stack exchange,提问作者OldChi
相关产品推荐
相关产品推荐

