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

vincbeck 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 747f69b1a8c Warn when db clean skips only some rows of a batch (#73568)
747f69b1a8c is described below

commit 747f69b1a8c1c715bb85df8aaf22b6e1667c9cdf
Author: D. Ferruzzi <[email protected]>
AuthorDate: Mon Oct 5 06:55:50 2026 -0700

    Warn when db clean skips only some rows of a batch (#73568)
    
    The guarded DELETE added in #66350 logged a warning only when ``deleted == 
0``, which fires when an entire batch is skipped by the skip_if_referenced 
guard.  The partial case is both the likely one and the harmful one: the 
archive has already committed a copy of every row the SELECT found, so a row 
the guard skips is archived and still live, and export-archived will emit a 
phantom copy of it.
---
 airflow-core/src/airflow/utils/db_cleanup.py     | 30 ++++++++----
 airflow-core/tests/unit/utils/test_db_cleanup.py | 62 ++++++++++++++++++------
 2 files changed, 67 insertions(+), 25 deletions(-)

diff --git a/airflow-core/src/airflow/utils/db_cleanup.py 
b/airflow-core/src/airflow/utils/db_cleanup.py
index ad41b482c44..bc09ba4ec8b 100644
--- a/airflow-core/src/airflow/utils/db_cleanup.py
+++ b/airflow-core/src/airflow/utils/db_cleanup.py
@@ -483,15 +483,27 @@ def _do_delete(
             # A guarded DELETE (skip_if_referenced) may delete fewer rows than 
the SELECT
             # found. The SELECT includes the same NOT EXISTS guard, so the 
skipped row is
             # excluded on the next pass too and the loop drains naturally. 
With --batch-size
-            # set, continuing lets subsequent batches clean rows unaffected by 
the race.
-            if deleted == 0:
-                logger.warning(
-                    "Some rows from %s are still referenced by another table 
and were not "
-                    "deleted; they remain in %s and will be retried on the 
next cleanup run.",
-                    source_table_name,
-                    target_table_name if not skip_archive else "the archive 
(which is being dropped)",
-                )
-                continue
+            # set, later batches still clean rows unaffected by the race.
+            #
+            # Compare against the archive rather than testing ``deleted == 
0``: the archive
+            # holds exactly the rows this pass found, so any shortfall is a 
skipped row. A
+            # partial skip is the likely case and is also the harmful one, 
because the
+            # archive has already committed a copy of a row that is still live.
+            if skip_if_referenced:
+                archived = 
session.scalars(select(func.count()).select_from(target_table)).one()
+                if deleted < archived:
+                    logger.warning(
+                        "%s of %s rows from %s are still referenced by another 
table and were "
+                        "not deleted; they remain live in %s and will be 
retried on the next "
+                        "cleanup run.%s",
+                        archived - deleted,
+                        archived,
+                        source_table_name,
+                        source_table_name,
+                        ""
+                        if skip_archive
+                        else f" {target_table_name} already holds an archived 
copy of them.",
+                    )
 
         except BaseException:
             error_raised = True
diff --git a/airflow-core/tests/unit/utils/test_db_cleanup.py 
b/airflow-core/tests/unit/utils/test_db_cleanup.py
index 375461b6848..e5a1d74730b 100644
--- a/airflow-core/tests/unit/utils/test_db_cleanup.py
+++ b/airflow-core/tests/unit/utils/test_db_cleanup.py
@@ -17,6 +17,7 @@
 # under the License.
 from __future__ import annotations
 
+import re
 import threading
 import time
 from contextlib import suppress
@@ -650,7 +651,16 @@ class TestDBCleanup:
         assert latest_id in remaining  # kept by keep_last
         assert orphan_id not in remaining  # old and unreferenced -> pruned
 
-    def test_do_delete_skip_if_referenced_guards_against_race(self):
+    @pytest.mark.parametrize(
+        ("extra_unreferenced", "expected_count"),
+        [
+            pytest.param(0, "1 of 1", id="whole-batch-skipped"),
+            pytest.param(1, "1 of 2", id="partial-skip"),
+        ],
+    )
+    def test_do_delete_skip_if_referenced_guards_against_race(
+        self, extra_unreferenced, expected_count, cap_structlog
+    ):
         """_do_delete must not issue a DELETE that violates an ON DELETE 
RESTRICT FK.
 
         Reproduces the real race: the dag_version row passes the SELECT filter 
and is
@@ -658,6 +668,12 @@ class TestDBCleanup:
         skip_if_referenced guard on the DELETE must skip the row instead of 
failing with
         IntegrityError, and the loop must still drain because the next SELECT 
pass
         re-evaluates the same NOT EXISTS guard and excludes it.
+
+        Parametrized over how many *unreferenced* rows share the pass, because 
the guard
+        skipping part of a batch is not the same as it skipping all of it.  
The archive has
+        already committed a copy of every row the SELECT found, so a skipped 
row is both
+        archived and still live and must be reported either way.  The partial 
case is the
+        one a ``deleted == 0`` condition stays silent for.
         """
         from airflow.utils.db import reflect_tables
 
@@ -671,26 +687,33 @@ class TestDBCleanup:
             session.add(DagModel(dag_id=dag_id, bundle_name=bundle_name))
             session.flush()
 
-            raced_old = DagVersion(
-                dag_id=dag_id,
-                version_number=1,
-                bundle_name=bundle_name,
-                created_at=base_date,
-                last_updated=base_date,
-            )
-            # dag_version is keep_last per dag_id, so a lone version is always 
the
-            # keep_last survivor and is never eligible for deletion.  A 
second, newer
-            # version takes that role and leaves raced_old as the deletion 
candidate.
+            # The first is the row the race pins mid-pass; any extras are 
unreferenced and
+            # must still be deleted in the same pass, which is what makes the 
skip partial.
+            eligible = [
+                DagVersion(
+                    dag_id=dag_id,
+                    version_number=n + 1,
+                    bundle_name=bundle_name,
+                    created_at=base_date.add(minutes=n),
+                    last_updated=base_date.add(minutes=n),
+                )
+                for n in range(1 + extra_unreferenced)
+            ]
+            # dag_version is keep_last per dag_id, so the newest version is 
always the
+            # keep_last survivor and is never eligible for deletion.  It is 
what leaves
+            # the versions above as deletion candidates.
             latest = DagVersion(
                 dag_id=dag_id,
-                version_number=2,
+                version_number=len(eligible) + 1,
                 bundle_name=bundle_name,
-                created_at=base_date.add(minutes=1),
-                last_updated=base_date.add(minutes=1),
+                created_at=base_date.add(minutes=len(eligible)),
+                last_updated=base_date.add(minutes=len(eligible)),
             )
-            session.add_all([raced_old, latest])
+            session.add_all([*eligible, latest])
             session.flush()
-            raced_old_id, latest_id = raced_old.id, latest.id
+            raced_old_id = eligible[0].id
+            unreferenced_ids = {version.id for version in eligible[1:]}
+            latest_id = latest.id
 
             # Query built while nothing references raced_old, so the first 
SELECT pass
             # returns it and _do_delete archives it.
@@ -742,6 +765,13 @@ class TestDBCleanup:
         assert raced, "the TI was never inserted mid-pass; the race was not 
reproduced"
         assert raced_old_id in remaining, "dag_version referenced by a 
task_instance must not be deleted"
         assert latest_id in remaining, "the keep_last survivor must not be 
deleted"
+        assert not unreferenced_ids & remaining, "unreferenced dag_versions 
must still be deleted"
+        # The count is what distinguishes a partial skip from a whole-batch 
one, and the
+        # partial case is what ``deleted == 0`` reported nothing for.
+        assert {
+            "event": re.compile(rf"{expected_count} rows from dag_version are 
still referenced"),
+            "log_level": "warning",
+        } in cap_structlog, cap_structlog.entries
 
     def test_table_config_skip_if_referenced_requires_pk_column(self):
         """A misconfigured skip_if_referenced (pk not in columns) must fail 
fast at construction."""

Reply via email to