代码之家  ›  专栏  ›  技术社区  ›  Damien

气流-无法在雪花中创建表格

  •  0
  • Damien  · 技术社区  · 4 年前

    我在跑步。2.2在AWS上使用MWAA服务

    我有下面的DAG

    import airflow
    from airflow import DAG
    from airflow.operators.python_operator import PythonOperator
    import time
    from datetime import timedelta
    
    import snowflake.connector
    from airflow.hooks.base_hook import BaseHook
    
    datestr = time.strftime("%Y%m%d")
    
    connection = BaseHook.get_connection('test_snowflake')
    snow_user = connection.login
    snow_pass = connection.password
    snow_host = connection.host
    snow_schema = connection.schema
    
    default_args = {
        "owner": "test",
        "depends_on_past": False,
        "email": ["info@test.com"],
        "email_on_failure": False,
        "email_on_retry": False,
        "retries": 1,
        "retry_delay": timedelta(minutes=5)
    }
    
    def create_external_table():
        source = "S3_CMS_CMISDB_CMSTRFVW"
        ctx = snowflake.connector.connect(
            user=snow_user,
            password=snow_pass,
            account=snow_host,
            database="aflow_test",
            schema=snow_schema
        )
    
        ctx.cursor().execute("USE ROLE DF_DEV_SYSADMIN_FR;")
        ctx.cursor().execute("USE WAREHOUSE DF_DEV_SERVICE_AIRFLOW_WH;")
        ctx.cursor().execute("USE DATABASE RAW_DB;")
        try:
            path = 'cms/cmisdb/cmstrfvw/dateday=' + datestr
            sql_create = "CREATE OR REPLACE EXTERNAL TABLE " + source + \
                " with location = @stage_devraw/" + path + \
                " auto_refresh = true file_format = (type = parquet);"
            print(f"sql_create: {sql_create}")
            ctx.cursor().execute(sql_create)
    
        finally:
            ctx.cursor().close()
        ctx.close()
    
    
    dag = DAG(
        "simple_dag",
        default_args=default_args,
        start_date=airflow.utils.dates.days_ago(1),
        description="generic simple dag",
        schedule_interval="@daily",
        catchup=False,
    )
    
    
    create_external_table_task = PythonOperator(
        task_id='create_external_table',
        python_callable=create_external_table,
        dag=dag,
    )
    
    # flow
    create_external_table_task
    

    当我执行这个DAG时,我得到以下错误 使用舞台区域时发生故障。原因:[您提供的AWS访问密钥Id无效。]

    aws_default的连接是使用有效的aws访问和密钥设置的 我为aws_access_key和aws_secret_key设置了变量

    有人对我需要做什么才能从气流连接到雪花有什么建议吗

    0 回复  |  直到 4 年前
        1
  •  1
  •   Elad Kalif    4 年前

    你得到的错误其实与气流无关——它来自雪花。提供的凭据很可能缺少列表权限。看到这个了吗 answer 获取有关它的信息。

    还要注意的是,如果你使用的是气流,那么你可以使用它的力量。不需要实现Airflow已经支持的代码。你创造了一个 PythonOperator 连接到snowflake并运行查询。对于这个用例,只需使用 SnowflakeOperator .

    例子:

    from airflow.providers.snowflake.operators.snowflake import SnowflakeOperator
    snowflake_op_sql_str = SnowflakeOperator(
        task_id='snowflake_sql',
        snowflake_conn_id = 'my_snowflake_conn'
        sql="Select 1",
        warehouse="my_warehouse",
        database="my_db",
        schema="my_schema",
        role="my_role",
    )
    

    使用操作员可以省去验证、使用光标等的麻烦。。。 您需要在气流中定义雪花连接。你可以用这个 doc 获取有关它的信息。

    如果出于某种原因 雪花操作员 如果没有所需的功能,则可以创建自定义运算符或使用 SnowflakeHook 与蟒蛇共舞。

    推荐文章