Python进程启动开销过高,寻求低延迟并发执行方案
我尝试在Python3中并行运行两个函数,每个函数耗时约30ms,但编写测试脚本后发现,后台进程的启动时间超过100ms,这一开销过高,我希望能避免。请问是否有更快的方式在Python3中实现并发执行(开销更低——理想为几毫秒到十几毫秒),同时仍能在主进程中获取函数的执行结果?
硬件环境:2019款MacBook Pro,Python版本3.10.9,处理器为2GHz四核Intel Core i5。
测试脚本:
import multiprocessing as mp import time import numpy as np def t(s): return (time.perf_counter() - s) * 1000 def run0(s): print(f"Time to reach run0: {t(s):.2f}ms") time.sleep(0.03) return np.ones((1,4)) def run1(s): print(f"Time to reach run1: {t(s):.2f}ms") time.sleep(0.03) return np.zeros((1,5)) def main(): s = time.perf_counter() with mp.Pool(processes=2) as p: print(f"Time to init pool: {t(s):.2f}ms") f0 = p.apply_async(run0, args=(time.perf_counter(),)) f1 = p.apply_async(run1, args=(time.perf_counter(),)) r0 = f0.get() r1 = f1.get() print(r0, r1) print(f"Time to run end-to-end: {t(s):.2f}ms") if __name__ == "__main__": main()
运行输出:
Time to init pool: 33.14ms Time to reach run0: 198.50ms Time to reach run1: 212.06ms [[1. 1. 1. 1.]] [[0. 0. 0. 0. 0.]] Time to run end-to-end: 287.68ms
注:我希望将输出中第2、3行的耗时降低10-20倍。我知道这个目标跨度很大,若无法实现也没关系,但想咨询专业人士是否有可行方法。谢谢!
1. 使用线程而非进程(threading模块)
进程的启动开销远高于线程,因为每个进程需要独立的Python解释器实例和内存空间。线程共享主进程的内存空间,启动和切换开销极低,完全适配你的IO密集型任务(time.sleep属于IO操作)。
修改后的代码示例:
import threading import time import numpy as np from queue import Queue def t(s): return (time.perf_counter() - s) * 1000 def run0(s, result_queue): print(f"Time to reach run0: {t(s):.2f}ms") time.sleep(0.03) result_queue.put(np.ones((1,4))) def run1(s, result_queue): print(f"Time to reach run1: {t(s):.2f}ms") time.sleep(0.03) result_queue.put(np.zeros((1,5))) def main(): s = time.perf_counter() result_queue = Queue() # 创建线程 thread0 = threading.Thread(target=run0, args=(time.perf_counter(), result_queue)) thread1 = threading.Thread(target=run1, args=(time.perf_counter(), result_queue)) print(f"Time to init threads: {t(s):.2f}ms") # 启动线程 thread0.start() thread1.start() # 等待线程结束并获取结果 thread0.join() thread1.join() r0 = result_queue.get() r1 = result_queue.get() print(r0, r1) print(f"Time to run end-to-end: {t(s):.2f}ms") if __name__ == "__main__": main()
预期效果:线程启动耗时通常在1ms以内,Time to reach run0/run1会降到几毫秒级别,完全符合你的预期。
2. 提前初始化进程池(复用进程)
如果必须使用进程(比如函数是CPU密集型,受GIL限制),可以在程序启动时就初始化进程池,而非每次需要并发时才创建。复用已有的进程能彻底避免重复的启动开销。
示例:
import multiprocessing as mp import time import numpy as np # 全局初始化进程池,程序启动时完成创建 pool = mp.Pool(processes=2) def t(s): return (time.perf_counter() - s) * 1000 def run0(s): print(f"Time to reach run0: {t(s):.2f}ms") time.sleep(0.03) return np.ones((1,4)) def run1(s): print(f"Time to reach run1: {t(s):.2f}ms") time.sleep(0.03) return np.zeros((1,5)) def main(): s = time.perf_counter() print(f"Time to init pool (already pre-initialized): {t(s):.2f}ms") f0 = pool.apply_async(run0, args=(time.perf_counter(),)) f1 = pool.apply_async(run1, args=(time.perf_counter(),)) r0 = f0.get() r1 = f1.get() print(r0, r1) print(f"Time to run end-to-end: {t(s):.2f}ms") if __name__ == "__main__": try: main() finally: pool.close() pool.join()
效果:第一次运行可能仍有初始化开销,但后续调用时,进程已存在,Time to reach会大幅降低(通常10ms以内)。
3. 使用concurrent.futures的线程池(更简洁的API)
concurrent.futures.ThreadPoolExecutor提供了更简洁的线程并发API,底层基于线程实现,开销同样极低,代码可读性更高。
示例:
from concurrent.futures import ThreadPoolExecutor import time import numpy as np def t(s): return (time.perf_counter() - s) * 1000 def run0(s): print(f"Time to reach run0: {t(s):.2f}ms") time.sleep(0.03) return np.ones((1,4)) def run1(s): print(f"Time to reach run1: {t(s):.2f}ms") time.sleep(0.03) return np.zeros((1,5)) def main(): s = time.perf_counter() with ThreadPoolExecutor(max_workers=2) as executor: print(f"Time to init executor: {t(s):.2f}ms") future0 = executor.submit(run0, time.perf_counter()) future1 = executor.submit(run1, time.perf_counter()) r0 = future0.result() r1 = future1.result() print(r0, r1) print(f"Time to run end-to-end: {t(s):.2f}ms") if __name__ == "__main__": main()
优势:代码简洁,线程池管理自动化,启动开销与原生threading一致。
关键说明
- 你的任务是IO密集型(
time.sleep),线程完全适用,且IO操作会释放GIL,不会存在并发阻塞问题。 - 如果是CPU密集型任务,线程无法提升效率,此时可考虑提前初始化进程池,或使用
multiprocessing的forkserver启动方式(仍有一定开销,但远低于每次创建新池)。 - 降低并发开销的核心原则是:尽量避免频繁创建销毁进程/线程池,复用已有的执行单元。
内容的提问来源于stack exchange,提问作者dinodeep

