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
相关产品推荐
相关产品推荐

