FastAPI多进程OCR服务问题:异步等待子进程处理结果
问题描述
计划基于FastAPI搭建多进程服务器,用户上传图片至指定端点后,由3个不同OCR模型并发处理识别文本,待所有模型处理完成后返回3种结果。由于OCR模型加载耗时极长,不想为每次请求创建新进程,需预先初始化模型,通过独立队列分发任务,从共享字典获取结果。用requests替代OCR编写的示例代码因asyncio Future无法跨进程使用,出现阻塞/无限循环问题,寻求非hack式异步等待子进程结果的方案,禁止使用asyncio.sleep轮询。
原示例代码:
import asyncio import multiprocessing from multiprocessing import Manager from multiprocessing.managers import DictProxy # noqa from multiprocessing.queues import Queue manager = Manager() url_queue = manager.Queue() url_shared_dict = manager.dict() def tess_process(url_queue: Queue, text_shared_dict: DictProxy): # Some heavy initialization import requests # starts waiting for tasks while True: if url_queue.empty(): continue url, future = url_queue.get() result = requests.get(url) # adding result to a shared dict, to get it in main process afterward text_shared_dict[url] = result # Here is the main problem # I need to somehow notify a main process that a job is done # and main process can take it. # While subprocess is doing its job, # the main process need to be able to await future.set_result('Task completed') async def async_func(): loop = asyncio.get_event_loop() future = loop.create_future() url = 'https://google.com' url_queue.put((url, future)) await future print(url_shared_dict[url]) if __name__ == '__main__': p = multiprocessing.Process( target=tess_process, args=(url_queue, url_shared_dict) ) p.start() asyncio.run(async_func())
解决方案
方案1:管道+asyncio事件监听
通过管道传递任务结果,将管道的可读事件注册到asyncio事件循环,实现异步通知主进程任务完成。每个任务用唯一ID标识,主进程维护任务ID与future的映射,收到结果后触发future。
import asyncio import multiprocessing from multiprocessing import Manager, Pipe def ocr_worker(task_queue, result_pipe): # 模拟模型初始化(耗时操作) import requests print("Worker 初始化完成") while True: task_id, url = task_queue.get() if task_id is None: break # 退出信号 # 模拟OCR处理逻辑 response = requests.get(url) result = response.text[:100] # 截取部分内容模拟识别结果 # 发送结果到主进程 result_pipe.send((task_id, result)) async def main(): manager = Manager() task_queue = manager.Queue() parent_pipe, child_pipe = Pipe() # 启动3个worker进程(对应3个OCR模型) workers = [] for _ in range(3): p = multiprocessing.Process(target=ocr_worker, args=(task_queue, child_pipe)) p.start() workers.append(p) # 存储任务ID与对应asyncio future的映射 task_futures = {} def handle_result(): try: task_id, result = parent_pipe.recv() future = task_futures.pop(task_id) future.set_result(result) except EOFError: pass # 将管道可读事件注册到asyncio事件循环 loop = asyncio.get_event_loop() loop.add_reader(parent_pipe.fileno(), handle_result) # 模拟FastAPI请求处理函数 async def process_image(url): task_id = id(url) # 生成唯一任务ID future = loop.create_future() task_futures[task_id] = future task_queue.put((task_id, url)) return await future # 测试:获取3个模型的识别结果 urls = ["https://google.com"] * 3 results = await asyncio.gather(*[process_image(url) for url in urls]) print("3个模型的识别结果:") for idx, res in enumerate(results): print(f"模型{idx+1}:{res}") # 关闭worker进程 for _ in range(3): task_queue.put((None, None)) for p in workers: p.join() loop.remove_reader(parent_pipe.fileno()) if __name__ == "__main__": asyncio.run(main())
方案2:ProcessPoolExecutor+asyncio.wrap_future
利用标准库concurrent.futures.ProcessPoolExecutor的initializer参数预先加载模型,复用进程处理任务。通过asyncio.wrap_future将进程池的future转换为asyncio兼容的future,实现异步等待。
import asyncio import requests from concurrent.futures import ProcessPoolExecutor # 全局变量存储模型(每个进程初始化一次) model = None def init_model(): # 模拟模型加载(耗时操作) global model print("加载OCR模型完成") model = "模拟OCR模型实例" def process_image(url): # 模拟模型处理逻辑 response = requests.get(url) return f"[{model}] 识别结果:{response.text[:100]}" async def main(): # 创建进程池,每个进程预先初始化模型 executor = ProcessPoolExecutor(max_workers=3, initializer=init_model) # 模拟FastAPI请求,等待3个模型结果 url = "https://google.com" # 提交3个任务到进程池 sync_futures = [executor.submit(process_image, url) for _ in range(3)] # 转换为asyncio future async_futures = [asyncio.wrap_future(fut) for fut in sync_futures] results = await asyncio.gather(*async_futures) print("3个模型的识别结果:") for res in results: print(res) executor.shutdown() if __name__ == "__main__": asyncio.run(main())
方案3:aiomultiprocess简化异步多进程
使用aiomultiprocess库封装异步多进程逻辑,语法贴近asyncio,支持进程初始化,无需手动处理通信细节。
import asyncio import requests from aiomultiprocess import Pool model = None def init_model(): global model print("加载OCR模型完成") model = "模拟OCR模型实例" async def process_image(url): response = requests.get(url) return f"[{model}] 识别结果:{response.text[:100]}" async def main(): async with Pool(initializer=init_model, processes=3) as pool: url = "https://google.com" results = await pool.map(process_image, [url]*3) print("3个模型的识别结果:") for res in results: print(res) if __name__ == "__main__": asyncio.run(main())
方案对比
- 方案1:完全手动控制进程与通信,灵活性最高,适合需要精细管理任务队列、模型分发的场景,易于集成到FastAPI中。
- 方案2:基于标准库实现,代码简洁,无需额外依赖,适合对进程控制要求不高的场景。
- 方案3:依赖第三方库,代码最简洁,异步语法一致性好,适合快速开发。
内容的提问来源于stack exchange,提问作者ReYaN WTF
相关产品推荐
相关产品推荐

