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

Dask Delayed与Pandas对比:并行失效及结果格式异常问题排查

Dask活跃客户数统计的问题排查与解决

背景

正在学习用Dask做大数据并行计算,目标是统计交易列表中各周期的活跃客户数(近90天有交易记录的客户),现有Pandas实现和Dask Delayed实现,但遇到两个问题:

  • Dask返回结果是包含元组列表的元组,和Pandas格式不一致
  • 未实现并行,仅单核心高负载

示例数据生成代码

import pandas as pd
import numpy as np
from datetime import date, timedelta, datetime
import dask.dataframe as dd
import dask 

num_variables = 10000
rng = np.random.default_rng()

df = pd.DataFrame({
    'id' :  np.random.randint(1,999999999,num_variables),
    'date' : [np.random.choice(pd.date_range(datetime(2021,6,1),datetime(2022,12,31))) for i in range(num_variables)],
    'product' : [np.random.choice(['giftcards', 'afiliates']) for i in range(num_variables)],
    'brand' : [np.random.choice(['brand_1', 'brand_2', 'brand_4', 'brand_6']) for i in range(num_variables)],
    'gmv': rng.random(num_variables) * 100,
    'revenue': rng.random(num_variables) * 100})

Pandas实现(参考)

def active_clients(df : pd.DataFrame , date : date):
    date1 = (date - timedelta(days=90))
    date2 = date
    clients_base = df.loc[(df['date'].dt.date >= date1) & (df['date'].dt.date <= date2),'id'].nunique()
    return (date, clients_base)

months = []
results = []

dates = df.date.dt.to_period('M').drop_duplicates()
for i in dates:
    test = pd.Period(i,freq='M').end_time.date()
    months.append(test)

for i in months:
    test = active_clients(df,i)
    results.append(test)

results

返回格式:

[(datetime.date(2022, 7, 31), 24),
 (datetime.date(2022, 10, 31), 48),
 (datetime.date(2022, 12, 31), 43),
 (datetime.date(2022, 8, 31), 42),
 (datetime.date(2022, 9, 30), 46),
 (datetime.date(2022, 11, 30), 46),
 (datetime.date(2022, 6, 30), 11)]

Dask Delayed实现的问题与修复

原代码问题分析

@dask.delayed
def active_clients(df : pd.DataFrame , date : date):
    date1 = (date - timedelta(days=90))
    date2 = date
    clients_base = df.loc[(df['date'].dt.date >= date1) & (df['date'].dt.date <= date2),'id'].nunique()
    return (date, clients_base)

months = []
results = []

dates = df.date.dt.to_period('M').drop_duplicates()
for i in dates:
    test = dask.delayed(pd.Period(i,freq='M').end_time.date())
    months.append(test)

for i in months:
    test = dask.delayed(active_clients(df,i))
    results.append(test)

resultados = dask.compute(results)

问题1:结果格式差异

dask.compute()接收可迭代对象时,会返回一个元组,每个元素对应输入的计算结果。这里传入的是results列表,所以返回包含该列表的元组。只需取元组第一个元素即可和Pandas结果格式一致:

resultados = dask.compute(results)[0]

问题2:未并行执行的原因

  1. 不必要的dask.delayed包装:pd.Period(i,freq='M').end_time.date()是轻量计算,无需用dask.delayed包装,反而会增加不必要的任务依赖。
  2. 全局DataFrame依赖:每个active_clients任务都接收整个Pandas DataFrame作为参数,Dask会将这个单一对象视为所有任务的依赖,导致任务只能串行执行。
  3. 未利用Dask DataFrame分区特性:直接使用Pandas DataFrame无法发挥Dask并行优势,应转换为Dask DataFrame并分区。

修复后的Dask实现

# 将Pandas DataFrame转为Dask DataFrame,按核心数设置分区数
ddf = dd.from_pandas(df, npartitions=4)

def get_month_end_dates(df):
    dates = df.date.dt.to_period('M').drop_duplicates()
    return [pd.Period(i,freq='M').end_time.date() for i in dates]

month_ends = get_month_end_dates(df)

# 生成每个日期对应的活跃客户数计算任务
tasks = []
for end_date in month_ends:
    start_date = end_date - timedelta(days=90)
    # 直接用Dask DataFrame API构建计算逻辑
    unique_clients = ddf[(ddf['date'].dt.date >= start_date) & (ddf['date'].dt.date <= end_date)]['id'].nunique()
    tasks.append((end_date, unique_clients))

# 统一调度计算所有任务
computed_results = dask.compute(*[t[1] for t in tasks])
# 组装成与Pandas一致的格式
final_results = [(month_ends[i], computed_results[i]) for i in range(len(month_ends))]

关键优化点

  • 使用Dask DataFrame替代原生Pandas DataFrame,利用分区实现数据并行
  • 移除不必要的dask.delayed包装,直接用Dask DataFrame API构建计算逻辑
  • 让Dask统一调度所有任务,避免单任务单独compute()导致的资源浪费

内容的提问来源于stack exchange,提问作者FábioRB

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.09 17:35:24