管理 Airflow 连接

Airflow 的 Connection 对象保存连接外部服务所需的凭据和其他信息。关于钩子与连接的总体介绍,参见 Connections & Hooks。

连接可以定义在环境变量、外部 Secrets Backend,或 Airflow 元数据库中;数据库中的连接可通过 CLI 或 Web UI 管理。

在环境变量中保存连接

环境变量命名为全大写的 AIRFLOW_CONN_{CONN_ID},注意 CONN 两侧各只有一个下划线。例如连接 ID 为 my_prod_db,变量名就是 AIRFLOW_CONN_MY_PROD_DB。值可以是 JSON,也可以是 Airflow URI 格式。

JSON 示例

JSON 支持在 2.3.0 中加入:

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"
    }
}'

生成连接的 JSON 表示

2.8.0 加入了 Connection.as_json(),方便生成连接 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%22r etries%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}}}}'

URI 示例

以 Airflow URI 序列化时:

export AIRFLOW_CONN_MY_PROD_DATABASE='my-conn-type://login:password@host:port/schema?param1=val1&param2=val2'

生成有效 URI 的方法见下文“URI 格式”。

在 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。

在数据库中保存连接

数据库之外,也可使用环境变量或外部秘密后端。保存到数据库时,可通过 Web UI 或 Airflow CLI 管理。

在 UI 中创建连接

打开 Admin->Connections,点击 Add Connection:

  1. 填写 Connection Id,建议使用小写字母,并用下划线分隔词语。
  2. 在 Connection Type 中选择类型。
  3. 填写剩余字段,不同类型字段的处理参见下文 extra。
  4. 点击 Save。

在 UI 中编辑连接

在 Admin->Connections 的连接列表中,点击目标连接旁的铅笔图标,修改属性后点击 Save。

使用 CLI 创建连接

从 2.3.0 开始,可以用 JSON 添加数据库连接:

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"
        }
    }'

也可以使用 URI:

airflow connections add 'my_prod_db' \
    --conn-uri '<conn-type>://<login>:<password>@<host>:<port>/<schema>?param1=val1&param2=val2&...'

或分别指定参数:

airflow connections 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 使用 Fernet 加密元数据库里的密码和其他可能敏感的数据。没有加密密钥,就不能读取或修改连接密码。配置方法见 Fernet 文档。

测试连接

出于安全原因,UI、API 与 CLI 默认禁用连接测试。强烈建议只有在确认具备“编辑连接”权限的 UI/API 用户均高度可信后,才启用此功能,参见已认证 UI 用户能力。

通过 airflow.cfg 的 [core] test_connection 或环境变量 AIRFLOW__CORE__TEST_CONNECTION 控制:

  • Disabled:禁用测试及 UI 按钮,也是默认值。
  • Enabled:启用测试和 UI 按钮。
  • Hidden:禁用测试并隐藏按钮。

启用后,可在 UI 新建/编辑连接页面使用,通过 Connections REST API,或运行 airflow connections test。

通过 UI 或 REST API 时,外部秘密后端中的连接不支持此功能。

测试会调用关联 Hook 类的 test_connection。如果连接类型没有 Hook,或 Hook 没有该方法实现,会显示错误;在 UI 中也可能直接禁用功能。

UI 测试从 Webserver 执行,受其网络出口规则约束。Webserver 与 Worker,或不同 CLI 所在机器/Pod,若安装的库或 Provider 不同,测试结果也可能不同。

异步测试:派发到 Worker

上述同步测试运行在 API Server。也可以把测试派发到 Worker,让凭据只在与任务执行相同的 Worker 上使用。连接只有 Worker 能访问、或希望避免在 API Server 使用凭据时,这种方式很有用。

它使用相同的 [core] test_connection 开关。通过 REST API 向 POST /connections/enqueue-test 提交请求,取得令牌后,以 GET /connections/enqueue-test 轮询,在 Airflow-Connection-Test-Token 请求头中传入令牌。

只有持有该令牌且获授权访问该连接的用户才能读取结果;多团队部署中,其他团队拥有的连接测试不会对你可见。

Worker 的授权与轮询令牌分别处理。Scheduler 为单次请求签发短期 JWT,主体为连接测试请求 ID,作用域为 workload。Worker 调用的 Execution API 强制执行 ct:self 检查:令牌主体必须与路径中的测试 ID 相同。因此它只能获取并报告这一次测试,不能访问其他连接测试、任务实例或连接。

