Postgres 使用 SQLExecuteQueryOperator 的操作指南¶
简介¶
Apache Airflow 拥有大量算子,可用于实现工作流中的各种任务。Airflow 本质上是一个有向无环图(Directed Acyclic Graph),由任务(节点)和依赖关系(边)组成。
由算子定义或实现的任务是数据管道中的一个工作单元。
本指南的目的在于展示如何使用 SQLExecuteQueryOperator 与 PostgreSQL 数据库进行交互。
注意
以前使用 PostgresOperator 来完成此类操作。该算子已被弃用并移除,请改用 SQLExecuteQueryOperator。
使用 SQLExecuteQueryOperator 的常见数据库操作¶
要使用 SQLExecuteQueryOperator 执行 PostgreSQL 请求,需要两个参数:sql 和 conn_id。这两个参数最终会传递给直接与 Postgres 数据库交互的 DbApiHook 对象。
创建 Postgres 数据库表¶
以下代码片段基于 Airflow‑2.0
# 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 中可用。虽然 parameters 与 params 都能实现动态传参,但它们的使用方式略有差别,下面的示例已作说明。
如果要查询两个日期之间所有宠物的出生日期,并在代码中直接写 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 参数,从而在连接建立时向服务器发送 命令行选项。
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 应该是下面的样子
# 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 文件,这样代码更简洁、易于维护。最后,我们展示了如何通过 parameters 或 params 属性在运行时动态传递参数,以及如何通过 hook_params 属性向会话传递选项。