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

如何实现线程/函数执行超10秒自动重启及主线程时长跟踪?

嘿,我来帮你搞定这两个问题~先提一句:你的示例代码用的是multiprocessing.Pool(进程池),不是线程池,所以下面的方案会基于进程来调整(如果确实需要线程的话,可以换成concurrent.futures.ThreadPoolExecutor,逻辑类似)。


问题1:如何重启运行时长超过10秒的任务(进程/线程)

核心思路是给每个任务设置超时阈值,一旦超时就终止当前任务并重新启动它。因为Pool.map()没法直接设置单任务超时,所以我们需要改用apply_async()来提交任务,这样能单独监控每个任务的执行时长,捕获超时后触发重启逻辑。


问题2:主线程跟踪write_file()执行时长,超10秒则重启

结合你的需求,我修改了原代码,加入了超时监控和重启逻辑,同时优化了文件操作的安全性(用with语句自动关闭文件):

from multiprocessing import Pool
from multiprocessing import TimeoutError

def write_file(file: str):
    # 用with语句自动管理文件,避免异常导致文件句柄泄漏
    with open(file, 'w') as f:
        for item in range(0, 1500000):
            f.write("%s\n" % item)
    print(f"任务 {file} 执行完成")

def run_with_timeout_retry(func, args, timeout=10, max_retries=3):
    """封装超时监控和重试逻辑的工具函数"""
    retries = 0
    while retries < max_retries:
        try:
            # 用单进程池提交任务,设置超时时间
            with Pool(1) as p:
                p.apply_async(func, args).get(timeout=timeout)
            return
        except TimeoutError:
            retries += 1
            print(f"任务 {args[0]} 执行超时(超过{timeout}秒),正在进行第{retries}次重试...")
        except Exception as e:
            print(f"任务 {args[0]} 执行出错:{str(e)},终止重试")
            raise
    print(f"任务 {args[0]} 重试{max_retries}次后仍超时,放弃执行")

if __name__ == '__main__':
    list_files = ['1.txt', '2.txt', '3.txt']
    # 逐个处理文件任务
    for file in list_files:
        run_with_timeout_retry(write_file, (file,))

代码说明:

  • run_with_timeout_retry():封装了超时判断和重试逻辑,你可以自定义timeout(超时阈值)和max_retries(最大重试次数,避免无限循环)。
  • 改用apply_async().get(timeout=):替代原代码的map(),实现单任务的超时监控,捕获TimeoutError后自动触发重试。
  • 文件操作优化:用with open(...)自动关闭文件,比手动调用close()更安全,能避免异常场景下的资源泄漏。

保持并行处理的版本:

如果你想保留原代码的并行性(同时处理3个文件),可以用concurrent.futures来实现,结合as_completed()监控每个任务的超时状态:

from concurrent.futures import ProcessPoolExecutor, TimeoutError, as_completed

def write_file(file: str):
    with open(file, 'w') as f:
        for item in range(0, 1500000):
            f.write("%s\n" % item)
    print(f"任务 {file} 执行完成")

def main():
    list_files = ['1.txt', '2.txt', '3.txt']
    timeout = 10
    max_retries = 3

    with ProcessPoolExecutor(max_workers=3) as executor:
        # 提交所有任务,建立任务与文件的映射关系
        future_to_file = {executor.submit(write_file, file): file for file in list_files}
        
        for future in as_completed(future_to_file):
            file = future_to_file[future]
            retries = 0
            while retries < max_retries:
                try:
                    # 等待任务完成,设置超时
                    future.result(timeout=timeout)
                    break
                except TimeoutError:
                    retries += 1
                    print(f"任务 {file} 超时,第{retries}次重试...")
                    # 重新提交任务
                    future = executor.submit(write_file, file)
                except Exception as e:
                    print(f"任务 {file} 出错:{str(e)}")
                    break
            if retries >= max_retries:
                print(f"任务 {file} 重试{max_retries}次仍超时,放弃")

if __name__ == '__main__':
    main()

这个版本能保持3个任务并行处理,同时监控每个任务的超时情况,超时后自动重启对应任务。


内容的提问来源于stack exchange,提问作者袙谢邪写懈屑懈褉

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.28 07:26:32