【发布时间】:2022-10-16 01:14:30
【问题描述】:
我需要创建一个 DAG,它会根据模式名称将 sql 插入到 db 表中。
DAG 示例:
from datetime import datetime
from airflow import DAG, utils
from airflow.providers.snowflake.operators.snowflake import SnowflakeOperator
CONNECTION_ID = ...
WAREHOUSE = ...
DATABASE = ...
ROLE = ...
SQL_STRING = (
"SELECT SCHEMA,TABLE FROM LOG_SCHEMA.LOG001;"
)
dag = DAG(
'my_test',
start_date=utils.dates.days_ago(1),
default_args={'connection_id': CONNECTION_ID},
catchup=False,
)
my_sql = SnowflakeOperator(
task_id='my_sql',
dag=dag,
sql=SQL_STRING,
warehouse=WAREHOUSE,
database=DATABASE,
role=ROLE,
)
my_sql
在我的示例中 my_sql 的输出只是模式名称和表名称。我想使用它来执行插入。 例子:
INSERT INTO SCHEMA.TABLE SELELECT * FROM SCHEMA.TABLE WHERE COL1=2;
我将使用模式名称导入我的变量,并根据我的需要选择模式 TEST,以便对此模式中的所有表执行插入操作。
INSERT INTO TEST.TABLE001 SELELECT * FROM TEST.TABLE001 WHERE COL1=2;
INSERT INTO TEST.TABLE002 SELELECT * FROM TEST.TABLE002 WHERE COL1=2;
INSERT INTO TEST.TABLE003 SELELECT * FROM TEST.TABLE003 WHERE COL1=2;
【问题讨论】:
-
谁能回答?
标签: airflow