Airflow调度运行误报Databricks任务失败问题求助
Airflow调度DAG误报Databricks多任务Job失败问题
问题背景
手动触发每日运行的Airflow DAG时,所有任务均成功完成;但调度自动运行时,Databricks中的多任务Job实际执行正常(Kubernetes Pod运行无异常),但Airflow仍标记DAG任务报错,显示一个任务被跳过、另一个完成成功。
DAG代码
from datetime import datetime from airflow.providers.databricks.operators.databricks import DatabricksRunNowOperator from airflow.operators.slack import SlackNotifier default_args = { 'owner': 'airflow', 'depends_on_past': False, 'email_on_failure': False, 'email_on_retry': False, 'retries': 0 } dag = DAG( 'bronze_to_silver_and_silver_to_gold', default_args=default_args, description='DAG to run the steps bronze to silver and silver to gold.', schedule_interval='0 8 * * *', start_date=datetime(2024, 1, 1), catchup=False, max_active_runs=1, concurrency=1, ) slack_failure_notifier = SlackNotifier( slack_conn_id="slack_conn", text=SLACK_FAILURE_MESSAGE, channel="pipeline-status", ) norden_complete_process = DatabricksRunNowOperator( task_id='complete_process', databricks_conn_id='databricks_default', job_id=get_latest_job_id(), python_params=[ "name", "", "tables", "" ], dag=dag, on_failure_callback=slack_failure_notifier, ) complete_process
Databricks Job配置JSON
{ "job_id": , "creator_user_name": "", "run_as_user_name": "", "run_as_owner": false, "settings": { "name": "", "email_notifications": {}, "webhook_notifications": {}, "timeout_seconds": 0, "max_concurrent_runs": 1, "tasks": [ { "task_key": "Bronze_to_silver", "run_if": "ALL_SUCCESS", "python_wheel_task": { "package_name": "data", "entry_point": "bronze_to_silver", "parameters": [ "--client-name", "example_client", "--client-tables", "example_table" ] }, "existing_cluster_id": "", "libraries": [ { "whl": "dbfs:/artifacts/3.0.0-py3-none-any.whl" } ], "timeout_seconds": 0, "email_notifications": {} }, { "task_key": "Silver_to_gold", "depends_on": [ { "task_key": "Bronze_to_silver" } ], "run_if": "ALL_SUCCESS", "python_wheel_task": { "package_name": "data", "entry_point": "silver_to_gold", "parameters": [ "--client-name", "example_client" ] }, "existing_cluster_id": "", "libraries": [ { "whl": "dbfs:/artifacts/3.0.0-py3-none-any.whl" } ], "timeout_seconds": 0, "email_notifications": {} } ], "format": "MULTI_TASK" }, "created_time": 1724685763834 }
具体问题表现
- 手动触发DAG:任务全部成功完成
- 调度自动运行:Databricks端两个任务均执行正常,但Airflow UI显示一个任务成功、一个任务被跳过,且标记DAG失败
- Airflow存在两次尝试后失败的记录,即使Databricks运行无异常
疑问
- 为何手动执行正常,调度运行时Airflow会误报失败?
- 如何阻止DAG重试,或让Airflow正确获取Databricks多任务Job的成功状态?
附图说明
- 图1:Airflow任务状态视图(一个任务跳过,另一个成功)
- 图2:被跳过任务的错误信息详情
- 图3:Airflow DAG调度配置界面
- 图4:Airflow两次尝试失败的错误记录
内容的提问来源于stack exchange,提问作者Antônio Farias
相关产品推荐
相关产品推荐

