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

如何通过SLURM在HPC系统提交作业配置Dask集群并保障生命周期一致?

解决方案:通过SBATCH绑定Dask客户端与集群生命周期

完全可以通过SBATCH指令直接调度整个Dask工作流,让客户端与集群共享同一生命周期,同时确保所有Worker就绪后再启动数据处理。以下是具体实现步骤:

1. 编写SBATCH提交脚本

创建一个名为submit_dask.sh的脚本,内容如下:

#!/bin/bash
#SBATCH --ntasks=500
#SBATCH --ntasks-per-node=25
#SBATCH --time=20:00:00
#SBATCH --job-name=dask-data-processing
# 加载你的Python环境(根据实际情况调整,比如conda或module)
module load python/3.9
source activate dask-env

# 运行数据处理Python脚本
python data_processing.py

2. 编写数据处理Python脚本

创建data_processing.py,在其中完成集群初始化、Worker就绪等待和数据处理逻辑:

from dask.distributed import Client
from dask_jobqueue import SLURMCluster

# 创建匹配SBATCH资源的Dask集群
cluster = SLURMCluster(
    cores=1,  # 每个Worker占用1核,对应总任务数500
    memory="2GB",  # 根据节点总内存调整,确保单节点内存不超配(25*2GB=50GB/节点)
    processes=1,
    walltime="20:00:00",
    scheduler_options={"dashboard_address": ":8787"}  # 可选,开启监控面板
)

# 启动所有Worker
cluster.scale(500)

# 连接客户端
client = Client(cluster)

# 等待所有Worker就绪
print("等待所有Worker启动...")
client.wait_for_workers(n=500)
print("所有Worker已就绪,开始数据处理...")

# ------------------- 数据处理核心逻辑 -------------------
# 示例:读取CSV并计算聚合结果
# import dask.dataframe as dd
# df = dd.read_csv("/path/to/your/data/*.csv")
# aggregated_result = df.groupby("category").sum().compute()
# aggregated_result.to_csv("/path/to/result.csv")
# -------------------------------------------------------

# 处理完成后主动释放资源(可选,SBATCH任务结束后会自动回收)
client.close()
cluster.close()

关键说明

  • 生命周期绑定:整个工作流由SBATCH统一管理,客户端与集群的生命周期完全与SBATCH任务一致,只有当任务达到指定时长或处理完成后才会终止,避免了原Jupyter客户端超时导致集群提前停止的问题。
  • 资源匹配:SLURMCluster的参数需与SBATCH指令严格匹配,比如cores和memory的设置要确保不超过单节点的资源上限。
  • Worker就绪等待:使用client.wait_for_workers(n=500)确保所有Worker启动完成后再开始数据处理,避免因资源未就绪导致的任务延迟或失败。

内容的提问来源于stack exchange,提问作者Axel Wang

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.13 13:01:06