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

Cloud Composer(Airflow)定时任务未触发 求助排查配置问题

Cloud Composer Airflow任务未执行排查方案

一、基础调度状态检查

  • 确认Cloud Composer调度器运行状态:在控制台查看调度器是否正常,无频繁重启或报错记录,调度器离线会直接导致任务无法触发。
  • 验证DAG启用状态:在Airflow UI中确认bn_trading_flow DAG处于启用状态(开关为打开状态),暂停的DAG不会触发任务。

二、时间配置验证

  • 你的DAG start_date为2022-10-28 05:00 UTC,schedule_interval=datetime.timedelta(days=1),第一个任务的触发时间应为start_date + schedule_interval,即2022-10-29 05:00 UTC。若当前时间未到该点,任务不会执行;若已过该时间但未触发,因设置了catchup=False,Airflow不会补跑错过的任务,需调整start_date为当前时间之前的合理时间点。

三、DAG解析与依赖检查

  • 检查DAG解析状态:在Airflow UI的Browse > DAGs页面,查看bn_trading_flow的Last Parsed时间是否正常,若解析失败会导致DAG无法调度。可在Cloud Composer日志中搜索DAG解析错误,重点排查:
    • includes/bn_buy_sell_summary_py.py文件是否存在于DAG目录的includes子文件夹中
    • summarize_buy_sell_trans函数导入是否正确,无语法错误
  • 确认Airflow变量function_url配置正确:在Airflow UI的Admin > Variables中检查该变量是否存在且值为有效的Cloud Function URL。

四、HTTP连接配置检查

  • 验证Airflow Connections中的cfn_crypto_trading_to_bq连接:确认连接的URL与function_url变量一致,无额外认证配置冲突(当前代码已自定义Authorization header,无需在连接中设置认证)。

五、代码逻辑潜在问题修正(提前规避执行失败)

当前代码中TOKEN生成逻辑存在风险:TOKEN在DAG定义阶段生成,而非任务运行时,会导致TOKEN过期(Google ID Token有效期1小时),即使任务触发也会执行失败。需修改为任务运行时动态生成:

import datetime
import airflow
from airflow.providers.http.operators.http import SimpleHttpOperator
from airflow.operators import python_operator
import google.oauth2.id_token
import google.auth.transport.requests
from includes.bn_buy_sell_summary_py import summarize_buy_sell_trans

start_date_run = datetime.datetime(2022,10,28,5,0,0)  # UTC Time

default_args = {
    'owner': 'Binance Trading Transaction',
    'depends_on_past': False,
    'email': [''],
    'email_on_failure': False,
    'email_on_retry': False,
    'retries': 1,
    'retry_delay': datetime.timedelta(minutes=5),
    'start_date': start_date_run,
}

def get_auth_headers():
    request = google.auth.transport.requests.Request()
    audience = airflow.models.Variable.get('function_url')
    token = google.oauth2.id_token.fetch_id_token(request, audience)
    return {'Authorization': f"Bearer {token}", "Content-Type": "application/json"}

with airflow.DAG(
        'bn_trading_flow',
        catchup=False,
        default_args=default_args,
        schedule_interval=datetime.timedelta(days=1)
) as dag:
    load_crypto_trading = SimpleHttpOperator(
        task_id= "crypto_trading_to_bq",
        method='POST',
        http_conn_id='cfn_crypto_trading_to_bq',
        data={},
        headers=get_auth_headers,  # 任务运行时动态调用生成headers
    )
    
    summarize_buy_sell_transaction = python_operator.PythonOperator(
        task_id='summarize_buy_sell_trans',
        python_callable=summarize_buy_sell_trans,
    )
    
    load_crypto_trading >> summarize_buy_sell_transaction

六、日志排查

  • 在Cloud Composer控制台的Logging页面,过滤资源类型为cloud-composer、日志名称为scheduler,搜索bn_trading_flow,查看调度器是否有相关报错或提示信息,定位具体未触发原因。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.14 19:40:58