Postgres 使用 SQLExecuteQueryOperator 的操作指南

简介

Apache Airflow 拥有大量算子,可用于实现工作流中的各种任务。Airflow 本质上是一个有向无环图(Directed Acyclic Graph),由任务(节点)和依赖关系(边)组成。

由算子定义或实现的任务是数据管道中的一个工作单元。

本指南的目的在于展示如何使用 SQLExecuteQueryOperator 与 PostgreSQL 数据库进行交互。

注意

以前使用 PostgresOperator 来完成此类操作。该算子已被弃用并移除,请改用 SQLExecuteQueryOperator。

使用 SQLExecuteQueryOperator 的常见数据库操作

要使用 SQLExecuteQueryOperator 执行 PostgreSQL 请求,需要两个参数:sqlconn_id。这两个参数最终会传递给直接与 Postgres 数据库交互的 DbApiHook 对象。

创建 Postgres 数据库表

以下代码片段基于 Airflow‑2.0

tests/system/postgres/example_postgres.py[source]



# create_pet_table, populate_pet_table, get_all_pets, and get_birth_date are examples of tasks created by
# instantiating the Postgres Operator

ENV_ID = os.environ.get("SYSTEM_TESTS_ENV_ID")
DAG_ID = "postgres_operator_dag"

with DAG(
    dag_id=DAG_ID,
    start_date=datetime.datetime(2020, 2, 2),
    schedule="@once",
    catchup=False,
) as dag:
    create_pet_table = SQLExecuteQueryOperator(
        task_id="create_pet_table",
        sql="""
            CREATE TABLE IF NOT EXISTS pet (
            pet_id SERIAL PRIMARY KEY,
            name VARCHAR NOT NULL,
            pet_type VARCHAR NOT NULL,
            birth_date DATE NOT NULL,
            OWNER VARCHAR NOT NULL);
          """,
    )

把所有 SQL 语句直接写进算子里既不美观,也会给后期维护带来麻烦。为了解决这个问题,Airflow 提供了一种优雅的方案:在 DAG 目录下创建一个名为 sql 的子目录,并将所有包含 SQL 查询的文件放入该目录。

你的 dags/sql/pet_schema.sql 应该是如下结构

-- create pet table
CREATE TABLE IF NOT EXISTS pet (
    pet_id SERIAL PRIMARY KEY,
    name VARCHAR NOT NULL,
    pet_type VARCHAR NOT NULL,
    birth_date DATE NOT NULL,
    OWNER VARCHAR NOT NULL);

现在让我们在 DAG 中重构 create_pet_table 任务

create_pet_table = SQLExecuteQueryOperator(
    task_id="create_pet_table",
    conn_id="postgres_default",
    sql="sql/pet_schema.sql",
)

向 Postgres 数据库表插入数据

假设我们已经在 dags/sql/pet_schema.sql 文件中写好了下面的 INSERT 语句

-- populate pet table
INSERT INTO pet VALUES ( 'Max', 'Dog', '2018-07-05', 'Jane');
INSERT INTO pet VALUES ( 'Susie', 'Cat', '2019-05-01', 'Phil');
INSERT INTO pet VALUES ( 'Lester', 'Hamster', '2020-06-23', 'Lily');
INSERT INTO pet VALUES ( 'Quincy', 'Parrot', '2013-08-11', 'Anne');

随后可以创建一个 SQLExecuteQueryOperator 任务来向 pet 表写入数据。

populate_pet_table = SQLExecuteQueryOperator(
    task_id="populate_pet_table",
    conn_id="postgres_default",
    sql="sql/pet_schema.sql",
)

从 Postgres 数据库表中获取记录

从 Postgres 数据库表中获取记录可以如此简洁

get_all_pets = SQLExecuteQueryOperator(
    task_id="get_all_pets",
    conn_id="postgres_default",
    sql="SELECT * FROM pet;",
)

向 SQLExecuteQueryOperator 传递参数(针对 Postgres)

SQLExecuteQueryOperator 提供了 parameters 属性,能够在运行时动态注入 SQL 请求中的值。BaseOperator 类自带的 params 属性同样在 SQLExecuteQueryOperator 中可用。虽然 parametersparams 都能实现动态传参,但它们的使用方式略有差别,下面的示例已作说明。

如果要查询两个日期之间所有宠物的出生日期,并在代码中直接写 SQL 语句,则使用 parameters 属性。

get_birth_date = SQLExecuteQueryOperator(
    task_id="get_birth_date",
    conn_id="postgres_default",
    sql="SELECT * FROM pet WHERE birth_date BETWEEN SYMMETRIC %(begin_date)s AND %(end_date)s",
    parameters={"begin_date": "2020-01-01", "end_date": "2020-12-31"},
)

