Hi! I'm trying to ingest mssql from airflow. I ge...
# ingestion
d
Hi! I'm trying to ingest mssql from airflow. I get the following error
Copy code
[2022-01-19, 16:43:56 UTC] {local_task_job.py:154} INFO - Task exited with return code Negsignal.SIGKILL
I'm using the inline approach like this
Dag:
Copy code
from datetime import timedelta

from airflow import DAG
from airflow.hooks.base import BaseHook
from airflow.operators.python import PythonOperator
from airflow.utils.dates import days_ago
from datahub.ingestion.run.pipeline import Pipeline

from kubernetes.client import models as k8s

default_args = {
    'owner': 'airflow',
    'depends_on_past': False,
    'email': ['<mailto:airflow@example.com|airflow@example.com>'],
    'email_on_failure': False,
    'email_on_retry': False,
    'start_date': days_ago(5),
    'schedule_interval': "0/60 11-23 * * 1-5",
    'max_active_runs': 1,
    'retries': 0,
    'retry_delay': timedelta(seconds=30),
    'execution_timeout': timedelta(minutes=120)
}


def ingest_from_sqlserver():
    conn = BaseHook.get_connection('dummy')
    host_port = f"{conn.host}:{conn.port}"
    user = conn.login
    password = conn.password
    database = "dummy-db"

    conn_datahub = BaseHook.get_connection('datahub_rest_default')
    datahub_rest = conn_datahub.host

    pipeline = Pipeline.create(
        # This configuration is analogous to a recipe configuration.
        {
            "source": {
                "type": "mssql",
                "config": {
                    "username": user,
                    "password": password,
                    "database": database,
                    "host_port": host_port,
                    "env": "DEV"
                },
            },
            "transformers": [
                {
                    "type": "simple_add_dataset_tags",
                    "config": {
                        "tag_urns": [
                            "urn:li:tag:NeedsDocumentation",
                            "urn:li:tag:sign",
                            "urn:li:tag:ingested_by_airflow"
                        ]
                    }
                }
            ],
            "sink": {
                "type": "datahub-rest",
                "config": {"server": datahub_rest},
            },
        }
    )
    pipeline.run()
    pipeline.raise_from_status()


with DAG(dag_id="mssql_metadata_example",
         default_args=default_args,
         catchup=False) as dag:
    ingest_task = PythonOperator(
        task_id="ingest_from_mysql",
        executor_config={
            "pod_template_file": "/opt/airflow/pod_templates/pod_template_file.yaml",
            "pod_override": k8s.V1Pod(
                spec=k8s.V1PodSpec(
                    containers=[
                        k8s.V1Container(
                            name="base",
                            image="us-east1-docker.pkg.dev/xxxxx/xxxxx/dataengineering-airflow:latest",
                            image_pull_policy="Always")]))
        },
        python_callable=ingest_from_sqlserver
    )
Task log:
Copy code
[2022-01-19, 16:43:48 UTC] {taskinstance.py:1035} INFO - Dependencies all met for <TaskInstance: mssql_metadata_example.ingest_from_mysql manual__2022-01-19T16:42:55.055617+00:00 [queued]>
[2022-01-19, 16:43:48 UTC] {taskinstance.py:1035} INFO - Dependencies all met for <TaskInstance: mssql_metadata_example.ingest_from_mysql manual__2022-01-19T16:42:55.055617+00:00 [queued]>
[2022-01-19, 16:43:48 UTC] {taskinstance.py:1241} INFO - 
--------------------------------------------------------------------------------
[2022-01-19, 16:43:48 UTC] {taskinstance.py:1242} INFO - Starting attempt 1 of 1
[2022-01-19, 16:43:48 UTC] {taskinstance.py:1243} INFO - 
--------------------------------------------------------------------------------
[2022-01-19, 16:43:48 UTC] {taskinstance.py:1262} INFO - Executing <Task(PythonOperator): ingest_from_mysql> on 2022-01-19 16:42:55.055617+00:00
[2022-01-19, 16:43:48 UTC] {standard_task_runner.py:52} INFO - Started process 31 to run task
[2022-01-19, 16:43:48 UTC] {standard_task_runner.py:76} INFO - Running: ['airflow', 'tasks', 'run', 'mssql_metadata_example', 'ingest_from_mysql', 'manual__2022-01-19T16:42:55.055617+00:00', '--job-id', '478', '--raw', '--subdir', 'DAGS_FOLDER/mssql_metadata_ingestion_demo.py', '--cfg-path', '/tmp/tmpfvcndl6x', '--error-file', '/tmp/tmpri6qou8y']
[2022-01-19, 16:43:48 UTC] {standard_task_runner.py:77} INFO - Job 478: Subtask ingest_from_mysql
[2022-01-19, 16:43:48 UTC] {logging_mixin.py:109} INFO - Running <TaskInstance: mssql_metadata_example.ingest_from_mysql manual__2022-01-19T16:42:55.055617+00:00 [running]> on host mssqlmetadataexampleingestfrommysql.d5e32f5a193b4b54a09eb155dda
[2022-01-19, 16:43:48 UTC] {taskinstance.py:1429} INFO - Exporting the following env vars:
AIRFLOW_CTX_DAG_EMAIL=airflow@example.com
AIRFLOW_CTX_DAG_OWNER=airflow
AIRFLOW_CTX_DAG_ID=mssql_metadata_example
AIRFLOW_CTX_TASK_ID=ingest_from_mysql
AIRFLOW_CTX_EXECUTION_DATE=2022-01-19T16:42:55.055617+00:00
AIRFLOW_CTX_DAG_RUN_ID=manual__2022-01-19T16:42:55.055617+00:00
[2022-01-19, 16:43:49 UTC] {base.py:79} INFO - Using connection to: id: dummy. Host: 10.224.0.5, Port: 1433, Schema: , Login: ***, Password: ***, extra: {}
[2022-01-19, 16:43:49 UTC] {base.py:79} INFO - Using connection to: id: datahub_rest_default. Host: <http://datahub-datahub-gms.datahub-dev.svc.cluster.local:8080>, Port: None, Schema: , Login: , Password: None, extra: {}
[2022-01-19, 16:43:56 UTC] {local_task_job.py:154} INFO - Task exited with return code Negsignal.SIGKILL
[2022-01-19, 16:43:56 UTC] {taskinstance.py:1280} INFO - Marking task as FAILED. dag_id=mssql_metadata_example, task_id=ingest_from_mysql, execution_date=20220119T164255, start_date=20220119T164348, end_date=20220119T164356
[2022-01-19, 16:43:56 UTC] {local_task_job.py:264} INFO - 0 downstream tasks scheduled from follow-on schedule check
l
Hi @damp-queen-61493! I’m so sorry for the delay - are you still running into the same issue?
s
I have seen this
Task exited with return code Negsignal.SIGKILL
before when using K8s pods. This usually happened when the resource requests and limits on K8s pod were not specified. If the cluster is running low on resources and the pod is using more than the default resource assigned by Kubernetes (mainly memory, CPU gets compressed) then Kubernetes will kill the pod. That is one possibility but I saw this happen frequently. Although those were ETL jobs not datahub ingestion. This is one avenue you can explore to fix this possibly.
d
Hello @little-megabyte-1074 and @square-activity-64562! Tks for your answers! I've already fixed the issue when setting up memory requests for the pods. But I need to set a 1gb request for the ingest tasks to work properly. I think it's too much.