【问题标题】:Airflow dag for sql insert statements用于 sql 插入语句的气流 dag
【发布时间】: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


【解决方案1】:

您有两种选择来实现这一目标:

  • 使用 Airflow:通过创建一个读取第一个查询结果并准备新查询的 python 任务,您可以在单独的SnowflakeOperator 任务中执行。
  • 使用雪花和SQL:SQL支持循环,你可以将结果存储在游标中,然后循环游标记录以调用新语句(我不确定雪花中的语法,所以也许你需要修复它):
DECLARE
    c1 CURSOR FOR SELECT SCHEMA,TABLE FROM LOG_SCHEMA.LOG001;
    schema_col varchar;
    table_col varchar;
    sql varchar default 'INSERT INTO ?.? SELECT * FROM ?.? WHERE COL1=2';
BEGIN
  FOR record IN c1 DO
      schema_col := record.SCHEMA;
      table_col := record.TABLE;
      execute immediate :query using(schema_col, table_col, schema_col, table_col)
  END FOR;      
  RETURN 'End';
END;

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2023-01-02
    • 1970-01-01
    • 2015-04-03
    • 1970-01-01
    • 2023-03-17
    • 1970-01-01
    相关资源
    最近更新 更多