如何用asyncio实现每次N个的异步HTTP请求分批并发处理?
异步分批处理数据库记录并控制请求并发
针对你的一次性任务,核心是用信号量控制并发请求数+数据库分批读取,既避免服务器封禁,又不会一次性加载过多数据到内存。以下是具体实现方案:
核心思路
- 用
asyncio.Semaphore限制同时发起的HTTP请求数(比如50个),确保不超过远程服务器的限制。 - 从数据库分批拉取未处理的记录(目标字段为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
相关产品推荐
相关产品推荐

