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

Dask Client引发内存溢出问题求助

Dask Client导致内存溢出及大型计算图警告的解决方案

问题背景

正在将部分Pandas代码迁移至Dask,多数代码运行正常,但启用Client()监控资源时出现内存溢出、Worker频繁重启的问题,小数据集测试还出现大型计算图警告。


未使用Client时

代码运行正常,耗时约30秒:

import dask.dataframe as dd
from dask.distributed import Client
from datetime import date, timedelta, datetime

def carregar ():
    df = dd.read_parquet(path=f'{pastafonte}\\Unificado',  parse_dates=['data'])
    df['parceiro'] = df.parceiro.cat.add_categories('Not_available').fillna('Not_available') 
    df['mci'] = df['mci'].fillna(0)
    df['sku'] = df['sku'].fillna(df['marca'].astype(str))
    df['cod_transacao'] = df['cod_transacao'].fillna('Not_available')
    return df
dftotal = carregar()
gerado = dftotal.loc[(dftotal.parceiro != 'Amazon') | ((dftotal.parceiro == 'Amazon') & (dftotal.status == 'indefinido'))] 

df = dftotal.groupby(['produto','parceiro', 'marca',dftotal.data.dt.to_period("M"), 'status'], dropna=False, observed=True).aggregate({'gmv': 'sum', 'receita': 'sum', 'cashback': 'sum'}).reset_index()

dfafiliados = dftotal.loc[(dftotal.produto == 'Afiliados')]
dfafiliados = dfafiliados.groupby(['produto','parceiro', 'marca',dfafiliados.data.dt.to_period("M"), 'status'], dropna=False, observed=True)['cod_transacao'].nunique().reset_index()
dfafiliados = dfafiliados.rename(columns={ 'cod_transacao' : 'qtde_vendas'})
dfafiliados = dfafiliados.loc[dfafiliados.qtde_vendas != 0] #Para substituir o "observed=True" que não está funcionando

dfoutros = dftotal.loc[(dftotal.produto != 'Afiliados')]
dfoutros.cod_transacao = dfoutros.cod_transacao.fillna('Não disponível')
dfoutros = dfoutros.groupby(['produto','parceiro', 'marca',dfoutros.data.dt.to_period("M"), 'status'], dropna=False, observed=True).aggregate({'cod_transacao' :'count'}).reset_index()
dfoutros = dfoutros.rename(columns={ 'cod_transacao' : 'qtde_vendas'})

dfnunique = dd.concat([dfafiliados,dfoutros], axis=0)#, ignore_order = True

df = df.merge(dfnunique,how='left', on=['produto','parceiro', 'marca','data', 'status'])

df['data'] = df['data'].astype({'data' : 'datetime64[ns]'}) # type: ignore
df = df[['produto','parceiro','marca','data','status','qtde_vendas','gmv','receita','cashback']]
df = df.compute()
df = df.sort_values( by = ['produto','parceiro', 'marca','data','status'])

添加Client后

代码出现内存溢出,Worker因内存占用超过95%预算频繁重启:

Client = Client()
Client

日志信息:

2023-10-24 23:22:14,645 - distributed.nanny.memory - WARNING - Worker tcp://127.0.0.1:52395 (pid=17668) exceeded 95% memory budget. Restarting...
2023-10-24 23:22:15,090 - distributed.nanny - WARNING - Restarting worker
2023-10-24 23:22:15,363 - distributed.nanny.memory - WARNING - Worker tcp://127.0.0.1:52401 (pid=4388) exceeded 95% memory budget. Restarting...
2023-10-24 23:22:15,584 - distributed.nanny - WARNING - Restarting worker
2023-10-24 23:22:19,561 - distributed.nanny.memory - WARNING - Worker tcp://127.0.0.1:52402 (pid=26016) exceeded 95% memory budget. Restarting...
2023-10-24 23:22:19,778 - distributed.nanny - WARNING - Restarting worker
2023-10-24 23:22:21,731 - distributed.nanny.memory - WARNING - Worker tcp://127.0.0.1:52398 (pid=18488) exceeded 95% memory budget. Restarting...
2023-10-24 23:22:22,074 - distributed.nanny - WARNING - Restarting worker

