From 8f246edaaeebe070a89db9220f53940b8e2fa077 Mon Sep 17 00:00:00 2001 From: srinivas chennupati Date: Thu, 9 Jul 2026 12:21:53 -0500 Subject: [PATCH 1/2] Fix ExternalTaskSensor log timezone --- .../providers/standard/sensors/external_task.py | 15 ++++++++++++--- .../standard/sensors/test_external_task_sensor.py | 8 ++++++++ 2 files changed, 20 insertions(+), 3 deletions(-) diff --git a/providers/standard/src/airflow/providers/standard/sensors/external_task.py b/providers/standard/src/airflow/providers/standard/sensors/external_task.py index 8270386754f0a..81fee9a1d5722 100644 --- a/providers/standard/src/airflow/providers/standard/sensors/external_task.py +++ b/providers/standard/src/airflow/providers/standard/sensors/external_task.py @@ -23,12 +23,14 @@ from collections.abc import Callable, Collection, Iterable, Sequence from typing import TYPE_CHECKING, ClassVar +from airflow import settings from airflow.models.dag import DagModel from airflow.providers.common.compat.sdk import ( AirflowSkipException, BaseOperatorLink, BaseSensorOperator, conf, + timezone, ) from airflow.providers.standard.exceptions import ( DuplicateStateError, @@ -280,6 +282,12 @@ def _get_dttm_filter(self, context: Context) -> Sequence[datetime.datetime]: def _serialize_dttm_filter(dttm_filter: Sequence[datetime.datetime]) -> str: return ",".join(dt.isoformat() for dt in dttm_filter) + @staticmethod + def _serialize_dttm_filter_for_log(dttm_filter: Sequence[datetime.datetime]) -> str: + return ",".join( + timezone.coerce_datetime(dt).astimezone(settings.TIMEZONE).isoformat() for dt in dttm_filter + ) + def poke(self, context: Context) -> bool: # delay check to poke rather than __init__ in case it was supplied as XComArgs if self.external_task_ids and len(self.external_task_ids) > len(set(self.external_task_ids)): @@ -287,6 +295,7 @@ def poke(self, context: Context) -> bool: dttm_filter = self._get_dttm_filter(context) serialized_dttm_filter = self._serialize_dttm_filter(dttm_filter) + log_dttm_filter = self._serialize_dttm_filter_for_log(dttm_filter) # Save as attribute - to be used by listeners self.external_dates_filter = serialized_dttm_filter @@ -295,7 +304,7 @@ def poke(self, context: Context) -> bool: "Poking for tasks %s in dag %s on %s ... ", self.external_task_ids, self.external_dag_id, - serialized_dttm_filter, + log_dttm_filter, ) if self.external_task_group_id: @@ -303,14 +312,14 @@ def poke(self, context: Context) -> bool: "Poking for task_group '%s' in dag '%s' on %s ... ", self.external_task_group_id, self.external_dag_id, - serialized_dttm_filter, + log_dttm_filter, ) if self.external_dag_id and not self.external_task_group_id and not self.external_task_ids: self.log.info( "Poking for DAG '%s' on %s ... ", self.external_dag_id, - serialized_dttm_filter, + log_dttm_filter, ) if AIRFLOW_V_3_0_PLUS: diff --git a/providers/standard/tests/unit/standard/sensors/test_external_task_sensor.py b/providers/standard/tests/unit/standard/sensors/test_external_task_sensor.py index ef6bd33a76423..aee2bedd137e5 100644 --- a/providers/standard/tests/unit/standard/sensors/test_external_task_sensor.py +++ b/providers/standard/tests/unit/standard/sensors/test_external_task_sensor.py @@ -23,6 +23,7 @@ from datetime import time, timedelta from unittest import mock +import pendulum import pytest from sqlalchemy import select @@ -491,6 +492,13 @@ def test_external_dag_sensor_log(self, caplog, dag_maker): op.run(start_date=DEFAULT_DATE, end_date=DEFAULT_DATE, ignore_ti_state=True) assert (f"Poking for DAG 'other_dag' on {DEFAULT_DATE.isoformat()} ... ") in caplog.messages + def test_external_dag_sensor_log_uses_configured_timezone(self, monkeypatch): + monkeypatch.setattr(settings, "TIMEZONE", pendulum.timezone("Asia/Seoul")) + + dttm_filter = [pendulum.datetime(2026, 7, 6, 21, tz="UTC")] + + assert ExternalTaskSensor._serialize_dttm_filter_for_log(dttm_filter) == "2026-07-07T06:00:00+09:00" + def test_external_dag_sensor_soft_fail_as_skipped(self, dag_maker, session): with dag_maker("other_dag", default_args=self.args, end_date=DEFAULT_DATE, schedule="@once"): pass From a14b3ba1960439f5fa35dde95fb5fdd59cbacde0 Mon Sep 17 00:00:00 2001 From: srinivas chennupati Date: Tue, 14 Jul 2026 09:50:56 -0500 Subject: [PATCH 2/2] Address ExternalTaskSensor timezone log review --- .../standard/sensors/external_task.py | 12 ++++++--- .../sensors/test_external_task_sensor.py | 27 ++++++++++++++++++- 2 files changed, 35 insertions(+), 4 deletions(-) diff --git a/providers/standard/src/airflow/providers/standard/sensors/external_task.py b/providers/standard/src/airflow/providers/standard/sensors/external_task.py index 81fee9a1d5722..0e4003d8f588f 100644 --- a/providers/standard/src/airflow/providers/standard/sensors/external_task.py +++ b/providers/standard/src/airflow/providers/standard/sensors/external_task.py @@ -284,9 +284,15 @@ def _serialize_dttm_filter(dttm_filter: Sequence[datetime.datetime]) -> str: @staticmethod def _serialize_dttm_filter_for_log(dttm_filter: Sequence[datetime.datetime]) -> str: - return ",".join( - timezone.coerce_datetime(dt).astimezone(settings.TIMEZONE).isoformat() for dt in dttm_filter - ) + formatted_dates = [] + for dt in dttm_filter: + serialized_dt = dt.isoformat() + timezone_dt = timezone.coerce_datetime(dt).astimezone(settings.TIMEZONE).isoformat() + if serialized_dt == timezone_dt: + formatted_dates.append(serialized_dt) + else: + formatted_dates.append(f"{serialized_dt} (default timezone: {timezone_dt})") + return ",".join(formatted_dates) def poke(self, context: Context) -> bool: # delay check to poke rather than __init__ in case it was supplied as XComArgs diff --git a/providers/standard/tests/unit/standard/sensors/test_external_task_sensor.py b/providers/standard/tests/unit/standard/sensors/test_external_task_sensor.py index aee2bedd137e5..17dae7dcc0336 100644 --- a/providers/standard/tests/unit/standard/sensors/test_external_task_sensor.py +++ b/providers/standard/tests/unit/standard/sensors/test_external_task_sensor.py @@ -497,7 +497,10 @@ def test_external_dag_sensor_log_uses_configured_timezone(self, monkeypatch): dttm_filter = [pendulum.datetime(2026, 7, 6, 21, tz="UTC")] - assert ExternalTaskSensor._serialize_dttm_filter_for_log(dttm_filter) == "2026-07-07T06:00:00+09:00" + assert ( + ExternalTaskSensor._serialize_dttm_filter_for_log(dttm_filter) + == "2026-07-06T21:00:00+00:00 (default timezone: 2026-07-07T06:00:00+09:00)" + ) def test_external_dag_sensor_soft_fail_as_skipped(self, dag_maker, session): with dag_maker("other_dag", default_args=self.args, end_date=DEFAULT_DATE, schedule="@once"): @@ -1421,6 +1424,28 @@ def test_external_task_sensor_execution_delta(self, dag_maker): ) assert op.external_dates_filter == expected_date.isoformat() + @pytest.mark.execution_timeout(10) + def test_external_dag_sensor_log_uses_configured_timezone(self, monkeypatch, caplog, dag_maker): + monkeypatch.setattr(settings, "TIMEZONE", pendulum.timezone("Asia/Seoul")) + logical_date = pendulum.datetime(2026, 7, 6, 21, tz="UTC") + self.context["logical_date"] = logical_date + self.context["ti"].get_dr_count.return_value = 0 + + with dag_maker("test_dag_child"): + op = ExternalTaskSensor( + task_id="test_external_dag_sensor_check", + external_dag_id="other_dag", + ) + + with caplog.at_level(logging.INFO, logger=op.log.name): + caplog.clear() + op.poke(context=self.context) + + assert ( + "Poking for DAG 'other_dag' on " + "2026-07-06T21:00:00+00:00 (default timezone: 2026-07-07T06:00:00+09:00) ... " + ) in caplog.messages + @pytest.mark.execution_timeout(10) def test_external_task_sensor_duplicate_task_ids(self, dag_maker): with dag_maker("test_dag_child"):