在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
相关产品推荐
相关产品推荐

