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."""