Kedro APIDataSet如何实现API请求限流与下载失败重试
api.APIDataSet原生能力完全可以实现失败重试+请求限流,不需要重写自定义数据集类,具体配置方式如下:
1. 失败自动重试+限流配置
APIDataSet底层基于requests库发送请求,原生支持传入自定义Session对象,你可以提前构造带重试、限流策略的会话实例,直接传给数据集即可。
第一步,在项目源码目录下新建工具文件(比如src/[你的项目包名]/api_utils.py),写入逻辑代码:
import time from requests import Session from requests.adapters import HTTPAdapter from urllib3.util.retry import Retry # 基础参数配置 RETRY_COUNT = 3 # 单请求最多重试3次 BACKOFF_FACTOR = 2.0 # 重试退避系数,休眠时间 = factor * (2 ** (重试次数-1)),即首次重试等2s,二次等4s,三次等8s RETRY_STATUS_CODES = (429, 500, 502, 503, 504) # 触发自动重试的HTTP状态码 REQUEST_INTERVAL = 3 # 两次API请求的最小间隔,单位秒 _last_request_ts = 0 def _rate_limit_wait(): """全局请求限流,保证两次请求间隔不小于设定值""" global _last_request_ts wait_sec = REQUEST_INTERVAL - (time.time() - _last_request_ts) if wait_sec > 0: time.sleep(wait_sec) _last_request_ts = time.time() def get_ibge_api_session() -> Session: session = Session() # 挂载重试策略 retry_strategy = Retry( total=RETRY_COUNT, backoff_factor=BACKOFF_FACTOR, status_forcelist=RETRY_STATUS_CODES, allowed_methods=["GET"] ) adapter = HTTPAdapter(max_retries=retry_strategy) session.mount("https://", adapter) session.mount("http://", adapter) # 挂载限流钩子,每次请求拿到响应后等待,控制下一次请求的发起时间 session.hooks["response"].append(lambda r, *args, **kwargs: _rate_limit_wait()) return session # 初始化全局复用的会话实例 ibge_session = get_ibge_api_session()
第二步,修改catalog.yml中的数据集配置,给每个API数据集传入提前构造好的会话:
external-safra-cana: type: api.APIDataSet url: https://apisidra.ibge.gov.br/values/t/6588/p/all/v/allxp/c48/39456/n3/all session: ${importlib.import_module('[你的项目包名].api_utils').ibge_session} external-safra-algodao: type: api.APIDataSet url: https://apisidra.ibge.gov.br/values/t/6588/p/all/v/allxp/c48/39429/n3/all session: ${importlib.import_module('[你的项目包名].api_utils').ibge_session} external-safra-arroz: type: api.APIDataSet url: https://apisidra.ibge.gov.br/values/t/6588/p/all/v/allxp/c48/39432/n3/all session: ${importlib.import_module('[你的项目包名].api_utils').ibge_session} external-safra-milho1: type: api.APIDataSet url: https://apisidra.ibge.gov.br/values/t/6588/p/all/v/allxp/c48/39441/n3/all session: ${importlib.import_module('[你的项目包名].api_utils').ibge_session} external-safra-milho2: type: api.APIDataSet url: https://apisidra.ibge.gov.br/values/t/6588/p/all/v/allxp/c48/39442/n3/all session: ${importlib.import_module('[你的项目包名].api_utils').ibge_session}
注意:配置中的
${importlib.import_module()}是Kedro默认启用的OmegaConf解析器,不需要额外安装依赖即可使用,替换掉路径里的[你的项目包名]为你项目实际的包名即可。运行包含API拉取的流水线时不要加--parallel并行参数,默认串行执行配合上述限流逻辑即可。
2. 根源优化建议
政府公开类API服务承载能力普遍有限,建议首次拉取数据后直接落地到本地存储,后续流水线运行直接读取本地文件,不需要反复请求远端服务,从根源上减少对服务端的压力。可以直接修改数据集配置,将API返回结果落地为本地文件:
external-safra-cana: type: pandas.ParquetDataSet filepath: data/01_raw/safra_cana.parquet layer: raw
首次运行时手动调用API拉取数据存入本地路径后,后续流水线运行会直接读取本地文件,不会再发起远端请求。
内容的提问来源于stack exchange,提问作者João Areias
相关产品推荐
相关产品推荐