现在让我们重构 get_birth_date 任务。与其在代码里硬写 SQL,不如把 SQL 放到单独的文件中。同时,这次我们使用从父类 BaseOperator 继承而来的 params 属性。

-- dags/sql/birth_date.sql
SELECT * FROM pet WHERE birth_date BETWEEN SYMMETRIC {{ params.begin_date }} AND {{ params.end_date }};
get_birth_date = SQLExecuteQueryOperator(
    task_id="get_birth_date",
    conn_id="postgres_default",
    sql="sql/birth_date.sql",
    params={"begin_date": "2020-01-01", "end_date": "2020-12-31"},
)

启用客户端日志记录数据库消息

SQLExecuteQueryOperator 提供了 hook_params 属性,可向 DbApiHook 额外传递参数。使用 enable_log_db_messages 可以记录数据库消息或由 RAISE 语句抛出的错误。

call_proc = SQLExecuteQueryOperator(
    task_id="call_proc",
    conn_id="postgres_default",
    sql="call proc();",
    hook_params={"enable_log_db_messages": True},
)

向 PostgresOperator 传递服务器配置参数

SQLExecuteQueryOperator 同样提供 hook_params,可用于向 DbApiHook 传递额外参数。通过此方式可以传入 options 参数,从而在连接建立时向服务器发送 命令行选项

tests/system/postgres/example_postgres.py[source]

    get_birth_date = SQLExecuteQueryOperator(
        task_id="get_birth_date",
        sql="SELECT * FROM pet WHERE birth_date BETWEEN SYMMETRIC %(begin_date)s AND %(end_date)s",
        parameters={"begin_date": "2020-01-01", "end_date": "2020-12-31"},
        hook_params={"options": "-c statement_timeout=3000ms"},
    )

完整的 Postgres Operator DAG

当我们把所有内容组合在一起时,Dag 应该是下面的样子

tests/system/postgres/example_postgres.py[source]



# create_pet_table, populate_pet_table, get_all_pets, and get_birth_date are examples of tasks created by
# instantiating the Postgres Operator

ENV_ID = os.environ.get("SYSTEM_TESTS_ENV_ID")
DAG_ID = "postgres_operator_dag"

with DAG(
    dag_id=DAG_ID,
    start_date=datetime.datetime(2020, 2, 2),
    schedule="@once",
    catchup=False,
) as dag:
    create_pet_table = SQLExecuteQueryOperator(
        task_id="create_pet_table",
        sql="""
            CREATE TABLE IF NOT EXISTS pet (
            pet_id SERIAL PRIMARY KEY,
            name VARCHAR NOT NULL,
            pet_type VARCHAR NOT NULL,
            birth_date DATE NOT NULL,
            OWNER VARCHAR NOT NULL);
          """,
    )
    populate_pet_table = SQLExecuteQueryOperator(
        task_id="populate_pet_table",
        sql="""
            INSERT INTO pet (name, pet_type, birth_date, OWNER)
            VALUES ( 'Max', 'Dog', '2018-07-05', 'Jane');
            INSERT INTO pet (name, pet_type, birth_date, OWNER)
            VALUES ( 'Susie', 'Cat', '2019-05-01', 'Phil');
            INSERT INTO pet (name, pet_type, birth_date, OWNER)
            VALUES ( 'Lester', 'Hamster', '2020-06-23', 'Lily');
            INSERT INTO pet (name, pet_type, birth_date, OWNER)
            VALUES ( 'Quincy', 'Parrot', '2013-08-11', 'Anne');
            """,
    )
    get_all_pets = SQLExecuteQueryOperator(task_id="get_all_pets", sql="SELECT * FROM pet;")
    get_birth_date = SQLExecuteQueryOperator(
        task_id="get_birth_date",
        sql="SELECT * FROM pet WHERE birth_date BETWEEN SYMMETRIC %(begin_date)s AND %(end_date)s",
        parameters={"begin_date": "2020-01-01", "end_date": "2020-12-31"},
        hook_params={"options": "-c statement_timeout=3000ms"},
    )

    create_pet_table >> populate_pet_table >> get_all_pets >> get_birth_date

结论

在本操作指南中,我们演示了如何使用 Apache Airflow 的 SQLExecuteQueryOperator 连接 PostgreSQL 数据库。下面快速回顾关键要点:建议在 dags 目录下创建一个名为 sql 的子目录,用来存放所有 SQL 文件,这样代码更简洁、易于维护。最后,我们展示了如何通过 parametersparams 属性在运行时动态传递参数,以及如何通过 hook_params 属性向会话传递选项。

此条目是否有帮助?