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

Airflow任务并行化实现及Google Ads API调用超时处理求助

解决Airflow中Google Ads API调用无限挂起的超时控制问题

问题背景

从Google Ads获取数据时,API调用偶尔会出现无限挂起的情况。尝试通过独立进程监控超时,但多种方案均存在问题:

  • Threading:无法强制终止挂起线程,标志位机制无效(API挂起后不会执行后续Python代码)
  • multiprocessing:单独运行正常,但Airflow任务中无法生成新进程
  • signals:单独运行有效,但告警信号会发送给Airflow主进程而非任务进程,错误捕获仅在自身任务中生效
  • execute_tasks_new_python_interpreter配置:设置为true后任务卡在调度阶段无法运行
  • PythonVirtualenvOperator:安装virtualenv后Airflow提示找不到包,但激活Airflow自身虚拟环境可正常导入

可行解决方案

方案1:使用subprocess调用独立脚本实现超时控制

通过subprocess创建完全独立的子进程执行API请求,利用其内置的超时参数强制终止挂起进程,Airflow环境下可正常运行。

步骤:

  1. 将Google Ads API调用逻辑封装为独立脚本(如fetch_google_ads_data.py),脚本接收结果存储路径作为参数,将获取到的数据写入该路径:
# fetch_google_ads_data.py
import sys
from google.ads.googleads.client import GoogleAdsClient

def main(output_path):
    # 初始化Google Ads客户端
    client = GoogleAdsClient.load_from_storage()
    # 执行API请求逻辑
    ads_data = client.get_service("GoogleAdsService").search(...)
    # 将结果写入文件
    with open(output_path, 'w') as f:
        f.write(str(ads_data))

if __name__ == "__main__":
    main(sys.argv[1])
  1. 在Airflow任务中调用该脚本,设置超时时间:
import subprocess
import tempfile
import os
from airflow.decorators import task

@task
def fetch_google_ads_task():
    # 创建临时文件存储API结果
    with tempfile.NamedTemporaryFile(mode='w+', delete=False) as temp_file:
        temp_path = temp_file.name
    
    try:
        # 调用独立脚本,设置超时时间(示例为300秒)
        subprocess.run(
            ['python', '/path/to/fetch_google_ads_data.py', temp_path],
            check=True,
            timeout=300,
            capture_output=True,
            text=True
        )
        # 读取并处理结果
        with open(temp_path, 'r') as f:
            ads_data = f.read()
        print("Google Ads数据获取成功")
        # 后续数据处理逻辑...
    except subprocess.TimeoutExpired:
        print("API调用超时,已终止进程")
        # 超时处理逻辑...
    except subprocess.CalledProcessError as e:
        print(f"API调用执行失败: {e.stderr}")
        # 执行错误处理逻辑...
    finally:
        # 清理临时文件
        os.unlink(temp_path)

方案2:基于CeleryExecutor的任务超时配置

如果Airflow使用CeleryExecutor,可通过Airflow和Celery的双重超时配置强制终止挂起任务。

  1. 在DAG中设置任务软超时:
from datetime import timedelta
from airflow import DAG
from airflow.operators.python import PythonOperator

default_args = {
    'owner': 'airflow',
    'execution_timeout': timedelta(minutes=10),  # 软超时10分钟,Airflow发送终止信号
}

with DAG(
    'google_ads_fetch_dag',
    default_args=default_args,
    schedule_interval='@daily',
) as dag:
    def fetch_ads_data():
        # Google Ads API调用逻辑
        client = GoogleAdsClient.load_from_storage()
        # 执行请求...
    
    fetch_task = PythonOperator(
        task_id='fetch_google_ads_data',
        python_callable=fetch_ads_data,
    )
  1. 在Celery配置文件(如celeryconfig.py)中设置硬超时:
task_time_limit = 600  # 硬超时10分钟,Celery强制终止进程
task_soft_time_limit = 540  # 软超时9分钟,提前发送告警信号

方案3:修复PythonVirtualenvOperator的环境问题

若倾向使用VirtualenvOperator,通过指定正确的Python路径解决包找不到的问题:

from airflow.operators.python import PythonVirtualenvOperator

def fetch_ads_data():
    from google.ads.googleads.client import GoogleAdsClient
    # 执行Google Ads API请求逻辑...

fetch_task = PythonVirtualenvOperator(
    task_id='fetch_google_ads_virtualenv',
    python_callable=fetch_ads_data,
    requirements=['google-ads'],
    python_bin='/path/to/airflow/venv/bin/python',  # 指定Airflow虚拟环境的Python路径
)

确保Airflow worker节点的该路径下已安装virtualenv包,或全局安装virtualenv后无需指定python_bin。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.08 08:15:35