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

