如何为DataFrame.apply()实现异步批量处理并控制请求速率?
百万级数据批量调用外部接口(限30RPS)的优化方案
针对百万条记录调用外部接口的场景,原逐行apply的同步方式效率极低,我们可以通过异步批量请求+速率控制的方式优化,同时严格保证30RPS的调用限制。
核心思路
- 替换同步请求:用
aiohttp实现异步HTTP请求,避免单线程阻塞等待响应的低效问题 - 速率控制:通过
asyncio.Semaphore控制并发数,配合固定间隔等待,把请求速率稳定在30RPS - 批量处理:将DataFrame数据转化为异步任务批量执行,最后统一合并结果到原DataFrame,保证数据对应关系
优化后代码实现
import pandas as pd import mysql.connector import asyncio import aiohttp from typing import List, Dict # 配置项 SERVICE_URL = "你的接口基础地址" HEADERS = {"自定义请求头key": "value"} RPS_LIMIT = 30 # 每秒请求数限制 CONCURRENCY_LIMIT = 30 # 并发数,与RPS匹配避免瞬间打满接口 DATABASE1 = [{"host": "xxx", "user": "xxx", "password": "xxx", "database": "xxx"}] SQL_QUERY = "SELECT service_id, user_id FROM 你的表名" # ---------------------- 数据库读取优化 ---------------------- def generate_df(): # 用列表存储临时DataFrame,最后一次性合并,避免多次concat的性能损耗 df_list = [] for db in DATABASE1: conn = mysql.connector.connect( host=db["host"], user=db["user"], password=db["password"], database=db["database"] ) temp_df = pd.read_sql(SQL_QUERY, conn) df_list.append(temp_df) conn.close() # 及时释放数据库连接 return pd.concat(df_list, ignore_index=True) # ---------------------- 单个请求异步处理 ---------------------- async def fetch_single_status(session: aiohttp.ClientSession, user_id: str, service_id: str) -> Dict: """处理单条记录的接口请求,返回结构化结果""" url = f"{SERVICE_URL}/{service_id}/{user_id}" try: async with session.get(url, headers=HEADERS, ssl=False) as resp: resp_code = resp.status if resp.ok: resp_json = await resp.json() has_active_service = resp_json["data"]["has_active_service"] if resp_json["success"] else "NO_DATA" else: resp_text = await resp.text() has_active_service = f"FAILED_TO_FETCH::{resp_code}::{resp_text}" return { "user_id": user_id, "service_id": service_id, "resp_code": resp_code, "has_active_service": has_active_service } except Exception as e: # 捕获所有异常,避免单个请求失败中断全局任务 return { "user_id": user_id, "service_id": service_id, "resp_code": 500, "has_active_service": f"EXCEPTION::{str(e)}" } # ---------------------- 批量异步调度+速率控制 ---------------------- async def batch_fetch_status(df: pd.DataFrame) -> pd.DataFrame: # 用信号量限制并发数 semaphore = asyncio.Semaphore(CONCURRENCY_LIMIT) async def bounded_fetch(user_id, service_id): async with semaphore: result = await fetch_single_status(session, user_id, service_id) # 控制RPS:每完成一个请求后等待固定时长,确保每秒请求数稳定在30 await asyncio.sleep(1 / RPS_LIMIT) return result async with aiohttp.ClientSession() as session: # 生成所有异步任务 tasks = [ bounded_fetch(row["user_id"], row["service_id"]) for _, row in df.iterrows() ] # 批量执行并收集结果 results = await asyncio.gather(*tasks) # 将结果转为DataFrame并与原数据合并,保证数据对应 result_df = pd.DataFrame(results) merged_df = pd.merge(df, result_df, on=["user_id", "service_id"], how="left") return merged_df # ---------------------- 主函数 ---------------------- if __name__ == "__main__": # 读取原始数据 data_df = generate_df() # 异步批量处理接口请求 processed_df = asyncio.run(batch_fetch_status(data_df)) # 后续的计算操作 # rest of calculative operations on df
关键优化点说明
- 数据库读取优化:用列表缓存临时DataFrame后一次性合并,避免多次
concat产生的内存碎片;及时关闭数据库连接,防止资源泄漏 - 并发与速率控制:
Semaphore限制同时运行的请求数,配合asyncio.sleep(1/RPS_LIMIT)确保请求速率稳定在30RPS,避免触发接口限流 - 异常容错:单个请求失败不会中断全局任务,异常结果会被标记,方便后续排查
- 数据一致性:通过
user_id和service_id作为关联键合并结果,保证原数据与接口返回一一对应
内存友好扩展方案
如果百万条数据导致内存压力过大,可以分块处理数据,避免内存溢出:
chunk_size = 10000 # 每块处理1万条 processed_chunks = [] # 直接用read_sql的chunksize参数分块读取数据库数据 for db in DATABASE1: conn = mysql.connector.connect(**db) for chunk in pd.read_sql(SQL_QUERY, conn, chunksize=chunk_size): processed_chunk = asyncio.run(batch_fetch_status(chunk)) processed_chunks.append(processed_chunk) conn.close() final_df = pd.concat(processed_chunks, ignore_index=True)
内容的提问来源于stack exchange,提问作者Poojan Naik
相关产品推荐
相关产品推荐

