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

API数据抽取函数优化请求:50.5万条数据拖慢Airflow DAG运行

优化API数据抽取函数的可行方案

针对你的extract_people函数在抽取50万条API数据时的性能问题,以下是具体优化方案,先修复核心bug再从并发、内存、效率等维度优化:


1. 修复分页逻辑核心bug

原代码存在两个致命问题:定义了offset但未在请求URL中传入,导致每次请求重复获取相同数据;limit变量未定义,分页逻辑完全失效。必须修正分页参数传递:

# 初始化分页参数,offset从0开始避免漏数据
offset = 0
limit = 500  # 每页获取的数量,根据API限流规则调整

# 构造带分页参数的URL
url = f"https://api.apilayer.com/unogs/search/people?person_type=Director&offset={offset}&limit={limit}"

2. 改用异步请求提升并发效率

串行请求单页数据的速度瓶颈明显,改用aiohttp异步发送请求,同时控制并发数避免触发API限流:

import aiohttp
import asyncio
import pandas as pd
from datetime import datetime
from tenacity import retry, stop_after_attempt, wait_exponential

@retry(
    stop=stop_after_attempt(3),
    wait=wait_exponential(multiplier=1, min=2, max=10)
)
async def fetch_page(session, offset, limit):
    url = f"https://api.apilayer.com/unogs/search/people?person_type=Director&offset={offset}&limit={limit}"
    headers = {"apikey": '******'}
    async with session.get(url, headers=headers) as response:
        response.raise_for_status()
        data = await response.json()
        # 批量提取字段,减少字典创建开销
        return pd.DataFrame({
            "netflix_id": [res["netflix_id"] for res in data["results"]],
            "full_name": [res["full_name"] for res in data["results"]],
            "person_type": [res["person_type"] for res in data["results"]],
            "title": [res["title"] for res in data["results"]]
        }) if data["results"] else None

async def extract_people_async():
    offset = 0
    limit = 500
    max_concurrent = 5  # 控制并发数,根据API限流规则调整
    csv_file_name = f'Netflix_people_{datetime.now().strftime("%Y%m%d%H%M%S")}.csv'
    first_write = True

    async with aiohttp.ClientSession() as session:
        tasks = []
        while True:
            # 控制并发任务数量
            if len(tasks) < max_concurrent:
                tasks.append(fetch_page(session, offset, limit))
                offset += limit
                continue
            
            # 等待任一任务完成,处理结果
            done, pending = await asyncio.wait(tasks, return_when=asyncio.FIRST_COMPLETED)
            for task in done:
                page_df = task.result()
                if page_df is None:
                    # 无数据返回,取消剩余任务并退出循环
                    for p in pending:
                        p.cancel()
                    break
                # 分块写入CSV,避免内存占用过高
                page_df.to_csv(
                    csv_file_name,
                    index=False,
                    mode='w' if first_write else 'a',
                    header=first_write
                )
                first_write = False
                print(f'Loaded {offset} records')
            
            tasks = list(pending)
            if page_df is None:
                break

    print(f"CSV file {csv_file_name} saved successfully.")
    return pd.read_csv(csv_file_name)

# 调用异步函数
# asyncio.run(extract_people_async())

3. 优化内存占用:分块写入CSV

50万条数据全部存在列表中会占用大量内存,改为每获取一页就写入CSV,无需等待全部数据收集完成:

# 初始化CSV表头
page_df.to_csv(csv_file_name, index=False, mode='w')
first_write = False

# 后续每页追加写入
page_df.to_csv(csv_file_name, index=False, mode='a', header=False)

4. 复用HTTP连接减少开销

使用requests.Session(同步)或aiohttp.ClientSession(异步)复用TCP连接,避免每次请求重新握手,提升请求速度:

# 同步方式示例
with requests.Session() as session:
    session.headers.update({"apikey": '******'})
    while True:
        response = session.get(url)
        # 处理响应...

5. 优化数据提取逻辑

避免循环中创建大量字典,直接按字段批量提取后构造DataFrame,比字典列表转DataFrame效率更高:

# 批量提取字段
netflix_ids = [res["netflix_id"] for res in data["results"]]
full_names = [res["full_name"] for res in data["results"]]
person_types = [res["person_type"] for res in data["results"]]
titles = [res["title"] for res in data["results"]]

# 直接构造DataFrame
page_df = pd.DataFrame({
    "netflix_id": netflix_ids,
    "full_name": full_names,
    "person_type": person_types,
    "title": titles
})

6. 添加错误处理与重试机制

API请求可能因网络波动、限流失败,添加重试机制避免DAG意外中断:

使用tenacity库实现指数退避重试:

from tenacity import retry, stop_after_attempt, wait_exponential, retry_if_exception_type

@retry(
    stop=stop_after_attempt(3),
    wait=wait_exponential(multiplier=1, min=2, max=10),
    retry=retry_if_exception_type(requests.exceptions.RequestException)
)
def fetch_page(session, offset, limit):
    # 请求逻辑...

完整优化后的同步版本代码

import requests
import pandas as pd
from datetime import datetime
from tenacity import retry, stop_after_attempt, wait_exponential, retry_if_exception_type

def extract_people():
    offset = 0
    limit = 500
    headers = {"apikey": '******'}
    csv_file_name = f'Netflix_people_{datetime.now().strftime("%Y%m%d%H%M%S")}.csv'
    first_write = True

    @retry(
        stop=stop_after_attempt(3),
        wait=wait_exponential(multiplier=1, min=2, max=10),
        retry=retry_if_exception_type(requests.exceptions.RequestException)
    )
    def fetch_page(session, offset, limit):
        url = f"https://api.apilayer.com/unogs/search/people?person_type=Director&offset={offset}&limit={limit}"
        response = session.get(url)
        response.raise_for_status()
        return response.json()

    with requests.Session() as session:
        session.headers.update(headers)
        while True:
            data = fetch_page(session, offset, limit)
            if not data["results"]:
                break
            
            # 批量提取字段
            netflix_ids = [res["netflix_id"] for res in data["results"]]
            full_names = [res["full_name"] for res in data["results"]]
            person_types = [res["person_type"] for res in data["results"]]
            titles = [res["title"] for res in data["results"]]

            page_df = pd.DataFrame({
                "netflix_id": netflix_ids,
                "full_name": full_names,
                "person_type": person_types,
                "title": titles
            })

            # 分块写入CSV
            page_df.to_csv(
                csv_file_name,
                index=False,
                mode='w' if first_write else 'a',
                header=first_write
            )
            first_write = False

            offset += limit
            print(f'Loaded {offset} records')

    print(f"CSV file {csv_file_name} saved successfully.")
    return pd.read_csv(csv_file_name)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.27 18:51:15