【问题标题】:Task exited with return code Negsignal.SIGABRT. Airflow task fails with SnowflakeOperator任务以返回码 Negsignal.SIGABRT 退出。气流任务因 SnowflakeOperator 而失败
【发布时间】:2021-05-06 18:30:12
【问题描述】:

在气流中运行此 DAG 会出现错误,因为任务退出并返回代码 Negsignal.SIGABRT。

我不知道我做错了什么


    from airflow import DAG
    from airflow.providers.snowflake.operators.snowflake import SnowflakeOperator
    from airflow.utils.dates import days_ago
    
    SNOWFLAKE_CONN_ID = 'snowflake_conn'
    # TODO: should be able to rely on connection's schema, but currently param required by S3ToSnowflakeTransfer
    # SNOWFLAKE_SCHEMA = 'schema_name'
    #SNOWFLAKE_STAGE = 'stage_name'
    SNOWFLAKE_WAREHOUSE = 'SF_TUTS_WH'
    SNOWFLAKE_DATABASE = 'KAFKA_DB'
    SNOWFLAKE_ROLE = 'sysadmin'
    SNOWFLAKE_SAMPLE_TABLE = 'sample_table'
    
    CREATE_TABLE_SQL_STRING = (
        f"CREATE OR REPLACE TRANSIENT TABLE {SNOWFLAKE_SAMPLE_TABLE} (name VARCHAR(250), id INT);"
    )
    
    SQL_INSERT_STATEMENT = f"INSERT INTO {SNOWFLAKE_SAMPLE_TABLE} VALUES ('name', %(id)s)"
    SQL_LIST = [SQL_INSERT_STATEMENT % {"id": n} for n in range(0, 10)]
    
    default_args = {
        'owner': 'airflow',
    }
    
    dag = DAG(
        'example_snowflake',
        default_args=default_args,
        start_date=days_ago(2),
        tags=['example'],
    )
    
    snowflake_op_sql_str = SnowflakeOperator(
        task_id='snowflake_op_sql_str',
        dag=dag,
        snowflake_conn_id=SNOWFLAKE_CONN_ID,
        sql=CREATE_TABLE_SQL_STRING,
        warehouse=SNOWFLAKE_WAREHOUSE,
        database=SNOWFLAKE_DATABASE,
      #  schema=SNOWFLAKE_SCHEMA,
        role=SNOWFLAKE_ROLE,
    )
    
    snowflake_op_with_params = SnowflakeOperator(
        task_id='snowflake_op_with_params',
        dag=dag,
        snowflake_conn_id=SNOWFLAKE_CONN_ID,
        sql=SQL_INSERT_STATEMENT,
        parameters={"id": 56},
        warehouse=SNOWFLAKE_WAREHOUSE,
        database=SNOWFLAKE_DATABASE,
     #   schema=SNOWFLAKE_SCHEMA,
        role=SNOWFLAKE_ROLE,
    )
    
    
    snowflake_op_sql_list = SnowflakeOperator(
        task_id='snowflake_op_sql_list', dag=dag, snowflake_conn_id=SNOWFLAKE_CONN_ID, sql=SQL_LIST
    )
    
    snowflake_op_sql_str >> [
        snowflake_op_with_params,
        snowflake_op_sql_list,]

在airFlow中获取LOGS如下::

 Reading local file: /Users/aashayjain/airflow/logs/snowflake_test/snowflake_op_with_params/2021-02-02T13:51:18.229233+00:00/1.log
[2021-02-02 19:21:38,880] {taskinstance.py:826} INFO - Dependencies all met for <TaskInstance: snowflake_test.snowflake_op_with_params 2021-02-02T13:51:18.229233+00:00 [queued]>
[2021-02-02 19:21:38,887] {taskinstance.py:826} INFO - Dependencies all met for <TaskInstance: snowflake_test.snowflake_op_with_params 2021-02-02T13:51:18.229233+00:00 [queued]>
[2021-02-02 19:21:38,887] {taskinstance.py:1017} INFO - 
--------------------------------------------------------------------------------
[2021-02-02 19:21:38,887] {taskinstance.py:1018} INFO - Starting attempt 1 of 1
[2021-02-02 19:21:38,887] {taskinstance.py:1019} INFO - 
--------------------------------------------------------------------------------
[2021-02-02 19:21:38,892] {taskinstance.py:1038} INFO - Executing <Task(SnowflakeOperator): snowflake_op_with_params> on 2021-02-02T13:51:18.229233+00:00
[2021-02-02 19:21:38,895] {standard_task_runner.py:51} INFO - Started process 16510 to run task
[2021-02-02 19:21:38,901] {standard_task_runner.py:75} INFO - Running: ['airflow', 'tasks', 'run', 'snowflake_test', 'snowflake_op_with_params', '2021-02-02T13:51:18.229233+00:00', '--job-id', '7', '--pool', 'default_pool', '--raw', '--subdir', 'DAGS_FOLDER/snowflake_test.py', '--cfg-path', '/var/folders/6h/1pzt4pbx6h32h6p5v503wws00000gp/T/tmp1w61m38s']
[2021-02-02 19:21:38,903] {standard_task_runner.py:76} INFO - Job 7: Subtask snowflake_op_with_params
[2021-02-02 19:21:38,933] {logging_mixin.py:103} INFO - Running <TaskInstance: snowflake_test.snowflake_op_with_params 2021-02-02T13:51:18.229233+00:00 [running]> on host 1.0.0.127.in-addr.arpa
[2021-02-02 19:21:38,954] {taskinstance.py:1232} INFO - Exporting the following env vars:
AIRFLOW_CTX_DAG_OWNER=airflow
AIRFLOW_CTX_DAG_ID=snowflake_test
AIRFLOW_CTX_TASK_ID=snowflake_op_with_params
AIRFLOW_CTX_EXECUTION_DATE=2021-02-02T13:51:18.229233+00:00
AIRFLOW_CTX_DAG_RUN_ID=manual__2021-02-02T13:51:18.229233+00:00
[2021-02-02 19:21:38,955] {snowflake.py:119} INFO - Executing: INSERT INTO TEST_TABLE VALUES ('name', %(id)s)
[2021-02-02 19:21:38,961] {base.py:74} INFO - Using connection to: id: snowflake_conn. Host: uva00063.us-east-1.snowflakecomputing.com, Port: None, Schema: , Login: aashay, Password: XXXXXXXX, extra: XXXXXXXX
[2021-02-02 19:21:38,963] {connection.py:218} INFO - Snowflake Connector for Python Version: 2.3.7, Python Version: 3.7.3, Platform: Darwin-19.5.0-x86_64-i386-64bit
[2021-02-02 19:21:38,964] {connection.py:769} INFO - This connection is in OCSP Fail Open Mode. TLS Certificates would be checked for validity and revocation status. Any other Certificate Revocation related exceptions or OCSP Responder failures would be disregarded in favor of connectivity.
[2021-02-02 19:21:38,964] {connection.py:785} INFO - Setting use_openssl_only mode to False
[2021-02-02 19:21:38,996] {local_task_job.py:118} INFO - Task exited with return code Negsignal.SIGABRT

