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

如何正确关闭Dask-MPI集群?求助解决运行终止异常问题

Dask-MPI集群正确终止与异常解决

问题场景

在超算InfiniBand网络环境下,通过dask-mpi搭建Dask集群运行示例任务后,终止时出现CommClosedError异常,部分机器还出现进程无法正常结束、工人持续等待连接的问题,且无法使用LocalCluster,需要解决集群正确关闭的问题。

提交的Slurm脚本

import os
import numpy as np
import matplotlib.pyplot as plt
import dask_mpi as dm
import dask.distributed as dd

if __name__ == '__main__':
    dm.initialize(local_directory=os.getcwd())
    with dd.Client() as client:
        x = np.linspace(0, 1, 100)
        y = x * x
        plt.plot(x, y)
        print('Computation complete! Stopping workers...')

    print('Execution complete!')

终止时的异常日志

Computation complete! Stopping workers...
Execution complete!
2023-08-04 15:58:32,268 - distributed.scheduler - INFO - Receive client connection: Client-c2d12c68-32df-11ee-bf3c-a0369f214ab8
2023-08-04 15:58:32,268 - distributed.core - INFO - Starting established connection to tcp://134.94.1.45:55692
2023-08-04 15:58:32,270 - distributed.worker - INFO - Run out-of-band function 'stop'
2023-08-04 15:58:32,271 - distributed.scheduler - INFO - Scheduler closing...
2023-08-04 15:58:32,271 - distributed.scheduler - INFO - Scheduler closing all comms
2023-08-04 15:58:32,271 - distributed.core - INFO - Connection to tcp://134.94.1.46:59274 has been closed.
2023-08-04 15:58:32,271 - distributed.scheduler - INFO - Remove worker <WorkerState 'tcp://134.94.1.46:40107', name: 2, status: running, memory: 0, processing: 0> (stimulus_id='handle-worker-cleanup-1691164712.27189')
2023-08-04 15:58:32,272 - distributed.worker - INFO - Stopping worker at tcp://134.94.1.46:40107. Reason: scheduler-close
2023-08-04 15:58:32,272 - distributed.scheduler - INFO - Lost all workers
2023-08-04 15:58:32,272 - distributed.core - INFO - Received 'close-stream' from tcp://134.94.1.44:39517; closing.
2023-08-04 15:58:32,272 - distributed.batched - INFO - Batched Comm Closed <TCP (closed) Worker->Scheduler local=tcp://134.94.1.46:59274 remote=tcp://134.94.1.44:39517>
Traceback (most recent call last):
  File "/home/strube1/dask/sc_venv_template/venv/lib/python3.10/site-packages/distributed/batched.py", line 115, in _background_send
    nbytes = yield coro
  File "/home/strube1/dask/sc_venv_template/venv/lib/python3.10/site-packages/tornado/gen.py", line 767, in run
    value = future.result()
  File "/home/strube1/dask/sc_venv_template/venv/lib/python3.10/site-packages/distributed/comm/tcp.py", line 269, in write
    raise CommClosedError()
distributed.comm.core.CommClosedError

进程无法结束的日志

Computation complete! Stopping workers...
Execution complete!
distributed.scheduler - INFO - Receive client connection: Client-1616fc2e-32e1-11ee-943d-080038c02b5f
distributed.core - INFO - Starting established connection
distributed.worker - INFO - Run out-of-band function 'stop'
distributed.scheduler - INFO - Scheduler closing...
distributed.scheduler - INFO - Scheduler closing all comms
distributed.worker - INFO -       Start worker at:   tcp://10.13.27.249:44981
distributed.worker - INFO -          Listening to:   tcp://10.13.27.249:44981
distributed.worker - INFO -          dashboard at:         10.13.27.249:43201
distributed.worker - INFO - Waiting to connect to:   tcp://10.13.23.246:45221
distributed.worker - INFO -       Start worker at:   tcp://10.13.27.249:40555
distributed.worker - INFO -          Listening to:   tcp://10.13.27.249:40555
distributed.worker - INFO - -------------------------------------------------
distributed.worker - INFO -               Threads:                          1
distributed.worker - INFO -          dashboard at:         10.13.27.249:44089
distributed.worker - INFO - Waiting to connect to:   tcp://10.13.23.246:45221
distributed.worker - INFO -                Memory:                 503.45 GiB
distributed.worker - INFO - -------------------------------------------------
distributed.worker - INFO -       Local Directory: /p/home/jusers/strube1/juwels/dask/dask-worker-space/worker-jw56i_pm
distributed.worker - INFO -               Threads:                          1
distributed.worker - INFO - -------------------------------------------------
distributed.worker - INFO -                Memory:                 503.45 GiB
distributed.worker - INFO -       Local Directory: /p/home/jusers/strube1/juwels/dask/dask-worker-space/worker-60rszoo4
distributed.worker - INFO - -------------------------------------------------
distributed.worker - INFO - Waiting to connect to:   tcp://10.13.23.246:45221
distributed.worker - INFO - Waiting to connect to:   tcp://10.13.23.246:45221
distributed.worker - INFO - Waiting to connect to:   tcp://10.13.23.246:45221
distributed.worker - INFO - Waiting to connect to:   tcp://10.13.23.246:45221

解决方案

1. 修改脚本,添加显式关闭逻辑

调整后的脚本加入了集群显式关闭、MPI进程终止、InfiniBand接口指定以及Matplotlib后端配置,彻底解决终止异常和进程残留问题:

import os
import numpy as np
import matplotlib
# 指定非交互式后端,避免超算无图形环境下的渲染错误
matplotlib.use('Agg')
import matplotlib.pyplot as plt
import dask_mpi as dm
import dask.distributed as dd
from mpi4py import MPI

if __name__ == '__main__':
    # 初始化Dask-MPI集群,指定InfiniBand接口(根据超算实际接口名调整,如ib0、hfi0)
    dm.initialize(
        local_directory=os.getcwd(),
        interface='ib0'
    )
    
    with dd.Client() as client:
        # 示例计算任务
        x = np.linspace(0, 1, 100)
        y = x * x
        plt.plot(x, y)
        # 保存图像(无图形环境下必须保存而非显示)
        plt.savefig('output.png')
        print('Computation complete! Stopping workers...')
        
        # 显式关闭所有工人和调度器,比with块自动关闭更彻底
        client.shutdown()
    
    # 强制所有MPI进程退出,避免工人进程因网络问题挂起
    MPI.COMM_WORLD.Abort(0)

2. 关键调整说明

  • client.shutdown():主动通知调度器终止所有工人进程并关闭自身,确保集群资源完全释放,避免残留连接导致的CommClosedError。
  • MPI.COMM_WORLD.Abort(0):通过MPI强制所有进程退出,解决部分工人进程因网络通信异常无法正常终止的问题。
  • 指定interface:让Dask使用InfiniBand高速网络通信,减少TCP连接断开时的异常概率,适配超算网络环境。
  • Matplotlib后端配置:超算Slurm任务无图形界面,使用Agg后端避免渲染错误,同时必须用savefig保存图像而非直接显示。

内容的提问来源于stack exchange,提问作者Alexandre Strube

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.14 03:22:04