damp-queen-61493
01/19/2022, 4:53 PM[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 thisdamp-queen-61493
01/19/2022, 4:56 PMfrom 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
)damp-queen-61493
01/19/2022, 4:58 PM[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 checklittle-megabyte-1074
square-activity-64562
01/25/2022, 5:57 AMTask 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.damp-queen-61493
01/26/2022, 12:18 PM