【发布时间】:2019-01-22 15:16:41
【问题描述】:
通过Admin -> Connections,我们能够创建/修改连接的参数,但我想知道是否可以通过 API 执行相同操作,以便以编程方式设置连接
airflow.models.Connection 似乎只处理实际连接到实例而不是将其保存到列表中。这似乎是一个应该实现的功能,但我不确定在哪里可以找到这个特定功能的文档。
【问题讨论】:
通过Admin -> Connections,我们能够创建/修改连接的参数,但我想知道是否可以通过 API 执行相同操作,以便以编程方式设置连接
airflow.models.Connection 似乎只处理实际连接到实例而不是将其保存到列表中。这似乎是一个应该实现的功能,但我不确定在哪里可以找到这个特定功能的文档。
【问题讨论】:
要使用session = settings.Session(),它假定气流数据库后端已启动。对于尚未为您的开发环境设置它的人,使用 Connection 类和环境变量的混合方法将是一种解决方法。
以下是设置 S3Hook 的示例
from airflow.providers.amazon.aws.hooks.s3 import S3Hook
from airflow.models.connection import Connection
import os
import json
aws_default = Connection(
conn_id="aws_default",
conn_type="aws",
login='YOUR-AWS-KEY-ID',
password='YOUR-AWS-KEY-SECRET',
extra=json.dumps({'region_name': 'us-east-1'})
)
os.environ["AIRFLOW_CONN_AWS_DEFAULT"] = aws_default.get_uri()
s3_hook = S3Hook(aws_conn_id='aws_default')
s3_hook.list_keys(bucket_name='YOUR-BUCKET', prefix='YOUR-FILENAME')
【讨论】:
您可以使用connection URI 格式populate connections using environment variables。
环境变量命名约定是AIRFLOW_CONN_
,全部大写。 因此,如果您的连接 ID 是 my_prod_db,那么变量名称应该是 AIRFLOW_CONN_MY_PROD_DB。
一般来说,Airflow 的 URI 格式是这样的:
my-conn-type://my-login:my-password@my-host:5432/my-schema?param1=val1¶m2=val2
请注意,以这种方式注册的连接不会显示在 Airflow UI 中。
【讨论】:
首先检查连接是否存在,然后使用from airflow.models import Connection 创建新连接:
def create_conn(conn_id, conn_type, host, login, password, port):
conn = Connection(
conn_id=conn_id,
conn_type=conn_type,
host=host,
login=login,
password=password,
port=port
)
session = settings.Session()
conn_name = session\
.query(Connection)\
.filter(Connection.conn_id == conn.conn_id)\
.first()
if str(conn_name) == str(conn_id):
return logging.info(f"Connection {conn_id} already exists")
session.add(conn)
session.commit()
logging.info(Connection.log_info(conn))
logging.info(f'Connection {conn_id} is created')
【讨论】:
如果您需要在 Python/Airflow 代码之外、通过 bash、在 Dockerfile 等中添加、删除和列出连接,您还可以从 Airflow CLI 添加、删除和列出连接。
airflow connections --add ...
用法:
airflow connections [-h] [-l] [-a] [-d] [--conn_id CONN_ID]
[--conn_uri CONN_URI] [--conn_extra CONN_EXTRA]
[--conn_type CONN_TYPE] [--conn_host CONN_HOST]
[--conn_login CONN_LOGIN] [--conn_password CONN_PASSWORD]
[--conn_schema CONN_SCHEMA] [--conn_port CONN_PORT]
https://airflow.apache.org/cli.html#connections
CLI 目前似乎不支持修改现有连接,但它存在一个 Jira 问题,在 GitHub 上有一个活动的打开 PR。
【讨论】:
连接实际上是一个模型,您可以使用它来查询和插入新连接
from airflow import settings
from airflow.models import Connection
conn = Connection(
conn_id=conn_id,
conn_type=conn_type,
host=host,
login=login,
password=password,
port=port
) #create a connection object
session = settings.Session() # get the session
session.add(conn)
session.commit() # it will insert the connection object programmatically.
【讨论】:
Connection 也可以删除