Apache Airflow 连接(Connection)管理完全指南:环境变量、数据库、Secrets Backend 与连接 URI 格式详解
【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow
本文系统讲解 Apache Airflow 中 Connection 对象的三种定义方式(环境变量、外部 Secrets Backend、元数据库)、连接 URI 与 JSON 两种序列化格式、连接测试(含 Airflow 3.3 新增的 worker 异步测试机制)以及通过provider.yaml声明式定义自定义连接字段的完整方法。读完后,你可以独立完成生产环境的连接配置、密钥安全落库、连接可用性验证,并为自定义 provider 提供带表单校验的连接类型。
一、Connection 的三种存储方式
Airflow 的Connection对象(airflow.sdk.Connection)用于保存访问外部服务所需的凭据和其他连接信息。官方文档 connection.rst 明确了三种定义方式:
- 环境变量:
AIRFLOW_CONN_{CONN_ID},运行时动态解析; - 外部 Secrets Backend:如 HashiCorp Vault、AWS SSM Parameter Store;
- Airflow 元数据库:通过 CLI(
airflow connections add)或 Web UI 管理。
三种方式的优先级与可见性各不相同,选择哪种方式取决于你的部署形态与安全边界。
二、用环境变量存储连接
2.1 命名规范
连接可以通过环境变量定义,命名规则为AIRFLOW_CONN_{CONN_ID},全部大写(注意CONN前后各有一个单下划线)。例如连接 ID 为my_prod_db时,环境变量名应为AIRFLOW_CONN_MY_PROD_DB。
从源码 environment_variables.py 可以确认前缀常量CONN_ENV_PREFIX = "AIRFLOW_CONN_",并且EnvironmentVariablesBackend.get_conn_value()会执行os.environ.get(CONN_ENV_PREFIX + conn_id.upper()),即 conn_id 会被强制转大写后拼接查询。此外源码还支持多团队(multi-team)模式下的命名空间变量,格式为AIRFLOW_CONN__<TEAM_ID>___<CONN_ID>(以三段下划线___分隔团队名),仅当core.multi_team开启时生效。
2.2 值的两种格式:JSON 与 URI
变量值可以是JSON(2.3.0 起支持)或Airflow URI 格式。
JSON 格式示例:
export AIRFLOW_CONN_MY_PROD_DATABASE='{ "conn_type": "my-conn-type", "login": "my-login", "password": "my-password", "host": "my-host", "port": 1234, "schema": "my-schema", "extra": { "param1": "val1", "param2": "val2" } }'URI 格式示例:
export AIRFLOW_CONN_MY_PROD_DATABASE='my-conn-type://login:password@host:port/schema?param1=val1¶m2=val2'2.3 用as_json便捷生成 JSON 表示(2.8.0+)
Connection类提供了便捷属性as_json,可自动输出标准环境变量内容:
>>> from airflow.sdk import Connection >>> c = Connection( ... conn_id="some_conn", ... conn_type="mysql", ... description="connection description", ... host="myhost.com", ... login="myname", ... password="mypassword", ... extra={"this_param": "some val", "that_param": "other val*"}, ... ) >>> print(f"AIRFLOW_CONN_{c.conn_id.upper()}='{c.as_json()}'") AIRFLOW_CONN_SOME_CONN='{"conn_type": "mysql", "description": "connection description", "host": "myhost.com", "login": "myname", "password": "mypassword", "extra": {"this_param": "some val", "that_param": "other val*"}}'同样方法还可以把已有 URI 格式的连接转换为 JSON:
>>> from airflow.sdk import Connection >>> c = Connection( ... conn_id="awesome_conn", ... description="Example Connection", ... uri="aws://YOUR_AWS_ACCESS_KEY_ID:YOUR_AWS_SECRET_ACCESS_KEY@/?__extra__=%7B%22region_name%22%3A+%22eu-central-1%22%2C+%22config_kwargs%22%3A+%7B%22retries%22%3A+%7B%22mode%22%3A+%22standard%22%2C+%22max_attempts%22%3A+10%7D%7D%7D", ... ) >>> print(f"AIRFLOW_CONN_{c.conn_id.upper()}='{c.as_json()}'") AIRFLOW_CONN_AWESOME_CONN='{"conn_type": "aws", "description": "Example Connection", "host": "", "login": "YOUR_AWS_ACCESS_KEY_ID", "password": "YOUR_AWS_SECRET_ACCESS_KEY", "schema": "", "extra": {"region_name": "eu-central-1", "config_kwargs": {"retries": {"mode": "standard", "max_attempts": 10}}}}'在源码 connection.py 中as_json属性与get_uri()(L282)相互对应,两者构成连接的双向序列化能力。
2.4 环境变量连接在 UI 和 CLI 中不可见
通过环境变量定义的连接不会显示在 Airflow UI 中,也不会被airflow connections list列出。原因是这些连接在运行时动态解析,通常发生在执行任务的worker进程内,而不是存储在元数据库中被 webserver 或 scheduler 加载。
这支持一种安全的部署模式:基于环境的密钥(如.env文件、Docker 或 Kubernetes secrets)只注入到 worker 等运行时组件,而不注入到 webserver 这类对用户暴露的组件。如果你需要连接在 UI 中可见或可编辑,应改用元数据库方式定义。
三、用外部 Secrets Backend 存储连接
你可以把连接存储在 HashiCorp Vault、AWS SSM Parameter Store 等外部密钥服务中,具体配置方式参见仓库文档 secrets-backend/index。使用外部后端时需注意后文提到的限制:连接测试功能对这类连接在 UI/REST API 场景下不可用。
四、用元数据库存储与管理连接
存入元数据库后,可用 Web UI 或 Airflow CLI 两种方式管理。
4.1 UI 创建连接
打开 UI 的Admin -> Connections页面,点击Add Connection:
- 在
Connection Id字段填写连接 ID(建议使用小写字母、下划线分词); - 通过
Connection Type字段选择连接类型; - 填写其余字段,各类型字段含义见官方连接类型参考;
- 点击
Save完成创建。
4.2 UI 编辑连接
在Admin -> Connections列表中找到目标连接,点击其旁的铅笔图标,修改属性后点击Save保存。
4.3 CLI 创建连接
可以从 CLI 向数据库添加连接,共三种写法(实现位于 connection_command.py):
方式一:JSON 格式(2.3.0 起支持)
airflow connections add 'my_prod_db' \ --conn-json '{ "conn_type": "my-conn-type", "login": "my-login", "password": "my-password", "host": "my-host", "port": 1234, "schema": "my-schema", "extra": { "param1": "val1", "param2": "val2" } }'方式二:Connection URI 格式
airflow connections add 'my_prod_db' \ --conn-uri '<conn-type>://<login>:<password>@<host>:<port>/<schema>?param1=val1¶m2=val2&...'方式三:逐参数指定
airflow commands add 'my_prod_db' \ --conn-type 'my-conn-type' \ --conn-login 'login' \ --conn-password 'password' \ --conn-host 'host' \ --conn-port 'port' \ --conn-schema 'schema' \ ...注:实际命令为
airflow connections add,逐参数写法用于不便构造完整 JSON/URI 的场景。
4.4 导出连接到文件
支持将数据库中的连接导出到文件(例如在两个环境间迁移连接),用法见 CLI 参考中的 Exporting Connections 小节。
4.5 数据库中的连接安全:Fernet 加密
对存储在元数据库中的连接,Airflow 使用Fernet对称加密算法加密 password 等敏感数据。这保证了在没有加密密钥的情况下,连接密码无法被读取或篡改。Fernet 密钥的配置参见安全文档中的 Fernet 章节(airflow-core/docs/security)。源码 connection.py 中的rotate_fernet_key()方法还支持用新密钥重新加密 password 与 extra,便于密钥轮换。
五、连接测试(Test Connection)
5.1 开关:test_connection配置项
出于安全考虑,连接测试功能在 UI、API 和 CLI 中默认禁用。是否启用由 Airflow 配置(airflow.cfg)core 段的test_connection标志控制,也等价于环境变量AIRFLOW__CORE__TEST_CONNECTION。
在 config.yml 中可确认其定义(2.7.0 引入,默认Disabled),接受三个取值:
| 取值 | 行为 |
|---|---|
Disabled | 禁用测试功能,UI 中的 Test Connection 按钮不可用(默认值) |
Enabled | 启用测试功能,UI 显示并激活 Test Connection 按钮 |
Hidden | 禁用测试功能,且隐藏 UI 中的 Test Connection 按钮 |
官方强烈建议:在确认只有高度受信任的 UI/API 用户拥有 "edit connection" 权限之前,不要开启此功能——连接测试可能被恶意利用。
5.2 测试的执行机制
开启后,连接测试可在以下入口触发:UI 的连接创建/编辑页面、Connections REST API、或airflow connections testCLI 命令。
执行时,Airflow 调用对应 Hook 类的test_connection方法并回报结果。源码 connection.py 中Connection.test_connection()先经get_hook()按conn_type从 provider 发现 Hook,再调用 Hook 的test_connection。两种失败情形:
- 该连接类型没有关联的 Hook,或 Hook 未实现
test_connection方法——会显示错误信息(UI 场景下则功能被禁用); - UI 测试从webserver发起,因此受 webserver 的网络出口规则约束;若 webserver 与 worker 安装的库/provider 不同,测试结果可能不同。
另外注意:对存放在外部 Secrets Backend中的连接,UI/REST API 场景下的测试功能不可用。
5.3 异步(worker 派发)连接测试
上述测试运行在 API server 上。Airflow 还支持把测试派发到worker上执行,使连接凭据只在 worker(任务实际使用凭据的地方)上被使用,而不是 API server 上——这在"连接只能从 worker 网络可达"或"不希望凭据在 API server 上被执行"时非常有用。
工作机制(使用与同步测试相同的[core] test_connection开关):
- 通过 Connections REST API 提交测试:
POST /connections/enqueue-test(返回一个 token); - 用
GET /connections/enqueue-test轮询结果,token 通过请求头Airflow-Connection-Test-Token传递; - 结果只能用该 token 读取,且仅限对该连接有权限的用户;多团队部署中,其他团队拥有的连接的测试对你不可见;
- 执行测试的 worker 由 scheduler 单独授权:为其签发一个短期 JWT,subject 是该 connection-test 请求 ID,scope 为
workload;worker 调用的 Execution API 端点强制ct:self校验(token subject 必须与请求路径中的 connection-test id 一致),确保该 worker 只能取用并回报那一个指定请求的连接测试,无法触达其他测试、任务实例或连接。
相关调优参数位于配置段[connection_test],源码 config.yml 中确认(均为 3.3.0 引入):
| 配置项 | 默认值 | 说明 |
|---|---|---|
connection_test.timeout | 60(秒) | 单个 worker 派发测试允许运行的最长时间,超时报超时;scheduler 的 reaper 用它加宽限期标记过期测试为失败 |
connection_test.max_concurrency | 4 | 同时可处于活跃(QUEUED + RUNNING)状态的测试数量上限,超出者保持 PENDING;该上限按 scheduler 实例生效而非全局,N 个 HA scheduler 时最坏情况每 tick 派发量为N * max_concurrency |
connection_test.reaper_interval | 30.0(秒) | scheduler 检查并标记超期(QUEUED/RUNNING 超过 timeout+宽限期)测试的频率 |
因为测试在 worker 上运行,其结果反映的是该 worker的库、provider 与网络访问情况,可能与 API server 不同。
六、自定义连接类型与自定义连接字段
6.1 通过 provider 声明自定义连接类型
Airflow 允许定义自定义连接类型(包括修改连接的添加/编辑表单)。自定义连接类型通常由社区 provider 定义,你也可以添加自己的自定义 provider。通过在provider.yaml的connection-types数组中暴露连接类型,可以实现:
- 添加自定义连接类型;
- 由连接类型自动创建 Hook;
- 添加自定义表单字段,用于显示和编辑连接 URL 中的自定义 "extra" 参数;
- 隐藏你的连接用不到的标准字段;
- 添加占位符(placeholder),展示字段的格式示例。
6.2 推荐方式:在provider.yaml中声明式定义(Airflow 3.2+)
从 Airflow 3.2 起,定义自定义连接字段与字段行为的推荐方式是在provider.yaml中声明式定义,无需在运行时导入flask_appbuilder或wtforms。旧的 Python Hook 方法get_connection_form_widgets()和get_ui_field_behaviour()仍作为 fallback 可用,且只会"在发出弃用公告并经过迁移窗口后"才会移除——旧 provider 继续工作,但新 provider 应使用 YAML 方式。
conn-fields——存储在Connection.extra中的自定义字段:
connection-types: - hook-class-name: airflow.providers.myservice.hooks.myservice.MyServiceHook connection-type: myservice conn-fields: workspace: label: Workspace schema: type: - string - 'null' project: label: Project ID schema: type: - string - 'null'ui-field-behaviour——对标准连接字段的定制(隐藏、重命名、占位符):
connection-types: - hook-class-name: airflow.providers.myservice.hooks.myservice.MyServiceHook connection-type: myservice ui-field-behaviour: hidden-fields: - port - host - login - schema relabeling: password: API Token placeholders: password: your-api-token workspace: My workspace gid project: My project gid字段 schema 类型遵循 JSON Schema 约定。
6.3 传统方式:Python Hook 方法(legacy)
实现 Hook 类上的get_connection_form_widgets()可添加自定义表单字段——键是字段存入extra字典时使用的字符串名,值必须是wtforms.fields.core.Field的子类:
@staticmethod def get_connection_form_widgets() -> dict[str, Any]: """Returns connection widgets to add to connection form""" from flask_appbuilder.fieldwidgets import BS3TextFieldWidget from flask_babel import lazy_gettext from wtforms import StringField return { "workspace": StringField(lazy_gettext("Workspace"), widget=BS3TextFieldWidget()), "project": StringField(lazy_gettext("Project"), widget=BS3TextFieldWidget()), }注意:Airflow 2.3 起自定义字段不再需要extra__<conn type>__前缀(2.3 之前必须加前缀,且值以该前缀存入extra字典)。
get_ui_field_behaviour()方法用于定制标准字段行为(隐藏、重命名、占位符):
@staticmethod def get_ui_field_behaviour() -> dict[str, Any]: """Returns custom field behaviour""" return { "hidden_fields": ["port", "host", "login", "schema"], "relabeling": {}, "placeholders": { "password": "Asana personal access token", "workspace": "My workspace gid", "project": "My project gid", }, }两个易踩的坑:
- 若某个
extra字段名与标准连接属性(login、password、host、schema、port、extra)冲突,其表单占位符必须以extra__<conn type>__前缀书写,例如extra__myservice__password; - 2.2.0 之前 provider 通过
hook-class-names数组暴露连接,现已被connection-types数组取代(后者在 worker 中按需加载更高效)。若 provider 需兼容低于 2.2.0 的 Airflow,两个数组应同时存在,CI 构建时会自动校验两者一致性。
七、Connection URI 格式详解
7.1 URI 总体格式
出于历史原因,Airflow 有一套特殊的 URI 格式用于把Connection对象序列化为字符串(2.3.0 起也可改用 JSON 序列化)。通用格式如下:
my-conn-type://my-login:my-password@my-host:5432/my-schema?param1=val1¶m2=val2上述 URI 产生的Connection对象等价于:
Connection( conn_id="", conn_type="my_conn_type", description=None, login="my-login", password="my-password", host="my-host", port=5432, schema="my-schema", extra=json.dumps(dict(param1="val1", param2="val2")), )7.2 用get_uri()生成 URI
Connection类的便捷方法get_uri()可自动生成 URI:
>>> import json >>> from airflow.sdk import Connection >>> c = Connection( ... conn_id="some_conn", ... conn_type="mysql", ... description="connection description", ... host="myhost.com", ... login="myname", ... password="mypassword", ... extra=json.dumps(dict(this_param="some val", that_param="other val*")), ... ) >>> print(f"AIRFLOW_CONN_{c.conn_id.upper()}='{c.get_uri()}'") AIRFLOW_CONN_SOME_CONN='mysql://myname:mypassword@myhost.com?this_param=some+val&that_param=other+val%2A'两点注意:
get_uri()返回的是Airflow 格式URI,不是SQLAlchemy 兼容 URI;数据库连接若需要 SQLAlchemy URI,应使用DbApiHook.sqlalchemy_url属性;- 连接创建后也可用
airflow connections get查看其字段与 URI:
$ airflow connections get sqlite_default Id: 40 Connection Id: sqlite_default Connection Type: sqlite Host: /tmp/sqlite_default.db Schema: null Login: null Password: null Port: null Is Encrypted: false Is Extra Encrypted: false Extra: {} URI: sqlite://%2Ftmp%2Fsqlite_default.db从源码 connection.py 看,get_uri()内部会先处理 conn_type 下划线转连字符(scheme 中含下划线会按 RFC3986 告警)、对 login/password/host/schema 做quote转义,再决定 extra 是扁平化为 query 参数还是整体放入__extra__。
7.3extra中任意 dict 的处理:__extra__查询参数
某些 JSON 结构无法无损失地 url-encode。对于这类 JSON,get_uri会把整个字符串存放在查询参数__extra__下:
>>> extra_dict = {"my_val": ["list", "of", "values"], "extra": {"nested": {"json": "val"}}} >>> c = Connection( ... conn_type="scheme", ... host="host/location", ... schema="schema", ... login="user", ... password="password", ... port=1234, ... extra=json.dumps(extra_dict), ... ) >>> uri = c.get_uri() >>> uri 'scheme://user:password@host%2Flocation:1234/schema?__extra__=%7B%22my_val%22%3A+%5B%22list%22%2C+%22of%22%2C+%22values%22%5D%2C+%22extra%22%3A+%7B%22nested%22%3A+%7B%22json%22%3A+%22val%22%7D%7D%7D'且往返解析可验证等价:
>>> new_c = Connection(uri=uri) >>> new_c.extra_dejson == extra_dict True而最常见的情形——extra 只存 key-value 对——则使用普通的 url 编码。URI 解析验证示例:
>>> from airflow.sdk import Connection >>> c = Connection(uri="my-conn-type://my-login:my-password@my-host:5432/my-schema?param1=val1¶m2=val2") >>> print(c.login) my-login >>> print(c.password) my-password7.4 特殊字符处理
手动构造 URI 时,某些字符需要特殊处理。例如密码中含/会导致解析失败:
>>> c = Connection(uri="my-conn-type://my-login:my-pa/ssword@my-host:5432/my-schema?param1=val1¶m2=val2") ValueError: invalid literal for int() with base 10: 'my-pa'解决办法是用urllib.parse.quote_plus编码该字符:
>>> c = Connection(uri="my-conn-type://my-login:my-pa%2Fssword@my-host:5432/my-schema?param1=val1¶m2=val2") >>> print(c.password) my-pa/ssword因此生成 URI 时强烈建议直接使用Connection.get_uri()便捷方法(它内部已做完整的百分号转义,如*会被编码为%2A、空格为+),而不是手工拼接字符串。
八、实践建议小结
- 开发/本地环境:优先用环境变量(JSON 或 URI),快速注入、不进数据库,但注意其在 UI/CLI 不可见;
- 生产环境:优先外部 Secrets Backend(Vault/SSM 等),其次是元数据库 + Fernet 加密;
- 需要 UI 可见、REST API 集成:存入元数据库,通过 UI 或
airflow connections add管理; - 凭据敏感或网络隔离:使用
[core] test_connection = Enabled开启同步测试,或采用 3.3 引入的 worker 派发异步测试,配合connection_test.*三个参数控制超时与并发; - 自定义 provider:新代码一律用
provider.yaml的connection-types+conn-fields+ui-field-behaviour声明式方案,Python Hook 方法仅作遗留兼容。
【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考