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(