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

在HTCondor集群使用Dask运行仿真时,如何处理单个Worker故障?

问题

我在HTCondor集群上用Dask运行仿真任务,当前遇到的问题是:只要有一个仿真任务失败,所有Worker都会被强制关闭。核心代码如下(已移除无关内容):

def simulator_dask(
    args_dict: dict,
    prior: DirectPosterior,
    dataset: DatasetMultichannelArray,
    device: torch.device,
) -> None:

    """
    Execute simulations based on the provided prior distribution in a multithreaded manner using the `Dask` package.

    Args:
        args_dict (Dictionary): Dictionary with the arguments.
        prior (DirectPosterior): Prior distribution.
        dataset (DatasetMultichannelArray): Stores statistics and scaling information used in the prior distribution.
        device (torch.device): Device used to run the script.

    Returns:
        None
    """
    # Create a list to hold delayed computations for each simulation.
    delayed_simulations = []
    run_simulation_delayed = dask.delayed(run_simulation_dask)

    for s in parameter_sets_gen:
        
        # Setting output folder path for each simulation.
        # Note that the numbering of the folders is limited to 6 digits here,
        # i.e., we can only generate simulations below 10 million.
        folder_name = f"{simulation_number:06}"
        
        simulation_args = argparse.Namespace(
            output_dir=folder_name,
        )

        # Create delayed computation for each simulation.
        delayed_simulations.append(
            run_simulation_delayed(
                simulation_args,
            )
        )


    log.info("")
    log.info("***************************************************************")
    log.info("Launching simulations")
    log.info("***************************************************************")

    # Compute the delayed computations, i.e., run the simulations in parallel with HTCondor.
    dask.compute(delayed_simulations)

故障原因在于部分Worker执行run_simulation_dask时失败:节点与主目录连接不稳定,调用shutil.copytree()将节点输出文件夹传输到主目录时找不到目标,老旧节点这个问题尤为突出。我需要避免所有Worker因单个任务失败而停止,让失败的任务能在其他可用Worker上重新执行,或是在同一Worker上重试。

解决方案

1. 给单个任务添加重试机制

直接在dask.delayed中配置重试参数,指定任务失败后的重试次数,单个任务失败时会自动重试,不会影响其他任务运行:

# 调整delayed任务定义,添加retries参数
run_simulation_delayed = dask.delayed(run_simulation_dask, retries=3)  # 重试次数可根据实际情况调整

如果需要针对特定异常(比如文件传输相关的FileNotFoundError)重试,可以在run_simulation_dask函数内部结合tenacity库自定义重试逻辑:

from tenacity import retry, stop_after_attempt, retry_if_exception_type

@retry(stop=stop_after_attempt(3), retry=retry_if_exception_type(FileNotFoundError))
def run_simulation_dask(simulation_args):
    # 原函数逻辑,包括shutil.copytree文件传输部分

2. 捕获任务异常,避免全局崩溃

修改dask.compute的调用方式,使用errors='continue'参数,这样单个任务失败只会返回异常对象,不会终止整个计算流程,其他任务可正常执行。之后可筛选失败任务重新提交:

# 执行任务,允许单个任务失败
results = dask.compute(*delayed_simulations, errors='continue')

# 筛选出失败的任务索引
failed_indices = [i for i, res in enumerate(results) if isinstance(res, Exception)]

# 重新提交失败任务
if failed_indices:
    log.info(f"重新提交 {len(failed_indices)} 个失败任务")
    failed_tasks = [delayed_simulations[i] for i in failed_indices]
    dask.compute(*failed_tasks, errors='continue')

3. 调整Dask集群容错配置

初始化Dask客户端时,配置Worker相关的容错参数,允许Worker失败后重新分配任务到其他节点:

from dask.distributed import Client

client = Client(
    # 你的HTCondor集群连接参数
    worker_ttl=3600,  # 设置Worker闲置超时时间,清理老旧节点资源
    tolerate_failures=True,  # 允许Worker失败,调度器自动将任务分配到其他可用Worker
)

4. 优化文件传输逻辑(从根源减少失败)

针对主目录连接不稳定的问题,优化文件传输环节:

  • 改用rsync替代shutil.copytree(),它支持断点续传,网络波动时更可靠
  • 先将节点上的输出文件打包成压缩包,再进行传输,缩短连接时长
  • 临时将输出存储在节点本地目录,待所有任务完成后再批量同步到主目录

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 07:48:11