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

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运行无异常

疑问

  1. 为何手动执行正常,调度运行时Airflow会误报失败?
  2. 如何阻止DAG重试,或让Airflow正确获取Databricks多任务Job的成功状态?

附图说明

  • 图1:Airflow任务状态视图(一个任务跳过,另一个成功)
  • 图2:被跳过任务的错误信息详情
  • 图3:Airflow DAG调度配置界面
  • 图4:Airflow两次尝试失败的错误记录

内容的提问来源于stack exchange,提问作者Antônio Farias

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.19 03:42:33