This is an automated email from the ASF dual-hosted git repository.
shahar1 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 d7872deade7 Fix duplicated logs with memory usage issue in GCS log
handler (#72322)
d7872deade7 is described below
commit d7872deade77e1b9c8db96b2010968854adad84e
Author: Ivan Nikolic <[email protected]>
AuthorDate: Wed Sep 16 07:31:55 2026 +0200
Fix duplicated logs with memory usage issue in GCS log handler (#72322)
GCSRemoteLogIO.upload() never truncates the local log file after a
successful upload. For any sensor running mode="reschedule",
try_number does not change between pokes, so every poke's close()
writes to the same remote log key. If a later poke lands on a worker
that still has the local log file from an earlier poke (long-lived
Celery/Local workers, or a shared logs volume), duplicate content
compounds every cycle: poke 1's lines get re-uploaded on every later
poke, poke 2's lines on every poke after that, and so on. For n
pokes, the total amount of duplicated content grows like n(n+1)/2.
This is the same bug class already fixed for the S3 (#67144) and WASB
(#70860) log handlers, both by truncating the local log after a
successful upload. GCS never received the equivalent fix.
Truncates the local log file after a successful upload when
delete_local_copy is False, mirroring the existing S3/WASB fix.
Updated test_upload to assert the local file is truncated after a
successful upload (when delete_local_copy is False) and left
untouched after a failed upload (so a retry can still send it).
---
.../providers/google/cloud/log/gcs_task_handler.py | 2 +
.../unit/google/cloud/log/test_gcs_task_handler.py | 61 ++++++++++++++++++++++
2 files changed, 63 insertions(+)
diff --git
a/providers/google/src/airflow/providers/google/cloud/log/gcs_task_handler.py
b/providers/google/src/airflow/providers/google/cloud/log/gcs_task_handler.py
index 26a83af4190..78c2001f231 100644
---
a/providers/google/src/airflow/providers/google/cloud/log/gcs_task_handler.py
+++
b/providers/google/src/airflow/providers/google/cloud/log/gcs_task_handler.py
@@ -113,6 +113,8 @@ class GCSRemoteLogIO(LoggingMixin): # noqa: D101
has_uploaded = self.write(log, remote_loc)
if has_uploaded and self.delete_local_copy:
shutil.rmtree(os.path.dirname(local_loc))
+ elif has_uploaded:
+ local_loc.write_text("")
@cached_property
def hook(self) -> GCSHook | None:
diff --git
a/providers/google/tests/unit/google/cloud/log/test_gcs_task_handler.py
b/providers/google/tests/unit/google/cloud/log/test_gcs_task_handler.py
index 9a226fd37a6..49f2e8251f9 100644
--- a/providers/google/tests/unit/google/cloud/log/test_gcs_task_handler.py
+++ b/providers/google/tests/unit/google/cloud/log/test_gcs_task_handler.py
@@ -209,12 +209,73 @@ class TestGCSRemoteLogIO:
mock_write_method.assert_called_once()
if delete_local_copy and mock_write_method_result:
mock_rmtree.assert_called_once_with(tmp_path.as_posix())
+ # shutil.rmtree is mocked here, so it never actually
removes the file from
+ # disk; the assertion above already confirms the delete
path was taken. We
+ # don't assert on-disk deletion in this branch since it
would be asserting
+ # against a mock's real side effects, which don't exist.
+ elif mock_write_method_result:
+ # Upload succeeded but delete_local_copy is False: the
local file must be
+ # truncated so a later lifecycle (e.g. the next poke of a
reschedule-mode
+ # sensor landing on the same worker) doesn't re-upload
already-stored content.
+ mock_rmtree.assert_not_called()
+ assert (tmp_path / "existing.log").read_text() == ""
else:
+ # Upload failed: keep the local content untouched so a
retry can still send it.
mock_rmtree.assert_not_called()
+ assert (tmp_path / "existing.log").read_text() == "log
content"
else:
mock_write_method.assert_not_called()
mock_rmtree.assert_not_called()
+ def test_upload_repeated_cycles_no_duplication(self, mock_creds, tmp_path:
Path):
+ """Simulate reschedule-mode sensor: each cycle appends to the local
log, then uploads.
+
+ Without truncation after upload, the GCS object accumulates duplicate
lines and
+ grows O(N^2). The correct behavior is that each line appears in GCS
exactly once.
+ """
+ remote_store: dict[str, str] = {}
+
+ class FakeBlob:
+ def __init__(self, remote_log_location):
+ self.remote_log_location = remote_log_location
+
+ def download_as_bytes(self):
+ if self.remote_log_location not in remote_store:
+ raise Exception("No such object: fake-bucket")
+ return remote_store[self.remote_log_location].encode()
+
+ def upload_from_string(self, content, content_type=None):
+ remote_store[self.remote_log_location] = content
+
+ gcs_remote_log_io = GCSRemoteLogIO(
+ remote_base=self.gcs_log_folder,
+ base_log_folder=tmp_path.as_posix(),
+ delete_local_copy=False,
+ )
+ local_log = tmp_path / "1.log"
+
+ with (
+ mock.patch("google.cloud.storage.Client"),
+ mock.patch("google.cloud.storage.Blob") as mock_blob,
+ ):
+ mock_blob.from_string.side_effect = lambda remote_log_location,
client: FakeBlob(
+ remote_log_location
+ )
+
+ for cycle in range(1, 4):
+ with open(local_log, "a") as f:
+ f.write(f"cycle {cycle}\n")
+ gcs_remote_log_io.upload(local_log, self.ti)
+
+ # GCSRemoteLogIO.write() joins old and new content with an
unconditional "\n" separator
+ # (pre-existing behavior, not part of this fix), so each upload after
the first adds a
+ # blank line. What matters here is that each cycle's line appears
exactly once.
+ final_content = next(iter(remote_store.values()))
+ assert final_content == "cycle 1\n\ncycle 2\n\ncycle 3\n"
+ for cycle in range(1, 4):
+ assert final_content.count(f"cycle {cycle}") == 1
+ assert local_log.read_text() == ""
+
@pytest.mark.parametrize(
"upload_success",
[pytest.param(True, id="upload-success"), pytest.param(False,
id="upload-fail")],