apache-airflow==2.0.0

python 3.7.3

期待在这方面的帮助。让我知道我需要提供更多详细信息。代码或气流........................................................

【问题讨论】:

    标签: airflow snowflake-cloud-data-platform airflow-scheduler snowflake-schema


    【解决方案1】:

    你正在执行:

    snowflake_op_with_params = SnowflakeOperator(
        task_id='snowflake_op_with_params',
        dag=dag,
        snowflake_conn_id=SNOWFLAKE_CONN_ID,
        sql=SQL_INSERT_STATEMENT,
        parameters={"id": 56},
        warehouse=SNOWFLAKE_WAREHOUSE,
        database=SNOWFLAKE_DATABASE,
     #   schema=SNOWFLAKE_SCHEMA,
        role=SNOWFLAKE_ROLE,
    )
    

    这会尝试在SQL_INSERT_STATEMENT 中运行sql。 所以它执行:

    f"INSERT INTO {SNOWFLAKE_SAMPLE_TABLE} VALUES ('name', %(id)s)"
    

    给出:

    INSERT INTO sample_table VALUES ('name', %(id)s)
    

    如你自己的日志所示:

    [2021-02-02 19:21:38,955] {snowflake.py:119} INFO - Executing: INSERT INTO TEST_TABLE VALUES ('name', %(id)s)
    

    这不是有效的 SQL 语句。

    我真的不知道你想执行什么 SQL。基于SQL_LIST 我可以假设%(id)s 假设是整数类型的id。

    【讨论】:

    • 我硬编码了 sql 语句 INFO - 正在执行:INSERT INTO KAFKA_DB.KAFKA_SCHEMA.TEST_TABLE VALUES('name', 123) INFO -Snowflake Connector for Python Version: 2.3.7, Python Version: 3.7.3 , 平台:Darwin-19.5.0-x86_64-i386-64bit INFO - 此连接处于 OCSP Fail Open Mode。将检查 TLS 证书的有效性和撤销状态。任何其他与证书吊销相关的异常或 OCSP 响应程序故障都将被忽略以支持连接。 INFO - 将 use_openssl_only 模式设置为 False 信息 - 任务退出并返回代码 Negsignal.SIGABRT
    • 仍然出现同样的错误。在对 SQL sql 语句进行硬编码之后。这个 SQL 是正确的。在雪花中测试 --> INSERT INTO KAFKA_DB.KAFKA_SCHEMA.TEST_TABLE VALUES('name', 123)
    • @Aashu 连接中定义的主机会不会出错?您正在使用 ocsp。 docs.snowflake.com/en/user-guide/... 我相信您的主机应该是:ocsp.*.snowflakecomputing.com:80
    • 我也猜这将是问题所在。我试过使用你的建议。根据我尝试过的文档——ocsp.snowflakecomputing.com:80 也仍然出现同样的错误。感谢您查看这个。非常感谢您的帮助。如果有什么可以解决这个问题,请告诉我您的想法..
    • @Aashu ocsp.snowflakecomputing.com:80 不能作为答案,因为它没有指定您尝试连接的帐户。我看到你已经在雪花 python 连接器 github 中打开了一个支持请求。您可能想为雪花打开一个实际的支持案例。我不会向他们提及 Airflow,这只会让他们感到困惑。您只是想用他们的 python 包装器运行查询。气流不是这里的问题。您可能还想阅读github.com/snowflakedb/snowflake-connector-python/issues/245 与您的问题不同但可能相关的问题
    猜你喜欢
    • 2020-09-12
    • 2020-04-29
    • 2020-08-29
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-08-24
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多