Testing Airflow DAG for errors prior to trigger

Viewed 38

Is there a quick/simple way to check an Airflow (1.10) DAG script for errors such as syntax errors, task errors etc from the command line (or any other way) prior to triggering the DAG ?

In our project we have to first push the DAG code to Stash and then deploy it thru Jenkins before we can actually run the DAG. Hence we can find out errors only after triggering it.

This becomes very tedious when we have to repeatedly go thru the same process to discover and fix every minor coding errors.

So checking if there is a quick way to check and fix for these errors before we trigger the DAG.

Airflow Version - 1.10

Thanks

1 Answers

There are some steps you can perform after pushing DAGs to production and avoid going back and forth.

DAG Loader Test

This test should ensure that your DAG does not contain a piece of code that raises error while loading. No additional code needs to be written by the user to run this test.

python your-dag-file.py

Running the above command without any error ensures your DAG does not contain any uninstalled dependency, syntax errors, etc.

You can look into Testing a DAG for details on how to test individual operators.

Unit tests

Unit tests ensure that there is no incorrect code in your DAG. You can write a unit test for your tasks as well as your DAG.

Unit test for loading a DAG:

from airflow.models import DagBag
import unittest

class TestHelloWorldDAG(unittest.TestCase):
   @classmethod
   def setUpClass(cls):
       cls.dagbag = DagBag()

   def test_dag_loaded(self):
       dag = self.dagbag.get_dag(dag_id='hello_world')
       self.assertDictEqual(self.dagbag.import_errors, {})
       self.assertIsNotNone(dag)
       self.assertEqual(len(dag.tasks), 1)

Unit test a DAG structure: This is an example test want to verify the structure of a code-generated DAG against a dict object

import unittest
class testClass(unittest.TestCase):
    def assertDagDictEqual(self,source,dag):
        self.assertEqual(dag.task_dict.keys(),source.keys())
        for task_id,downstream_list in source.items():
            self.assertTrue(dag.has_task(task_id), msg="Missing task_id: {} in dag".format(task_id))
            task = dag.get_task(task_id)
            self.assertEqual(task.downstream_task_ids, set(downstream_list),
                             msg="unexpected downstream link in {}".format(task_id))
    def test_dag(self):
        self.assertDagDictEqual({
          "DummyInstruction_0": ["DummyInstruction_1"],
          "DummyInstruction_1": ["DummyInstruction_2"],
          "DummyInstruction_2": ["DummyInstruction_3"],
          "DummyInstruction_3": []
        },dag)

Unit test for custom operator:

import unittest
from airflow.utils.state import State

DEFAULT_DATE = '2019-10-03'
TEST_DAG_ID = 'test_my_custom_operator'

class MyCustomOperatorTest(unittest.TestCase):
   def setUp(self):
       self.dag = DAG(TEST_DAG_ID, schedule_interval='@daily', default_args={'start_date' : DEFAULT_DATE})
       self.op = MyCustomOperator(
           dag=self.dag,
           task_id='test',
           prefix='s3://bucket/some/prefix',
       )
       self.ti = TaskInstance(task=self.op, execution_date=DEFAULT_DATE)

   def test_execute_no_trigger(self):
       self.ti.run(ignore_ti_state=True)
       self.assertEqual(self.ti.state, State.SUCCESS)
       #Assert something related to tasks results

Self-Checks

You can also implement checks in a DAG to make sure the tasks are producing the results as expected. As an example, if you have a task that pushes data to S3, you can implement a check in the next task. For example, the check could make sure that the partition is created in S3 and perform some simple checks to see if the data is correct or not.

Similarly, if you have a task that starts a microservice in Kubernetes or Mesos, you should check if the service has started or not using airflow.sensors.http_sensor.HttpSensor.

task = PushToS3(...)
check = S3KeySensor(
   task_id='check_parquet_exists',
   bucket_key="s3://bucket/key/foo.parquet",
   poke_interval=0,
   timeout=0
)
task >> check

Staging environment

If possible, keep a staging environment to test the complete DAG run before deploying in the production. Make sure your DAG is parameterized to change the variables, e.g., the output path of S3 operation or the database used to read the configuration. Do not hard code values inside the DAG and then change them manually according to the environment.

You can use environment variables to parameterize the DAG.

import os

dest = os.environ.get(
   "MY_DAG_DEST_PATH",
   "s3://default-target/path/"
)

Source: Apache airflow docs.

Related