如何实现线程/函数执行超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,提问作者袙谢邪写懈屑懈褉
相关产品推荐
相关产品推荐

