Repository navigation
Read the AWS task logs once more after the fetcher is stopped - #73211
Merged
Merged
Conversation
`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: apache#73210
2 tasks done
o-nikolas
approved these changes
Sep 24, 2026
Wait on the stop event instead of sleeping for fetch_interval, so stop() ends the wait at once and the final read runs without the extra delay.
|
Awesome work, congrats on your first merged pull request! You are invited to check our Issue Tracker for additional contributions. |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
AwsTaskLogFetcher.runchecks the stop flag at the top of its loop, then sleeps, then fetches. Whenstop()arrives while the thread is inside a fetch, or in the gap between a fetch finishing and the next check, the loop ends with no read after the stop, and the events the container wrote at the end of the task are never forwarded to the task log.Both callers stop the fetcher as soon as the task or job has ended, which is exactly when those events appear:
EcsRunTaskOperator.execute(finally: self.task_log_fetcher.stop()) andBatchClientHook.wait_for_job(finally: batch_log_fetcher.stop()).This forwards the events once more after leaving the loop. The continuation token keeps the extra read from repeating events already seen, and
AwsLogsHook.get_log_eventsreturns as soon as the stream is exhausted, so the cost is oneget_log_eventscall at thread exit. The body of the loop moved into_forward_log_eventsunchanged.test_run_forwards_the_events_written_before_it_was_stoppedcovers it: it fails on the current code (one fetch, the last event never logged) and passes here. The two existingruntests each get one more empty page in theirside_effect, sincerunnow always performs that final fetch.Verified locally with
apache-airflow 3.1.8and the provider installed from this branch:tests/unit/amazon/aws/operators/test_ecs.py,test_batch.pyandhooks/test_batch_client.pypatchAwsTaskLogFetcheras a whole and never exerciserun, so they are unaffected; CI runs them.closes: #73210
Was generative AI tooling used to co-author this PR?
Generated-by: Claude Code following the guidelines
The behaviour was reproduced against the released provider before the change, the fix and its test were reviewed line by line, and the tests and static checks above were run locally.