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