在 [connection_test] 中调整行为:

  • timeout:测试超时前允许的时长,默认 60 秒。
  • max_concurrency:并发测试数量,默认 4;其他请求等待空位。
  • reaper_interval:检查并终止超时测试的间隔,默认 30.0 秒。

Worker 测试反映该 Worker 的库、Provider 与网络环境,结果可能不同于 API Server。

自定义连接类型

Airflow 允许自定义连接类型,包括调整新建/编辑表单。社区 Provider 可定义它们,也可以创建自己的 Provider,参见 Providers。

在 provider.yaml 的 connection-types 数组中公开连接类型,可以添加类型、按类型自动创建 Hook、为 URL 的额外参数添加表单字段、隐藏无用标准字段,以及添加格式示例占位文本。

自定义连接字段

从 Airflow 3.2 开始,推荐在 provider.yaml 中声明式定义 UI 元数据。 这种方式不需运行时导入 flask_appbuilder 或 wtforms。旧的 Python Hook 方法 get_connection_form_widgets() 和 get_ui_field_behaviour() 仍可作为回退;只有经过弃用通知和迁移窗口后才会移除,旧 Provider 仍能工作。新 Provider 应使用下面的 YAML 方式。

在 provider.yaml 中定义 UI 元数据

在 connection-types 下配置两部分。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。完整的 conn-fields schema 选项参考 Use Params to Provide a Trigger UI Form。

在 Python 中定义 UI 元数据(旧方式)

该方式仍有效,不会在没有弃用通知的情况下移除;新 Provider 应使用 YAML。

在 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>__ 前缀,保存到 extra 时也保留前缀。从 2.3 开始不再需要。

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、scheme、port、extra)重名,指定表单占位文本时仍须加 extra__<conn type>__,例如 extra__myservice__password。可参考 Provider 中的 JdbcHook。

已弃用的 hook-class-names: 2.2.0 之前,Provider 通过元数据的 hook-class-names 数组公开连接;该方式在 Worker 使用个别 Hook 时效率不足,已被 connection-types 替代。Provider 若仍支持低于 2.2.0 的 Airflow,需要同时保留两者;CI 自动检查会验证两个数组的一致性。

URI 格式

从 2.3.0 起,也可以使用 JSON 序列化连接。由于历史原因,Airflow 有自己的特殊 URI 格式,用来将 Connection 序列化为字符串:

my-conn-type://my-login:my-password@my-host:5432/my-schema?param1=val1&param2=val2

上面的 URI 等价于:

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")),
)

生成 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'

它返回 Airflow 格式的 URI,并非 SQLAlchemy 兼容 URI。数据库连接需要 SQLAlchemy URI 时,使用 sqlalchemy_url 属性。

创建连接后,也可使用 airflow connections get:

$ 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

extra 中的任意字典

部分 JSON 结构不能无损 URL 编码。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

最常见的纯键值对仍使用普通 URL 编码。可这样检查 URI 解析:

>>> from airflow.sdk import Connection

>>> c = Connection(uri="my-conn-type://my-login:my-password@my-host:5432/my-schema?param1=val1&param2=val2")
>>> print(c.login)
my-login
>>> print(c.password)
my-password

连接参数中的特殊字符

生成连接应优先用 Connection.get_uri();本节用于解释手动构造 URI 的注意事项。

特殊字符需要特殊处理。例如,密码包含 / 时,下面会失败:

>>> c = Connection(uri="my-conn-type://my-login:my-pa/ssword@my-host:5432/my-schema?param1=val1&param2=val2")
ValueError: invalid literal for int() with base 10: 'my-pa'

原文用 quote_plus() 编码来解决,并展示:

>>> c = Connection(uri="my-conn-type://my-login:my-pa%2Fssword@my-host:5432/my-schema?param1=val1&param2=val2")
>>> print(c.password)
my-pa/ssword

原文:Apache Airflow 3.3.2:Managing Connections。Copyright © The Apache Software Foundation. 本版本翻译正文并调整版式,保留原文 API 和输出,代码中的网页提取空格已按 3.3.2 原始文档恢复。采用 Apache License 2.0,许可全文随稿包保留。Apache、Airflow 及其标志为 Apache Software Foundation 的商标或注册商标;其他产品名称和品牌归各自权利人所有。

© 版权声明
THE END
喜欢就支持一下吧
点赞0 分享
评论 抢沙发

请登录后发表评论

    暂无评论内容