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

如何在Synapse中用Python线程池限制Azure Function并发请求?

解决方案:分批控制请求量避免Azure Function过载

你的核心问题是未对请求做分批限流——原代码会一次性把所有任务提交到线程池,哪怕线程数设为5,当array_of_objects规模较大时,仍会持续向Azure Function发送请求,最终导致服务崩溃。要实现“每200个请求一批,处理完暂停再继续”的需求,需要手动拆分任务批次,控制每批请求量,处理完一批后再执行暂停。

修改后的代码示例

import time
import threading
from concurrent.futures import ThreadPoolExecutor

def send_request_to_azure_func(item):
    try:
        # 替换为你向Azure Function发送请求的实际逻辑
        print(f"处理请求: {item}, 线程: {threading.currentThread().getName()}")
        # 示例请求逻辑(替换成你的真实调用)
        # import requests
        # response = requests.post("你的Azure Function URL", json=item)
        # response.raise_for_status()
        print(f"请求 {item} 处理完成")
    except Exception as e:
        print(f"处理请求 {item} 出错: {str(e)}")

def process_in_batches(items, batch_size=200, pause_time=60, max_workers=5):
    # 拆分任务为多个批次
    total_batches = (len(items) + batch_size - 1) // batch_size
    for batch_idx in range(total_batches):
        start_idx = batch_idx * batch_size
        end_idx = start_idx + batch_size
        batch = items[start_idx:end_idx]
        print(f"开始处理第 {batch_idx+1}/{total_batches} 批,共 {len(batch)} 个请求")
        
        # 用线程池处理当前批次
        with ThreadPoolExecutor(max_workers=max_workers) as executor:
            executor.map(send_request_to_azure_func, batch)
        
        # 非最后一批处理完成后暂停
        if batch_idx + 1 < total_batches:
            print(f"第 {batch_idx+1} 批处理完成,暂停 {pause_time} 秒...")
            time.sleep(pause_time)

try:
    # 替换为你的实际请求体数组
    array_of_objects = [...]
    # 配置参数:每批200个,暂停60秒,线程池最大5个线程
    process_in_batches(array_of_objects, batch_size=200, pause_time=60, max_workers=5)
except Exception as e:
    print("全局错误:", str(e))

关键改动说明

  • 手动分批:通过切片将大数组拆分为每200个元素的批次,确保同一时间仅处理一批请求,从根源上控制并发量。
  • 批次间暂停:每批处理完成后(最后一批除外)执行休眠,给Azure Function留出资源恢复的时间。
  • 线程池作用于单批次:线程池仅负责当前批次的并发处理,避免一次性提交所有任务导致请求过载。
  • 职责拆分:将原readarray拆分为专门的请求发送函数,逻辑更清晰,便于后续维护和修改请求逻辑。

额外优化建议

  • 根据Azure Function的并发上限调整max_workers和batch_size:如果Function的并发配额较低,可适当减小线程数或批次大小。
  • 添加重试机制:在请求函数中加入重试逻辑(比如捕获超时/5xx错误后重试2-3次),提升请求成功率。
  • 记录请求日志:将请求结果(成功/失败)写入Synapse的存储或日志服务,方便后续排查问题。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.26 05:41:15