From 33870f9c66a5c322720ec3dfce2c47cd5a758258 Mon Sep 17 00:00:00 2001 From: Hussein Awala Date: Mon, 19 Jun 2023 20:54:18 +0200 Subject: [PATCH 1/2] Add a check on none TIs for the current execution date Signed-off-by: Hussein Awala --- airflow/sensors/external_task.py | 13 +++++++++---- tests/sensors/test_external_task_sensor.py | 20 ++++++++++++++++++++ 2 files changed, 29 insertions(+), 4 deletions(-) diff --git a/airflow/sensors/external_task.py b/airflow/sensors/external_task.py index 959ebe5131bd3..badda224a4561 100644 --- a/airflow/sensors/external_task.py +++ b/airflow/sensors/external_task.py @@ -367,10 +367,15 @@ def get_count(self, dttm_filter, session, states) -> int: elif self.external_task_group_id: external_task_group_task_ids = self.get_external_task_group_task_ids(session, dttm_filter) count = ( - self._count_query(TI, session, states, dttm_filter) - .filter(tuple_in_condition((TI.task_id, TI.map_index), external_task_group_task_ids)) - .scalar() - ) / len(external_task_group_task_ids) + 0 + if not external_task_group_task_ids + else ( + self._count_query(TI, session, states, dttm_filter) + .filter(tuple_in_condition((TI.task_id, TI.map_index), external_task_group_task_ids)) + .scalar() + ) + / len(external_task_group_task_ids) + ) else: count = self._count_query(DR, session, states, dttm_filter).scalar() return count diff --git a/tests/sensors/test_external_task_sensor.py b/tests/sensors/test_external_task_sensor.py index 616c9c5cabbcc..a5259084b13f6 100644 --- a/tests/sensors/test_external_task_sensor.py +++ b/tests/sensors/test_external_task_sensor.py @@ -808,6 +808,26 @@ def test_external_task_group_with_mapped_tasks_failed_states(self): ): op.run(start_date=DEFAULT_DATE, end_date=DEFAULT_DATE, ignore_ti_state=True) + def test_external_task_group_when_there_is_no_TIs(self): + """Test that the sensor does not fail when there are no TIs to check.""" + self.add_time_sensor() + self.add_dummy_task_group_with_dynamic_tasks(State.FAILED) + op = ExternalTaskSensor( + task_id="test_external_task_sensor_check", + external_dag_id=TEST_DAG_ID, + external_task_group_id=TEST_TASK_GROUP_ID, + failed_states=[State.FAILED], + dag=self.dag, + poke_interval=1, + timeout=3, + ) + with pytest.raises(AirflowSensorTimeout): + op.run( + start_date=DEFAULT_DATE + timedelta(hours=1), + end_date=DEFAULT_DATE + timedelta(hours=1), + ignore_ti_state=True, + ) + def test_external_task_sensor_check_zipped_dag_existence(dag_zip_maker): with dag_zip_maker("test_external_task_sensor_check_existense.py") as dagbag: From 9364a4eb3a227d141f5e945dcee192363933cf26 Mon Sep 17 00:00:00 2001 From: Hussein Awala Date: Tue, 20 Jun 2023 20:31:27 +0200 Subject: [PATCH 2/2] replace inline if-else by old one Signed-off-by: Hussein Awala --- airflow/sensors/external_task.py | 12 +++++------- 1 file changed, 5 insertions(+), 7 deletions(-) diff --git a/airflow/sensors/external_task.py b/airflow/sensors/external_task.py index badda224a4561..69ba41ef09354 100644 --- a/airflow/sensors/external_task.py +++ b/airflow/sensors/external_task.py @@ -366,16 +366,14 @@ def get_count(self, dttm_filter, session, states) -> int: ) / len(self.external_task_ids) elif self.external_task_group_id: external_task_group_task_ids = self.get_external_task_group_task_ids(session, dttm_filter) - count = ( - 0 - if not external_task_group_task_ids - else ( + if not external_task_group_task_ids: + count = 0 + else: + count = ( self._count_query(TI, session, states, dttm_filter) .filter(tuple_in_condition((TI.task_id, TI.map_index), external_task_group_task_ids)) .scalar() - ) - / len(external_task_group_task_ids) - ) + ) / len(external_task_group_task_ids) else: count = self._count_query(DR, session, states, dttm_filter).scalar() return count