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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.21 03:34:54