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

基于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)

原代码问题评估

  1. 并行性未充分利用:在read_and_cal_corr内部调用dask.compute(df_ls)[0],会提前触发该组文件读取任务的同步计算,打断Dask的全局任务调度,失去自动并行读取的优势。
  2. 手动管理文件与日期冗余:通过循环拼接日期生成文件名,代码繁琐且易出错,没有利用Dask内置的批量文件读取能力。
  3. 内存压力风险:最终用pd.concat合并所有组结果,会将所有数据加载到本地内存,数据量变大时容易出现内存溢出。
  4. 任务粒度不合理:将整个组的读取和计算打包成单个延迟任务,Dask无法对组内的单个文件任务进行细粒度调度优化。

优化建议

  1. 改用dask.dataframe.read_csv批量读取文件,自动实现并行读取,替代手动延迟任务。
  2. 避免在延迟任务内部调用dask.compute,让Dask统一管理整个任务图的执行。
  3. 利用Dask的groupby+apply实现分组内的滚动计算,自动调度各组并行执行。
  4. 保留结果为Dask DataFrame,可选择直接保存为分区文件,避免一次性加载所有数据到内存。
  5. 确保时间序列的顺序正确性,避免滚动计算出错。

优化后的代码

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.28 01:18:01