Empresas
Empleos
  • Sobre nosotros
  • Soluciones
    • Publicación de vacantes
      Publica tu vacante y recibe candidatos calificados en 48h.
    • Evaluación de candidatos
      500+ pruebas técnicas y psicológicas, más anti-fraude.
    • Headhunting
      Búsqueda ejecutiva a la medida de principio a fin.
    • Nómina + EOR
      Dispersión de nómina y EOR en más de 15 países de LATAM.
  • Precios
  • Empleos

0

616
Vistas
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 Respuestas
Responde la pregunta

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 Denunciar
Responde la pregunta
Encuentra empleos remotos

¡Descubre la nueva forma de encontrar empleo!

Top de empleos
Top categorías de empleo
Empresas
Publicar vacante Precios Comercial
Legal
Términos y condiciones Política de privacidad
© 2026 PeakU Inc. All Rights Reserved.
Andres GPT
Recomiéndame algunas ofertas
Necesito ayuda