Business
Jobs
  • About Us
  • Solutions
    • Job Postings
      Post your job and receive qualified candidates in 48h.
    • Candidate Assessments
      500+ technical and psychological tests, plus anti-fraud.
    • Headhunting
      Tailor-made executive search from start to finish.
    • Payroll + EOR
      Payroll dispersal and EOR across 15+ LATAM countries.
  • Pricing
  • Jobs

0

605
Views
No se encontró la ejecución de Dag cuando la unidad prueba un operador personalizado en Airflow

Escribí un operador personalizado (DataCleaningOperator), que corrige los datos JSON en función de un esquema proporcionado.

Las pruebas unitarias funcionaban anteriormente cuando no tenía que crear una instancia de TaskInstance y proporcionar un contexto al operador. Sin embargo, actualicé el operador recientemente para incluir un contexto (para que pueda usar xcom_push).

Aquí hay un ejemplo de una de las pruebas:

 DEFAULT_DATE = datetime.today() class TestDataCleaningOperator(unittest.TestCase): """ Class to execute unit tests for the operator 'DataCleaningOperator'. """ def setUp(self) -> None: super().setUp() self.dag = DAG( dag_id="test_dag_data_cleaning", schedule_interval=None, default_args={ "owner": "airflow", "start_date": DEFAULT_DATE, "output_to_xcom": True, }, ) self._initialise_test_data() def _initialize_test_data() -> None: # Test data set here as class variables such as self.test_data_correct ... def test_operator_cleans_dataset_which_matches_schema(self) -> None: """ Test: Attempt to clean a dataset which matches the provided schema. Verification: Returns the original dataset, unchanged. """ task = DataCleaningOperator( task_id="test_operator_cleans_dataset_which_matches_schema", schema_fields=self.test_schema_nest, data_file_object=deepcopy(self.test_data_correct), dag=self.dag, ) ti = TaskInstance(task=task, execution_date=DEFAULT_DATE) result: List[dict] = task.execute(ti.get_template_context()) self.assertEqual(result, self.test_data_correct)

Sin embargo, cuando se ejecutan las pruebas, se genera el siguiente error:

 airflow.exceptions.DagRunNotFound: DagRun for 'test_dag_data_cleaning' with date 2022-02-22 12:09:51.538954+00:00 not found

Esto está relacionado con la línea en la que se instancia una instancia de tarea en test_operator_cleans_dataset_which_matches_schema.

¿Por qué Airflow no puede localizar el DAG test_dag_data_cleaning? ¿Hay alguna configuración específica que me haya perdido? ¿Debo crear también una instancia de ejecución de DAG o agregar el DAG a la bolsa de dag manualmente si este dag de prueba está fuera de mi directorio de DAG estándar? Todos los dags normales (que no son de prueba) en mi directorio dag se ejecutan correctamente.

En caso de que ayude, mi versión actual de Airflow es 2.2.3 y la estructura de mi proyecto es:

 airflow ├─ dags ├─ plugins | ├─ ... | └─ operators | ├─ ... | └─ data_cleaning_operator.py | └─ tests ├─ ... └─ operators └─ test_data_cleaning_operator.py
over 4 years ago · Santiago Trujillo
1 answers
Answer question

0

El código que se ha escrito utiliza el formato Airflow 2.0 de prueba unitaria . Entonces, cuando actualizó a Airflow 2.2.3, la prueba unitaria requiere que cree un dagrun antes de crear una ejecución de prueba.

A continuación se muestra el código de muestra que funcionó para mí:

 import unittest import pendulum from airflow import DAG from airflow.utils.state import DagRunState from airflow.utils.types import DagRunType from operators.test_operator import EvenNumberCheckOperator DEFAULT_DATE = pendulum.datetime(2022, 3, 4, tz='America/Toronto') TEST_DAG_ID = "my_custom_operator_dag" TEST_TASK_ID = "my_custom_operator_task" class TestEvenNumberCheckOperator(unittest.TestCase): def setUp(self): super().setUp() self.dag = DAG('test_dag4', default_args={'owner': 'airflow', 'start_date': DEFAULT_DATE}) self.even = 10 self.odd = 11 EvenNumberCheckOperator( task_id=TEST_TASK_ID, my_operator_param=self.even, dag=self.dag ) def test_even(self): """Tests that the EvenNumberCheckOperator returns True for 10.""" dagrun = self.dag.create_dagrun(state=DagRunState.RUNNING, execution_date=DEFAULT_DATE, #data_interval=DEFAULT_DATE, start_date=DEFAULT_DATE, run_type=DagRunType.MANUAL) ti = dagrun.get_task_instance(task_id=TEST_TASK_ID) ti.task = self.dag.get_task(task_id=TEST_TASK_ID) result = ti.task.execute(ti.get_template_context()) assert result is True
over 4 years ago · Santiago Trujillo Report
Answer question
Find remote jobs

Discover the new way to find a job!

Top jobs
Top job categories
Business
Post vacancy Pricing Sales
Legal
Terms and conditions Privacy policy
© 2026 PeakU Inc. All Rights Reserved.
Andres GPT
Show me some job opportunities
There's an error!