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

如何为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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.20 18:40:10