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

Airflow报错NameError: name 'ti'未定义,求解决方案

问题原因
  • 直接在SimpleHttpOperator的data参数中调用ti.xcom_pull会触发报错,因为DAG解析阶段(静态)还不存在任务实例ti,只有任务运行时(动态)才能访问该对象。
  • 全局变量prog_no无法跨任务共享,Airflow的任务运行在独立进程/worker中,全局变量不会在任务间传递,且解析阶段也获取不到运行时的dag_run.conf。
解决方案

方案1:利用Airflow模板化功能

Airflow支持Jinja2模板,SimpleHttpOperator的data参数默认支持模板渲染,直接用Jinja语法引用运行时变量即可:

sparkway_request = SimpleHttpOperator(
    task_id='sparkway_request',
    endpoint=Variable.get('sparkway_api_endpoint'),
    method="POST",
    http_conn_id="SPARKWAY_CONN",
    # 用Jinja模板引用执行日期和XCom值
    data="""
    {
        "cmd": "sparkway-submit --master kubernetes --job-name spark-allowance-calculation --class com.xxx.CalculationApplication --spark-app s3a://xyz.jar --arguments SuspendedProgramStatus --num-executors 2 --executor-cores 2 --executor-memory 3G --driver-memory 3G",
        "arguments": "{{ execution_date }},{{ ti.xcom_pull(task_ids='parameterized_task') }}",
        "type": "job"
    }
    """,
    headers={
        "Authorization": f"Bearer {Variable.get('sparkway_token')}",
        "Content-Type": "application/json",
        "X-CSRF-TOKEN": Variable.get('sparkway_csrf_token')
    },
    response_check=lambda response: handle_sparkway_response(response),
    log_response=True,
    dag=dag
)

说明:

  • {{ execution_date }}会自动获取当前任务的运行日期,替代静态的dag.latest_execution_date。
  • {{ ti.xcom_pull(task_ids='parameterized_task') }}会在任务运行时拉取指定任务的XCom返回值。
  • 直接使用Jinja格式的JSON字符串作为data,Airflow会自动完成模板渲染并解析为合法JSON。

方案2:改用PythonOperator结合HttpHook(更灵活)

如果需要处理复杂逻辑(比如参数校验),可以用PythonOperator构造请求,直接通过上下文获取运行时变量:

from airflow.providers.http.hooks.http import HttpHook
import json

def send_spark_request(**context):
    ti = context['ti']
    execution_date = context['execution_date']
    # 拉取XCom中的program_no值
    prog_no = ti.xcom_pull(task_ids='parameterized_task')
    
    # 构造请求 payload
    payload = {
        "cmd": "sparkway-submit --master kubernetes --job-name spark-allowance-calculation --class com.xxx.CalculationApplication --spark-app s3a://xyz.jar --arguments SuspendedProgramStatus --num-executors 2 --executor-cores 2 --executor-memory 3G --driver-memory 3G",
        "arguments": f"{execution_date},{prog_no}",
        "type": "job"
    }
    
    # 初始化HttpHook发送请求
    http_hook = HttpHook(http_conn_id='SPARKWAY_CONN', method='POST')
    headers = {
        "Authorization": f"Bearer {Variable.get('sparkway_token')}",
        "Content-Type": "application/json",
        "X-CSRF-TOKEN": Variable.get('sparkway_csrf_token')
    }
    
    response = http_hook.run(
        endpoint=Variable.get('sparkway_api_endpoint'),
        data=json.dumps(payload),
        headers=headers
    )
    
    # 执行响应检查
    handle_sparkway_response(response)

sparkway_request = PythonOperator(
    task_id='sparkway_request',
    python_callable=send_spark_request,
    provide_context=True,
    dag=dag
)
额外优化建议
  • 删除get_conf函数中的全局变量prog_no,直接返回值即可,XCom会自动存储该返回值供后续任务调用。
  • 为避免dag_run.conf缺失program_no的情况,在get_conf中添加异常处理:
def get_conf(**context):
    dag_run = context['dag_run']
    if not dag_run or 'program_no' not in dag_run.conf:
        raise ValueError("dag_run.conf中缺少必填参数program_no")
    return dag_run.conf['program_no']

内容的提问来源于stack exchange,提问作者Ayush Goyal

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.29 12:42:37