不可靠网络下Airflow调用HTTP接口分页拉取数据与超时重试方案咨询
Apache Airflow 实现方案
完全可以实现你提到的所有需求,以下是具体实现思路和示例代码:
核心实现逻辑
推荐使用PythonOperator实现自定义拉取逻辑,灵活度更高,可以同时覆盖重试、分页、超时限制三个需求:
- 504错误重试:可通过
tenacity库的重试装饰器,指定仅针对504网关超时异常进行重试,支持配置重试间隔(比如指数退避避免触发API限流) - 分页拉取:通过循环逐页发起GET请求,每次解析返回的
meta字段判断是否还有下一页,直到currentPage等于totalPage时终止循环 - 总耗时限制:给对应任务设置
execution_timeout参数为15分钟,Airflow会在任务运行超过阈值时自动终止,避免无限重试占用资源
示例代码
from airflow import DAG from airflow.operators.python import PythonOperator from datetime import datetime, timedelta import requests import time from tenacity import retry, stop_after_attempt, wait_exponential, retry_if_result # 定义504错误重试判断 def is_504_error(response): return response.status_code == 504 # 封装带重试的单页GET请求 @retry( retry=retry_if_result(is_504_error), wait=wait_exponential(multiplier=1, min=2, max=10), # 重试间隔2s、4s、8s...最大10s ) def fetch_page(page_num): resp = requests.get( "https://something.com/api/data", params={"page": page_num, "size": 50}, # 按接口实际要求调整分页参数 timeout=30 ) resp.raise_for_status() return resp.json() def pull_full_data(**context): current_page = 1 total_page = 1 # 初始值,第一次请求后会更新 while current_page <= total_page: page_data = fetch_page(current_page) # 解析分页信息 current_page = page_data["meta"]["currentPage"] total_page = page_data["meta"]["totalPage"] # 此处自定义当前页数据处理逻辑,比如写入数据库、存储到对象存储等 process_data(page_data["data"]) # 可选:请求间隔休眠1s,避免触发API限流 time.sleep(1) current_page += 1 default_args = { "owner": "airflow", "start_date": datetime(2024, 1, 1), } with DAG( "pull_api_full_data", default_args=default_args, schedule_interval="0 1 * * *", # 按自身需求调整调度周期 catchup=False ) as dag: pull_task = PythonOperator( task_id="pull_api_data", python_callable=pull_full_data, execution_timeout=timedelta(minutes=15), # 总耗时不超过15分钟 retries=0 # 单任务级别不重试,重试逻辑在单页请求中控制 )
补充注意事项
- 如果需要断点续传,可将已拉取的最大页码存入Airflow Variable或者外部存储,任务重启后从上次失败的页码继续拉取,避免重复拉取全量数据
- 可根据API实际要求调整请求头、参数、单次请求超时时间等配置
- 如果拉取的数据量较大,建议分批写入目标存储,不要全部存在内存中,避免OOM问题
内容的提问来源于stack exchange,提问作者Timothy
相关产品推荐
相关产品推荐

