Cloud Composer(Airflow)定时任务未触发 求助排查配置问题
Cloud Composer Airflow任务未执行排查方案
一、基础调度状态检查
- 确认Cloud Composer调度器运行状态:在控制台查看调度器是否正常,无频繁重启或报错记录,调度器离线会直接导致任务无法触发。
- 验证DAG启用状态:在Airflow UI中确认
bn_trading_flowDAG处于启用状态(开关为打开状态),暂停的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
相关产品推荐
相关产品推荐

