This is an automated email from the ASF dual-hosted git repository.

o-nikolas pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/airflow.git


The following commit(s) were added to refs/heads/main by this push:
     new 10221c300eb Read the AWS task logs once more after the fetcher is 
stopped (#73211)
10221c300eb is described below

commit 10221c300eb9266795863846309737cf48164ff9
Author: Filippo Scotti <[email protected]>
AuthorDate: Tue Sep 29 00:13:37 2026 +0200

    Read the AWS task logs once more after the fetcher is stopped (#73211)
    
    `AwsTaskLogFetcher.run` checks the stop flag at the top of its loop, so a
    `stop()` that arrives while the thread is inside a fetch, or between a fetch
    and the next check, ends the loop with no read after it. The events the
    container wrote at the end of the task are then never forwarded to the task
    log. `EcsRunTaskOperator.execute` and `BatchClientHook.wait_for_job` both
    stop the fetcher as soon as the task or job has ended, which is exactly when
    those last events appear.
    
    Forward the events once more after leaving the loop. The continuation token
    keeps that read from repeating events already seen.
    
    closes: #73210
---
 .../providers/amazon/aws/utils/task_log_fetcher.py | 38 ++++++++++++----------
 .../unit/amazon/aws/utils/test_task_log_fetcher.py | 36 ++++++++++++++++++--
 2 files changed, 54 insertions(+), 20 deletions(-)

diff --git 
a/providers/amazon/src/airflow/providers/amazon/aws/utils/task_log_fetcher.py 
b/providers/amazon/src/airflow/providers/amazon/aws/utils/task_log_fetcher.py
index dc0c35aedb2..7b4f8563cc7 100644
--- 
a/providers/amazon/src/airflow/providers/amazon/aws/utils/task_log_fetcher.py
+++ 
b/providers/amazon/src/airflow/providers/amazon/aws/utils/task_log_fetcher.py
@@ -100,24 +100,26 @@ class AwsTaskLogFetcher(Thread):
 
     def run(self) -> None:
         continuation_token = AwsLogsHook.ContinuationToken()
-        while not self.is_stopped():
-            time.sleep(self.fetch_interval.total_seconds())
-            log_events = self._get_log_events(continuation_token)
-            prev_timestamp_event = None
-            for log_event in log_events:
-                current_timestamp_event = datetime.fromtimestamp(
-                    log_event["timestamp"] / 1000.0, tz=timezone.utc
-                )
-                if current_timestamp_event == prev_timestamp_event:
-                    # When multiple events have the same timestamp, somehow, 
only one event is logged
-                    # As a consequence, some logs are missed in the log group 
(in case they have the same
-                    # timestamp)
-                    # When a slight delay is added before logging the event, 
that solves the issue
-                    # See https://github.com/apache/airflow/issues/40875
-                    time.sleep(0.001)
-                level = _parse_log_level(log_event["message"])
-                self.logger.log(level, self.event_to_str(log_event))
-                prev_timestamp_event = current_timestamp_event
+        while not self._event.wait(self.fetch_interval.total_seconds()):
+            self._forward_log_events(continuation_token)
+        # `stop()` is called once the task or job has ended, so the events 
written between the last
+        # fetch above and that moment have not been read yet: read them before 
the thread exits.
+        self._forward_log_events(continuation_token)
+
+    def _forward_log_events(self, continuation_token: 
AwsLogsHook.ContinuationToken) -> None:
+        prev_timestamp_event = None
+        for log_event in self._get_log_events(continuation_token):
+            current_timestamp_event = 
datetime.fromtimestamp(log_event["timestamp"] / 1000.0, tz=timezone.utc)
+            if current_timestamp_event == prev_timestamp_event:
+                # When multiple events have the same timestamp, somehow, only 
one event is logged
+                # As a consequence, some logs are missed in the log group (in 
case they have the same
+                # timestamp)
+                # When a slight delay is added before logging the event, that 
solves the issue
+                # See https://github.com/apache/airflow/issues/40875
+                time.sleep(0.001)
+            level = _parse_log_level(log_event["message"])
+            self.logger.log(level, self.event_to_str(log_event))
+            prev_timestamp_event = current_timestamp_event
 
     def _get_log_events(self, skip_token: AwsLogsHook.ContinuationToken | None 
= None) -> Generator:
         if skip_token is None:
diff --git 
a/providers/amazon/tests/unit/amazon/aws/utils/test_task_log_fetcher.py 
b/providers/amazon/tests/unit/amazon/aws/utils/test_task_log_fetcher.py
index e1a18faf8d3..6125d0bc3cf 100644
--- a/providers/amazon/tests/unit/amazon/aws/utils/test_task_log_fetcher.py
+++ b/providers/amazon/tests/unit/amazon/aws/utils/test_task_log_fetcher.py
@@ -59,10 +59,11 @@ class TestAwsTaskLogFetcher:
                 ]
             ),
             iter([]),
+            iter([]),
         ),
     )
     def test_run(self, get_log_events_mock):
-        with mock.patch.object(self.log_fetcher._event, "is_set", 
side_effect=(False, False, False, True)):
+        with mock.patch.object(self.log_fetcher._event, "wait", 
side_effect=(False, False, False, True)):
             self.log_fetcher.run()
 
         self.logger_mock.log.assert_has_calls(
@@ -73,6 +74,36 @@ class TestAwsTaskLogFetcher:
             ]
         )
 
+    @mock.patch(
+        "airflow.providers.amazon.aws.hooks.logs.AwsLogsHook.get_log_events",
+        side_effect=(
+            iter([{"timestamp": 1617400267123, "message": "First"}]),
+            iter([{"timestamp": 1617400467789, "message": "Written just before 
the task ended"}]),
+        ),
+    )
+    def test_run_forwards_the_events_written_before_it_was_stopped(self, 
get_log_events_mock):
+        """The callers stop the fetcher once the task ended, after the last 
events were written."""
+        with mock.patch.object(self.log_fetcher._event, "wait", 
side_effect=(False, True)):
+            self.log_fetcher.run()
+
+        assert get_log_events_mock.call_count == 2
+        self.logger_mock.log.assert_has_calls(
+            [
+                mock.call(logging.INFO, "[2021-04-02 21:51:07,123] First"),
+                mock.call(logging.INFO, "[2021-04-02 21:54:27,789] Written 
just before the task ended"),
+            ]
+        )
+
+    
@mock.patch("airflow.providers.amazon.aws.hooks.logs.AwsLogsHook.get_log_events",
 return_value=iter([]))
+    def test_stop_ends_the_wait_for_the_next_fetch(self, get_log_events_mock):
+        self.log_fetcher.fetch_interval = timedelta(hours=1)
+        self.log_fetcher.daemon = True
+        self.log_fetcher.start()
+        self.log_fetcher.stop()
+        self.log_fetcher.join(timeout=5)
+
+        assert not self.log_fetcher.is_alive()
+
     @mock.patch(
         "airflow.providers.amazon.aws.hooks.logs.AwsLogsHook.get_log_events",
         side_effect=ClientError({"Error": {"Code": 
"ResourceNotFoundException"}}, None),
@@ -162,10 +193,11 @@ class TestAwsTaskLogFetcher:
                     },
                 ]
             ),
+            iter([]),
         ),
     )
     def test_run_with_log_level_detection(self, get_log_events_mock):
-        with mock.patch.object(self.log_fetcher._event, "is_set", 
side_effect=(False, True)):
+        with mock.patch.object(self.log_fetcher._event, "wait", 
side_effect=(False, True)):
             self.log_fetcher.run()
 
         self.logger_mock.log.assert_has_calls(

Reply via email to