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

Azure Form Recognizer每日千份PDF异步处理优化方案咨询

优化Azure Form Recognizer异步批量处理方案

当前方案的核心问题

一次性启动1000个协程会瞬间打满Form Recognizer的请求配额,引发大量429限流错误;固定1秒休眠的重试逻辑过于僵化,无法适配服务状态的波动,处理确定性差。

推荐优化方案

1. 自适应并发控制(逐步提升+动态调整)

从低并发起步,根据服务响应动态调整并发数,既避免瞬间压垮服务,又能最大化处理效率:

  • 初始并发设为Form Recognizer默认的15TPS(或更保守的10)
  • 每完成一批请求后,若未出现429错误,逐步提升并发量(比如每次加5或10,上限可设为100)
  • 一旦检测到429错误,立即降低并发量(比如减半),同时启用指数退避重试
  • 维护任务队列,持续从队列取任务分配给空闲协程

核心代码示例:

import asyncio
from azure.core.exceptions import HttpResponseError

async def process_doc(doc):
    try:
        result = await analyse_async(doc)
        return (doc, True, result)
    except HttpResponseError as e:
        if e.status_code == 429:
            return (doc, False, "429")
        else:
            return (doc, False, str(e))
    except Exception as e:
        return (doc, False, str(e))

async def adaptive_concurrent_processor(documents, initial_concurrent=15, max_concurrent=100, step=5):
    current_concurrent = initial_concurrent
    task_queue = asyncio.Queue()
    for doc in documents:
        await task_queue.put(doc)
    
    async def worker():
        while not task_queue.empty():
            doc = await task_queue.get()
            status = await process_doc(doc)
            task_queue.task_done()
            return status
    
    while not task_queue.empty():
        worker_count = min(current_concurrent, task_queue.qsize())
        workers = [asyncio.create_task(worker()) for _ in range(worker_count)]
        results = await asyncio.gather(*workers)
        
        # 统计429错误数量
        error_429_count = sum(1 for res in results if res[2] == "429")
        
        if error_429_count > 0:
            # 出现限流,降低并发并重新入队失败任务
            current_concurrent = max(initial_concurrent, current_concurrent // 2)
            for res in results:
                if res[2] == "429":
                    await task_queue.put(res[0])
            # 指数退避等待
            await asyncio.sleep(2 ** (initial_concurrent // current_concurrent))
        else:
            # 无限流,尝试提升并发
            if current_concurrent < max_concurrent:
                current_concurrent = min(max_concurrent, current_concurrent + step)
        
        # 清理已完成的worker任务
        for worker_task in workers:
            worker_task.cancel()

# 调用示例
await adaptive_concurrent_processor(documents)

2. 利用Azure Form Recognizer批量处理API

Form Recognizer原生支持批量分析API,由服务端自动调度处理、管理限流和重试,无需手动控制并发:

  • 使用异步版本begin_analyze_documents_async方法,一次性传入多个文档的SAS URL
  • 通过操作ID轮询或等待处理完成,统一获取所有文档结果

核心代码示例:

from azure.ai.formrecognizer import DocumentAnalysisClient
from azure.core.credentials import AzureKeyCredential
import asyncio

async def batch_process(documents):
    # 初始化客户端
    document_analysis_client = DocumentAnalysisClient(
        endpoint="YOUR_ENDPOINT",
        credential=AzureKeyCredential("YOUR_KEY"),
        api_version="2023-07-31"  # 使用最新API版本
    )
    
    # 生成Azure Storage文档的SAS URL列表
    doc_urls = [get_azure_storage_sas_url(doc) for doc in documents]
    
    # 启动批量分析
    poller = await document_analysis_client.begin_analyze_documents_async(
        "prebuilt-read",  # 根据业务需求选择对应模型
        doc_urls
    )
    
    # 等待处理完成并获取结果
    result = await poller.result()
    
    # 遍历处理每个文档的结果
    for doc_result in result.documents:
        print(f"文档{doc_result.name}处理完成")

# 调用示例
await batch_process(documents)

3. 结合指数退避的Semaphore并发控制

如果要保留自主协程管理逻辑,用asyncio.Semaphore限制最大并发,同时为单个任务添加指数退避重试:

  • 为每个任务设置最大重试次数
  • 重试等待时间按指数增长(1秒→2秒→4秒…)
  • 避免批量休眠,让成功任务快速完成,失败任务单独重试

核心代码示例:

import asyncio
from azure.core.exceptions import HttpResponseError

async def analyse_with_retry(doc, max_retries=5):
    retry_delay = 1
    for attempt in range(max_retries):
        try:
            return await analyse_async(doc)
        except HttpResponseError as e:
            if e.status_code == 429 and attempt < max_retries -1:
                await asyncio.sleep(retry_delay)
                retry_delay *= 2
            else:
                raise
        except Exception as e:
            raise

async def process_all(documents, max_concurrent=100):
    semaphore = asyncio.Semaphore(max_concurrent)
    
    async def bounded_process(doc):
        async with semaphore:
            return await analyse_with_retry(doc)
    
    tasks = [asyncio.create_task(bounded_process(doc)) for doc in documents]
    await asyncio.gather(*tasks)

# 调用示例
await process_all(documents)

方案选择建议

  • 追求省心和稳定性:优先选批量处理API,服务端原生支持批量调度,无需手动管理并发和重试
  • 需要精细控制并发:选自适应并发控制,能根据服务状态动态调整并发数,平衡效率与稳定性
  • 最小改动优化现有代码:选Semaphore+指数退避重试,改动成本低,能有效缓解限流问题

内容的提问来源于stack exchange,提问作者orville

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.10 21:36:01