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

Airflow技术问询:如何将BaseOperator脚本变量传入失败回调

Solution: Pass Runtime Variable from Containerized Script to Airflow Callback

Great question! Let's break down how to solve this step by step—your core challenge is getting a variable generated during task execution (inside your container running script1.py) to a callback that triggers after the task fails. Your original setup tries to use the variable during DAG parsing (way before the task runs), and directly calling the callback immediately executes it on DAG load, not when the task fails.

Here's the fix using Airflow's built-in XCom (the standard way to share data between tasks, tasks, and callbacks):

Step 1: Push variable1 to XCom from script1.py

First, modify your container script to send variable1 to Airflow's XCom storage. Choose one of these reliable methods:

Add this to script1.py to explicitly push the variable to XCom. Airflow sets environment variables in containers that let you reference the current DAG and task IDs:

from airflow.models.xcom import XCom
from airflow.utils.session import create_session
import os

# Your existing logic to generate variable1
variable1 = "your_dynamic_value_here"

# Fetch DAG/task IDs from Airflow's auto-set environment variables
dag_id = os.environ.get("AIRFLOW_CTX_DAG_ID")
task_id = os.environ.get("AIRFLOW_CTX_TASK_ID")

# Push variable1 to XCom
with create_session() as session:
    XCom.set(
        key="variable1",
        value=variable1,
        task_id=task_id,
        dag_id=dag_id,
        session=session
    )

Note: Ensure your container has access to Airflow's metadata database (the same DB your Airflow instance uses) so the XCom API can write data to it.

Method 2: Print to stdout (Simpler for basic use cases)

If your operator supports auto-capturing stdout as XCom (like KubernetesPodOperator with do_xcom_push=True), just print variable1 at the end of script1.py:

# Your existing logic to generate variable1
variable1 = "your_dynamic_value_here"
print(variable1)

Airflow will capture this output as the return_value XCom key.

Step 2: Update the Callback Function to Fetch from XCom

Rewrite your callback to accept Airflow's context parameter—this gives you access to the task instance and its stored XCom data. Don't call the callback directly in the DAG definition—pass the function reference instead:

def callback_function(context):
    # Get the task instance from the context object
    task_instance = context["task_instance"]
    
    # Fetch variable1 from XCom (use "return_value" if you used Method 2 above)
    variable1 = task_instance.xcom_pull(task_ids="first-task", key="variable1")
    
    # Add your custom callback logic here
    print(f"Task failed! Received variable1 from the task: {variable1}")
    # Example: Send an alert with variable1, log to an external system, etc.

with DAG(
    dag_id=DagName,
    default_args=default_args,
    schedule_interval='12 * * * *',
    on_failure_callback=callback_function,  # Pass the function reference, not a direct call
) as dag:
    first_task = BaseOperator(
        task_id="first-task",  # Corrected: Airflow operators use task_id, not name
        image=image,
        cmd=f'python script1.py {arg1}',
        do_xcom_push=True,  # Enable this if you used Method 2 (stdout capture)
    )

Key Fixes from Your Original Code

  • No direct callback execution: on_failure_callback=callback_function (not callback_function(variable1)) ensures the callback runs only when the task fails, not during DAG parsing.
  • XCom bridges runtime data: XCom stores the variable in Airflow's metadata DB, making it accessible to the callback after the task runs.
  • Correct operator parameter: Changed name to task_id—this is the standard identifier for Airflow tasks.

内容的提问来源于stack exchange,提问作者Milos Borisavljevic

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.11 08:16:57