使用Airflow BashOperator执行Python脚本时如何传递变量
解决方案
方法1:通过命令行参数传递Airflow变量
在DAG的BashOperator中,利用Airflow模板语法将变量作为命令行参数传给Python脚本,再在脚本里通过sys.argv读取参数。
修改DAG中的BashOperator:
with DAG(dag_id='test', default_args=default_args, schedule_interval='0 1 * * *', catchup=False) as dag: task_1 = BashOperator( task_id='tsk1', # 用Airflow模板语法引用变量,作为命令行参数传递 bash_command=f'python {PATH}/update_files.py {{ var.value.usr }} {{ var.value.pws }}' )
修改update_files.py脚本:
import wsclient import sys def function_one(): ... return df def function_two(df): ... return df_transformed if __name__ == "__main__": # 从命令行参数中提取用户名和密码 usr = sys.argv[1] pws = sys.argv[2] ws = WSClient("https://") if ws.do_login(usr, pws): raw_dataframe = function_one() transformed_dataframe = function_two(raw_dataframe ) for record in transformed_dataframe.to_dict(orient='records'): ws.do_create('Clients', record)
方法2:通过环境变量传递Airflow变量
借助Airflow模板语法将变量注入环境变量,Python脚本通过os.environ读取对应环境变量。
修改DAG中的BashOperator:
with DAG(dag_id='test', default_args=default_args, schedule_interval='0 1 * * *', catchup=False) as dag: task_1 = BashOperator( task_id='tsk1', # 先设置环境变量,再执行Python脚本 bash_command=f'export USR="{{ var.value.usr }}" && export PWS="{{ var.value.pws }}" && python {PATH}/update_files.py' )
修改update_files.py脚本:
import wsclient import os def function_one(): ... return df def function_two(df): ... return df_transformed if __name__ == "__main__": # 从环境变量中获取用户名和密码 usr = os.environ.get('USR') pws = os.environ.get('PWS') ws = WSClient("https://") if ws.do_login(usr, pws): raw_dataframe = function_one() transformed_dataframe = function_two(raw_dataframe ) for record in transformed_dataframe.to_dict(orient='records'): ws.do_create('Clients', record)
注意:如果变量包含特殊字符或需要更高安全性,可使用
{{ var.json.xxx }}语法,或配置Airflow的Fernet密钥启用变量加密,避免敏感信息泄露。
内容的提问来源于stack exchange,提问作者eponkratova
相关产品推荐
相关产品推荐

