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")],

Reply via email to