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

Airflow(GCP Composer)调用Cloud Run任务后状态不同步的解决方法

解决Airflow DAG调用Cloud Run后无法同步状态的问题

你的核心问题是:Airflow的PythonOperator仅触发了Cloud Run任务,但没有正确感知任务的执行结果,导致DAG任务一直处于Running状态。以下是针对性的解决方案:

方案1:同步调用Cloud Run(适用于短耗时任务)

如果你的Cloud Run任务执行时间在Cloud Run的超时限制内(最长60分钟),可以修改invoke_with_http函数,让它等待Cloud Run任务执行完成,并根据响应结果判断成功/失败:

修改后的invoke_with_http代码

import requests
from google.auth.transport.requests import Request
from google.oauth2 import id_token

def invoke_with_http(method, service_name, json):
    # 替换为你的Cloud Run服务URL(可从GCP控制台获取)
    service_url = f"https://{service_name}-{ENV}-xxx.a.run.app"
    
    # 生成身份验证用的ID Token(如果Cloud Run开启了身份验证)
    auth_request = Request()
    id_token_obj = id_token.fetch_id_token(auth_request, service_url)
    headers = {"Authorization": f"Bearer {id_token_obj}"}

    # 发送同步POST请求,等待任务完成
    response = requests.request(
        method=method,
        url=service_url,
        json=json,
        headers=headers,
        timeout=3600  # 设置超时,不超过Cloud Run的最大超时时间
    )

    # 非2xx状态码直接抛出异常,Airflow会标记任务为Failed
    response.raise_for_status()

    # 可选:校验返回内容是否符合预期
    if response.text.strip() != "OK":
        raise ValueError("Cloud Run任务未返回预期结果")

    return response.text

当Cloud Run任务执行完成并返回200+状态码时,PythonOperator会正常结束,DAG任务自动标记为Success;如果Cloud Run返回错误状态码或内容不符合预期,函数抛出异常,Airflow会将任务标记为Failed。

方案2:异步触发+轮询状态(适用于长耗时任务)

如果Cloud Run任务执行时间超过同步超时限制,需要改为异步触发+轮询状态的模式:

第一步:修改Cloud Run服务代码

让Cloud Run接收到请求后立即返回任务ID,后台异步执行任务,并提供状态查询接口:

from flask import Flask, request, jsonify
import uuid
import threading
import os

app = Flask(__name__)
# 生产环境建议用Cloud Firestore/Cloud Storage替代字典存储状态
task_status_store = {}

def execute_task(task_id, date, env):
    """后台执行实际任务逻辑"""
    try:
        # 替换为你的任务代码
        print(f"执行任务:日期={date},环境={env}")
        # 模拟任务耗时
        # time.sleep(30)
        task_status_store[task_id] = "SUCCESS"
    except Exception as e:
        task_status_store[task_id] = "FAILED"
        print(f"任务失败:{str(e)}")

@app.route("/", methods=['POST'])
def trigger_task():
    request_data = request.get_json()
    task_id = str(uuid.uuid4())
    # 启动线程异步执行任务
    threading.Thread(
        target=execute_task,
        args=(task_id, request_data['date'], request_data['env'])
    ).start()
    return jsonify({"task_id": task_id}), 202

@app.route("/status/<task_id>", methods=['GET'])
def get_task_status(task_id):
    status = task_status_store.get(task_id, "PENDING")
    return jsonify({"task_id": task_id, "status": status})

if __name__ == "__main__":
    app.run(host="0.0.0.0", port=int(os.environ.get("PORT", 8080)))

第二步:修改Airflow的invoke_with_http函数

触发任务后,定期轮询Cloud Run的状态接口,直到任务完成或超时:

import requests
import time
from google.auth.transport.requests import Request
from google.oauth2 import id_token

def invoke_with_http(method, service_name, json):
    service_url = f"https://{service_name}-{ENV}-xxx.a.run.app"
    auth_request = Request()
    id_token_obj = id_token.fetch_id_token(auth_request, service_url)
    headers = {"Authorization": f"Bearer {id_token_obj}"}

    # 1. 触发任务,获取任务ID
    trigger_resp = requests.post(
        service_url,
        json=json,
        headers=headers,
        timeout=30
    )
    trigger_resp.raise_for_status()
    task_id = trigger_resp.json()["task_id"]

    # 2. 轮询任务状态
    max_poll_times = 60  # 最多轮询60次
    poll_interval = 10   # 每次间隔10秒(总等待时间10分钟)

    for _ in range(max_poll_times):
        status_resp = requests.get(
            f"{service_url}/status/{task_id}",
            headers=headers,
            timeout=30
        )
        status_resp.raise_for_status()
        task_status = status_resp.json()["status"]

        if task_status == "SUCCESS":
            return "任务执行完成"
        elif task_status == "FAILED":
            raise ValueError("Cloud Run任务执行失败")
        # 任务未完成,继续等待
        time.sleep(poll_interval)

    # 轮询超时,抛出异常
    raise TimeoutError("任务执行超时,未在预期时间内完成")

关键注意事项

  • 身份验证:如果Cloud Run服务开启了"需要身份验证",必须携带ID Token请求,否则会返回401错误,上述代码已包含该逻辑。
  • 超时设置:Airflow任务的execution_timeout要大于轮询总时间,避免Airflow先标记任务超时。可在PythonOperator中添加参数:execution_timeout=timedelta(minutes=15)。
  • 状态持久化:示例中用字典存储任务状态,生产环境必须改用Cloud Firestore、Cloud Storage等持久化服务,避免Cloud Run实例重启后丢失状态。
  • 异常处理:确保invoke_with_http中所有错误场景都抛出异常,Airflow才能正确标记任务状态。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.12 21:07:13