---

## 小数据集测试
用小数据集复现问题,未出现溢出,但出现大型计算图警告:

### 未使用Client时
运行耗时极短:
```python
from datetime import datetime
import pandas as pd
import dask.dataframe as dd
from dask.distributed import Client
import numpy as np

num_variables = 1_000_000
rng = np.random.default_rng()

data = pd.DataFrame({
    'id' :  np.random.randint(1,99999,num_variables),
    'date' : [np.random.choice(pd.date_range(datetime(2021,1,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', np.nan]) for i in range(num_variables)],
    'gmv' : rng.random(num_variables) * 100,
    'revenue' : rng.random(num_variables) * 100})

data = data.astype({'product': 'category', 'brand':'category'})
ddf = dd.from_pandas(data, npartitions=5)

df = ddf.groupby([ddf.date.dt.to_period('M'), 'product','brand'], dropna=False, observed=True).aggregate({'id' : 'count'}).reset_index()
df = df.compute()

使用Client时

出现警告:

from datetime import datetime
import pandas as pd
import dask.dataframe as dd
from dask.distributed import Client
import numpy as np

Client = Client()
Client

num_variables = 1_000_000
rng = np.random.default_rng()

data = pd.DataFrame({
    'id' :  np.random.randint(1,99999,num_variables),
    'date' : [np.random.choice(pd.date_range(datetime(2021,1,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', np.nan]) for i in range(num_variables)],
    'gmv' : rng.random(num_variables) * 100,
    'revenue' : rng.random(num_variables) * 100})

data = data.astype({'product': 'category', 'brand':'category'})
ddf = dd.from_pandas(data, npartitions=5)

df = ddf.groupby([ddf.date.dt.to_period('M'), 'product','brand'], dropna=False, observed=True).aggregate({'id' : 'count'}).reset_index()
df = df.compute()

警告信息:

UserWarning: Sending large graph of size 28.62 MiB.
This may cause some slowdown.
Consider scattering data ahead of time and using futures.
warnings.warn(


解决方案

1. 优化计算图大小

  • 提前持久化数据:在加载数据后调用persist(),将数据缓存到Worker内存,避免重复计算和传递大型计算图。例如在carregar()函数末尾添加return df.persist()。
  • 简化链式操作:合并连续操作,减少计算图节点数量。比如把groupby、rename、filter操作尽量合并,避免生成过多中间节点。
  • 提前分发本地数据:小数据集测试中,避免在Client创建后生成大型Pandas DataFrame,可先将数据scatter到Worker,或直接用Dask生成数据替代Pandas。

2. 调整Worker内存配置

  • 设置合理内存限制:创建Client时指定内存参数,根据机器内存调整:
    from dask.distributed import Client
    client = Client(n_workers=2, threads_per_worker=2, memory_limit='8GB')
    
  • 调整内存预警阈值:降低内存监控阈值,给系统留缓冲空间:
    client.run(lambda dask_worker: dask_worker.memory_manager.set_target(0.8))
    

3. 优化分组与聚合操作

  • 预先生成分组键列:将dftotal.data.dt.to_period("M")作为单独列添加到数据中,避免在groupby中重复计算:
    dftotal['mes'] = dftotal.data.dt.to_period("M")
    dftotal = dftotal.persist()
    # 后续groupby使用'mes'列
    df = dftotal.groupby(['produto','parceiro', 'marca','mes', 'status'], ...)
    
  • 精简分组参数:仅在需要包含空值分组时使用dropna=False和observed=True,减少分组数量和内存占用。

4. 调整分区策略

  • 增加分区数量:加载数据时指定npartitions,或用repartition()调整,确保每个分区大小在100MB-1GB之间:
    df = dd.read_parquet(path=f'{pastafonte}\\Unificado', parse_dates=['data'], npartitions=20)
    
  • 对齐分区键:merge前确保两个DataFrame的分区键一致,避免shuffle过程中过度占用内存。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.07 21:32:01