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
相关产品推荐
相关产品推荐

