如何在Airflow中设置操作系统环境变量供下游任务及dbt调用
核心原因说明
Airflow中每个任务(不管是BashOperator还是PythonOperator)都会启动独立的进程运行,进程之间的环境变量是完全隔离的:
- 你在第一个Bash任务里执行export设置的变量,只会在当前任务的bash进程生命周期内有效,第二个任务启动的是全新的bash进程,不会继承前一个任务的自定义环境变量
- PythonOperator修改
os.environ也是同理,仅对当前任务的Python进程生效,下一个任务的新进程不会保留这个值
适配dbt场景的解决方案(最推荐)
不需要单独做任务导出变量,直接在运行dbt命令的BashOperator的env参数中注入Airflow变量即可,dbt运行的进程可以直接读取到对应环境变量,完全适配你的yaml配置:
from airflow import DAG from airflow.operators.bash import BashOperator from datetime import datetime from airflow.models import Variable import os default_args = { 'owner': 'airflow', 'depends_on_past': False, 'start_date': datetime(2020,8,1), 'retries': 0 } with DAG('sample', default_args=default_args, schedule_interval='@once') as dag: dbt_run_task = BashOperator( task_id='dbt_run', # 替换为你实际的dbt命令,比如dbt seed、dbt run等 bash_command='dbt run', env={ **os.environ.copy(), # 保留系统默认环境变量 "dbt_user": Variable.get("dbt_user"), # 敏感变量也可以直接传入,Airflow会自动对日志做脱敏处理 "dbt_password": Variable.get("dbt_password") } )
其他场景解决方案
方案2:同一会话执行多命令
如果你的操作必须拆分在同一个bash会话里执行,就把所有命令合并到同一个任务的bash_command中,用&&连接,保证所有命令都在同一个bash进程里运行:
task = BashOperator( task_id='dbt_run_single_task', bash_command='export dbt_user={{ var.value.dbt_user }} && export dbt_password={{ var.value.dbt_password }} && dbt run' )
方案3:跨多任务共享变量
如果有多个下游任务都需要用到这些变量,可以用XCom传递变量值,每个下游任务自行读取后注入当前进程的环境变量:
from airflow import DAG from airflow.operators.python import PythonOperator from airflow.operators.bash import BashOperator from datetime import datetime from airflow.models import Variable import os default_args = { 'owner': 'airflow', 'depends_on_past': False, 'start_date': datetime(2020,8,1), 'retries': 0 } def push_dbt_vars(ti): # 把变量推送到XCom共享存储 ti.xcom_push(key="dbt_user", value=Variable.get("dbt_user")) ti.xcom_push(key="dbt_password", value=Variable.get("dbt_password")) def use_vars_python(ti): # 下游Python任务读取XCom,设置当前进程的环境变量 dbt_user = ti.xcom_pull(key="dbt_user", task_ids="push_dbt_vars") os.environ["dbt_user"] = dbt_user print(os.environ["dbt_user"]) with DAG('sample', default_args=default_args, schedule_interval='@once') as dag: push_vars_task = PythonOperator( task_id='push_dbt_vars', python_callable=push_dbt_vars ) python_use_task = PythonOperator( task_id='use_vars_python', python_callable=use_vars_python ) bash_use_task = BashOperator( task_id='use_vars_bash', bash_command='echo $dbt_user', env={ **os.environ.copy(), # Bash任务直接用模板语法读取XCom "dbt_user": "{{ ti.xcom_pull(key='dbt_user', task_ids='push_dbt_vars') }}" } ) push_vars_task >> [python_use_task, bash_use_task]
内容的提问来源于stack exchange,提问作者Adam
相关产品推荐
相关产品推荐

