Empresas
Empregos
  • Sobre nós
  • Soluções
    • Publicação de vagas
      Publique sua vaga e receba candidatos qualificados em 48h.
    • Avaliações de candidatos
      Mais de 500 testes técnicos e psicológicos, mais anti-fraude.
    • Headhunting
      Busca executiva personalizada do início ao fim.
    • Folha de Pagamento + EOR
      Dispersão de folha e EOR em mais de 15 países da LATAM.
  • Preços
  • Empregos

0

606
Visualizações
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 Respostas
Responde à pergunta

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 Relatório
Responde à pergunta
Encontrar trabalhos remotos

Descubra a nova forma de encontrar um emprego!

melhores empregos
Principais categorias de trabalho
Empresas
Postar vaga Preços Comercial
Jurídico
Termos e Condições Política de privacidade
© 2026 PeakU Inc. All Rights Reserved.
Andres GPT
Recomende algumas ofertas para mim
Preciso de ajuda