如何在Airflow中通过编程方式设置connections与variables?
Airflow 编程设置 Connections 和 Variables 的实现方法
以下方法仅适用于本地调试场景,生产环境请勿使用,避免敏感信息泄露。
设置 Variables
直接调用 Airflow 内置的 Variable 类接口即可:
from airflow.models import Variable # 新增/修改变量 Variable.set(key="debug_var", value="test_value") # 存储字典/列表等可序列化对象时加 serialize_json 参数 Variable.set(key="debug_config", value={"env": "test", "retry": 3}, serialize_json=True) # 读取变量测试 print(Variable.get("debug_var"))
设置 Connections
需要结合 Connection 类和 Airflow 数据库会话实现,避免手动管理会话可以用官方提供的 provide_session 装饰器:
from airflow.models import Connection from airflow.utils.session import provide_session @provide_session def add_conn(conn_params, session=None): # 先检查是否存在同ID的连接,避免重复添加报错 exist_conn = session.query(Connection).filter(Connection.conn_id == conn_params["conn_id"]).first() if exist_conn: # 已存在则更新参数 for k, v in conn_params.items(): setattr(exist_conn, k, v) else: # 不存在则新增 new_conn = Connection(**conn_params) session.add(new_conn) session.commit() # 调用示例 conn_config = { "conn_id": "my_test_mysql", "conn_type": "mysql", "host": "127.0.0.1", "login": "root", "password": "test123", "port": 3306, "schema": "test_database" } add_conn(conn_config)
注意事项
- 执行代码的用户需要有 Airflow 元数据库的读写权限
- Airflow 1.x 和 2.x 的 Connection 可选字段略有差异,可根据自身版本调整入参
内容的提问来源于stack exchange,提问作者safex
相关产品推荐
相关产品推荐

