【问题标题】:Is there a way to create/modify connections through Airflow API有没有办法通过 Airflow API 创建/修改连接
【发布时间】:2019-01-22 15:16:41
【问题描述】:

通过Admin -> Connections,我们能够创建/修改连接的参数,但我想知道是否可以通过 API 执行相同操作,以便以编程方式设置连接

airflow.models.Connection 似乎只处理实际连接到实例而不是将其保存到列表中。这似乎是一个应该实现的功能,但我不确定在哪里可以找到这个特定功能的文档。

【问题讨论】:

    标签: python airflow


    【解决方案1】:

    要使用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')
    

    【讨论】:

      【解决方案2】:

      您可以使用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&param2=val2
      

      请注意,以这种方式注册的连接不会显示在 Airflow UI 中。

      【讨论】:

        【解决方案3】:

        首先检查连接是否存在,然后使用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')
        

        【讨论】:

          【解决方案4】:

          如果您需要在 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。

          【讨论】:

            【解决方案5】:

            连接实际上是一个模型,您可以使用它来查询和插入新连接

            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.
            

            【讨论】:

            • 感谢您的快速回复。这正是我正在寻找的。一旦 stackoverflow 允许我,我会接受这个答案
            • 有没有先列出连接,然后检查连接是否已经存在?
            • @mad_ 我想可能是在使用上面的设置,我们可以访问连接列表并从中检查。是的,我们可以直接使用 bashoperator,然后使用气流连接 -l,但我不太愿意为它创建另一个任务
            • 如何从此会话中删除连接
            • 同样型号Connection 也可以删除
            猜你喜欢
            • 1970-01-01
            • 1970-01-01
            • 1970-01-01
            • 1970-01-01
            • 1970-01-01
            • 1970-01-01
            • 2015-06-12
            • 1970-01-01
            • 1970-01-01
            相关资源
            最近更新 更多