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

如何用asyncio实现每次N个的异步HTTP请求分批并发处理?

异步分批处理数据库记录并控制请求并发

针对你的一次性任务,核心是用信号量控制并发请求数+数据库分批读取,既避免服务器封禁,又不会一次性加载过多数据到内存。以下是具体实现方案:

核心思路

  1. 用asyncio.Semaphore限制同时发起的HTTP请求数(比如50个),确保不超过远程服务器的限制。
  2. 从数据库分批拉取未处理的记录(目标字段为NULL的),每批处理完成后再取下一批,避免内存溢出。

代码实现示例

1. 基础依赖导入

import asyncio
import aiohttp
import asyncpg  # 以PostgreSQL为例,其他数据库用对应的异步驱动(如aiomysql)

2. 并发控制与请求处理函数

# 最大并发请求数
MAX_CONCURRENT = 50
semaphore = asyncio.Semaphore(MAX_CONCURRENT)

async def process_record(session, db_conn, record_id, source_value):
    async with semaphore:
        try:
            # 发送HTTP请求
            async with session.get(f"https://your-api-url.com/process?value={source_value}") as resp:
                if resp.status != 200:
                    raise Exception(f"API返回错误状态码: {resp.status}")
                api_result = await resp.json()
                
                # 更新数据库(仅当目标字段仍为NULL时,避免重复更新)
                await db_conn.execute(
                    "UPDATE your_table SET target_column = $1 WHERE id = $2 AND target_column IS NULL",
                    api_result['result'], record_id
                )
        except Exception as e:
            # 异常处理:记录日志或标记失败记录,方便后续重试
            print(f"记录ID {record_id}处理失败: {str(e)}")

3. 分批处理主逻辑

async def process_batches(db_conn, batch_size=1000):
    while True:
        # 拉取一批未处理的记录(用LIMIT避免一次性加载过多数据)
        # 若数据量极大,建议用游标或按ID分段查询(WHERE id > last_id LIMIT batch_size),性能更优
        records = await db_conn.fetch(
            "SELECT id, source_column FROM your_table WHERE target_column IS NULL LIMIT $1",
            batch_size
        )
        
        if not records:
            break  # 无更多未处理记录,结束任务
        
        # 复用HTTP会话,提升请求效率
        async with aiohttp.ClientSession() as session:
            # 创建当前批次的所有任务并等待完成
            tasks = [
                process_record(session, db_conn, rec['id'], rec['source_column'])
                for rec in records
            ]
            await asyncio.gather(*tasks)
        
        print(f"已完成一批{len(records)}条记录的处理")

async def main():
    # 建立数据库异步连接
    db_conn = await asyncpg.connect(
        user="your_user",
        password="your_password",
        database="your_db",
        host="your_db_host"
    )
    
    try:
        await process_batches(db_conn, batch_size=1000)
    finally:
        await db_conn.close()

if __name__ == "__main__":
    asyncio.run(main())

关键细节说明

  • 信号量控制:Semaphore(50)确保同一时间最多有50个HTTP请求在执行,直接限制并发数,避免触发服务器封禁规则。
  • 数据库分批读取:每批拉取1000条(可根据内存调整),避免数百万条数据一次性加载到内存导致OOM。如果数据量极大,建议用游标查询或按ID分段,替代LIMIT,避免偏移量过大导致的性能问题。
  • 异常隔离:单个记录处理失败不会中断整个批次,异常会被捕获并记录,方便后续排查或重试。
  • 会话复用:每批用一个ClientSession,复用TCP连接,减少握手开销,提升请求效率。

额外优化(可选)

如果需要严格限制每秒请求数(而非并发数),可以在信号量内添加延迟:

async with semaphore:
    await asyncio.sleep(1/50)  # 确保每秒最多50个请求
    # 后续请求逻辑...

但这种方式会牺牲部分并发能力,建议优先用并发数限制,根据实际测试调整MAX_CONCURRENT的值。

内容的提问来源于stack exchange,提问作者Marco C. Stewart

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.21 06:43:26