使用Timer装饰器统计多进程耗时结果异常的原因咨询
多进程计时异常问题分析与修复
问题原因
计时器提前打印的核心问题是multi_process函数里的进程等待逻辑存在bug,导致函数本身在第一个进程结束后就提前返回,装饰器的计时代码因此立刻执行,而非等到所有进程完成。
具体来说,你写的while循环计数逻辑有严重漏洞:
- 初始
l设为进程列表的长度 - 每次进入循环,会遍历所有进程,只要进程已结束就给
l减1 - 这意味着同一个已结束的进程会在每轮循环中被重复触发减l操作,比如第一个进程结束后,每轮循环都会让
l减1,很快l就会降到<=0,此时其他进程可能还在运行,但函数已经跳出循环返回了。
修复方案
方案1:使用标准的join()方法(推荐)
这是Python多进程等待所有子进程完成的标准写法,逻辑清晰且无bug:
import functools import time from multiprocessing import Process import pandas as pd # this decorator is used to record the running time. def timer(func): @functools.wraps(func) def wrapper(*args, **kwargs): start_time = time.time() func(*args, **kwargs) stop_time = time.time() cost_time = stop_time-start_time print(f'cost time: {cost_time} s!') return wrapper # caculate_and_save is my target function. @timer def multi_process(): process_list = [] gzdhb_df = pd.read_excel(io='./raw_data/各站点海拔.xlsx') for province in province_list: province_array = gzdhb_df[gzdhb_df['省份']==province].values p = Process(target=caculate_and_save, kwargs={'province': province, 'province_data': province_array}) process_list.append(p) p.start() # 替换原while循环,等待所有进程结束 for p in process_list: p.join()
join()方法会阻塞当前进程,直到调用该方法的子进程结束,遍历所有进程调用join()就能确保所有子进程都完成后,multi_process函数才会返回,装饰器的计时逻辑也会在此时执行。
方案2:修正自定义等待循环逻辑
如果一定要自己实现等待逻辑,需要避免重复计数已结束的进程,比如用集合跟踪已完成的进程:
import functools import time from multiprocessing import Process import pandas as pd # this decorator is used to record the running time. def timer(func): @functools.wraps(func) def wrapper(*args, **kwargs): start_time = time.time() func(*args, **kwargs) stop_time = time.time() cost_time = stop_time-start_time print(f'cost time: {cost_time} s!') return wrapper # caculate_and_save is my target function. @timer def multi_process(): process_list = [] gzdhb_df = pd.read_excel(io='./raw_data/各站点海拔.xlsx') for province in province_list: province_array = gzdhb_df[gzdhb_df['省份']==province].values p = Process(target=caculate_and_save, kwargs={'province': province, 'province_data': province_array}) process_list.append(p) p.start() finished_processes = set() # 等待直到所有进程都完成 while len(finished_processes) < len(process_list): for p in process_list: if not p.is_alive() and p not in finished_processes: finished_processes.add(p) time.sleep(0.1) # 加延迟减少CPU占用
内容的提问来源于stack exchange,提问作者hua XtH
相关产品推荐
相关产品推荐

