This is an automated email from the ASF dual-hosted git repository.
ephraimbuddy 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 d6cd4595dbd Fix airflow db clean hang on MySQL when delete fails
(#66296)
d6cd4595dbd is described below
commit d6cd4595dbd55f53d0e0d99a48515c5f0544a610
Author: deepinsight coder <[email protected]>
AuthorDate: Tue Jul 14 01:53:55 2026 -0700
Fix airflow db clean hang on MySQL when delete fails (#66296)
* Fix airflow db clean hang on MySQL when delete fails (#66177)
* Rename newsfragment to match PR number
The check-newsfragment-pr-number CI job requires the newsfragment
filename to match the PR number, not the linked issue number.
* Remove newsfragment per review
* Suppress rollback errors so the original delete error propagates
* Explain why the archive table drop reuses the session connection
* Force the FK-failure path in backend regression tests via config copy
Keeps the hang-regression tests meaningful once dag_version cleanup learns
to skip FK-pinned rows (apache/airflow#68339).
* Guard archive-table cleanup so it can't mask the original delete error
When _do_delete unwinds from a failed DELETE, dropping the archive table in
the finally block could itself raise (e.g. a dead connection on MySQL).
Python
makes a finally-raised error the top-level exception, so that cleanup error
would replace the original delete failure -- and an OperationalError from
the
drop could then be downgraded to a warning by _suppress_with_logging,
silently
masking the real IntegrityError. Guard the cleanup so it only propagates on
the
success path; on the failure path it is logged and the original error
survives.
* Update airflow-core/tests/unit/utils/test_db_cleanup.py
Co-authored-by: Ephraim Anierobi <[email protected]>
* Retry CI after provider compat startup failure
The previous run failed because the Mongo provider test fixture could not
connect to its temporary mongod container. No source change is needed; this
commit retriggers CI so the PR can use a fresh signal.
---------
Co-authored-by: Ephraim Anierobi <[email protected]>
Co-authored-by: Ramachandra Nalam <[email protected]>
---
airflow-core/src/airflow/utils/db_cleanup.py | 44 +++-
airflow-core/tests/unit/utils/test_db_cleanup.py | 318 ++++++++++++++++++++++-
2 files changed, 353 insertions(+), 9 deletions(-)
diff --git a/airflow-core/src/airflow/utils/db_cleanup.py
b/airflow-core/src/airflow/utils/db_cleanup.py
index 8197d56c42d..6deb4e78f68 100644
--- a/airflow-core/src/airflow/utils/db_cleanup.py
+++ b/airflow-core/src/airflow/utils/db_cleanup.py
@@ -27,7 +27,7 @@ import csv
import logging
import os
from collections.abc import Generator
-from contextlib import contextmanager
+from contextlib import contextmanager, suppress
from dataclasses import dataclass
from types import SimpleNamespace
from typing import TYPE_CHECKING, Any
@@ -286,6 +286,10 @@ def _do_delete(
target_table_name =
f"{ARCHIVE_TABLE_PREFIX}{orm_model.name}__{timestamp_str}{suffix}"
print(f"Moving data to table {target_table_name}")
target_table = None
+ # Lets the ``finally`` cleanup below tell the failure path (don't let a
+ # cleanup error mask the original) from the success path (a cleanup
error
+ # is a real problem and must propagate).
+ error_raised = False
try:
if dialect_name == "mysql":
@@ -323,13 +327,41 @@ def _do_delete(
session.execute(delete)
session.commit()
- except BaseException as e:
- raise e
+ except BaseException:
+ error_raised = True
+ # Roll back the failed transaction so its locks are released before
+ # the archive table is dropped in the ``finally`` block below.
+ # ``rollback()`` itself can raise (e.g. the connection died);
suppress
+ # it so it does not shadow the original error being re-raised.
+ with suppress(Exception):
+ session.rollback()
+ raise
finally:
if target_table is not None and skip_archive:
- bind = session.get_bind()
- target_table.drop(bind=bind)
- session.commit()
+ # Drop the archive table on the session's own connection.
Binding
+ # the drop to ``session.get_bind()`` (the Engine) would check
out a
+ # *second* pooled connection, and on MySQL its ``DROP TABLE``
blocks
+ # indefinitely on the metadata lock still held by this
session's
+ # open transaction when the DELETE above failed -- the ``db
clean``
+ # hang reported in #66177.
+ try:
+ target_table.drop(bind=session.connection())
+ session.commit()
+ except Exception:
+ # If we are already unwinding from a delete failure, a
cleanup
+ # error here must not replace the original exception
(Python
+ # makes a ``finally``-raised error the top-level one). Log
and
+ # let the original delete error keep propagating. On the
success
+ # path (no delete error), a drop/commit failure is a real
+ # problem, so re-raise it.
+ if not error_raised:
+ raise
+ logger.warning(
+ "Failed to drop archive table %s while cleaning up
after a "
+ "delete failure; propagating the original delete error
instead.",
+ target_table_name,
+ exc_info=True,
+ )
print("Finished Performing Delete")
diff --git a/airflow-core/tests/unit/utils/test_db_cleanup.py
b/airflow-core/tests/unit/utils/test_db_cleanup.py
index 4d66a96b1a3..8271d38968b 100644
--- a/airflow-core/tests/unit/utils/test_db_cleanup.py
+++ b/airflow-core/tests/unit/utils/test_db_cleanup.py
@@ -17,18 +17,21 @@
# under the License.
from __future__ import annotations
+import threading
+import time
from contextlib import suppress
from importlib import import_module
from io import StringIO
from pathlib import Path
-from unittest.mock import MagicMock, mock_open, patch
+from unittest.mock import MagicMock, call, mock_open, patch
from uuid import uuid4
import pendulum
import pytest
-from sqlalchemy import func, inspect, select, text
-from sqlalchemy.exc import OperationalError, SQLAlchemyError
+from sqlalchemy import Column, Integer, MetaData, Table, func, inspect,
literal, select, text
+from sqlalchemy.exc import IntegrityError, OperationalError, SQLAlchemyError
from sqlalchemy.ext.declarative import DeclarativeMeta
+from sqlalchemy.orm import Session
from airflow import DAG
from airflow._shared.timezones import timezone
@@ -46,6 +49,7 @@ from airflow.utils.db_cleanup import (
_build_query,
_cleanup_table,
_confirm_drop_archives,
+ _do_delete,
_dump_table_to_file,
_get_archived_table_names,
_TableConfig,
@@ -558,6 +562,167 @@ class TestDBCleanup:
# "id" intentionally omitted from extra_columns
)
+ def test_do_delete_rolls_back_before_drop_on_failure(self):
+ session = MagicMock(spec=Session)
+ session.get_bind.return_value.dialect.name = "mysql"
+ session.connection.return_value = object()
+ session.scalars.return_value.one.side_effect = [1, 0]
+ tracker = MagicMock()
+ tracker.attach_mock(session, "session")
+
+ metadata, source_table, target_table, query =
_build_do_delete_test_objects()
+ delete_failure = IntegrityError("DELETE FROM dag_version", {},
Exception("fk violation"))
+ session.execute.side_effect = [None, None, delete_failure]
+
+ with (
+ patch("airflow.utils.db_cleanup.reflect_tables",
return_value=metadata),
+ patch("airflow.utils.db_cleanup.timezone.utcnow",
return_value=_delete_test_timestamp()),
+ patch.object(target_table, "drop") as drop_mock,
+ ):
+ tracker.attach_mock(drop_mock, "drop")
+ with pytest.raises(IntegrityError) as exc_info:
+ _do_delete(
+ query=query,
+ orm_model=source_table,
+ skip_archive=True,
+ session=session,
+ batch_size=None,
+ )
+
+ assert exc_info.value is delete_failure
+ session.rollback.assert_called_once_with()
+ session.connection.assert_called_once_with()
+ assert session.get_bind.call_count == 1
+ drop_mock.assert_called_once_with(bind=session.connection.return_value)
+
+ rollback_call_index = tracker.mock_calls.index(call.session.rollback())
+ drop_call_index =
tracker.mock_calls.index(call.drop(bind=session.connection.return_value))
+ commit_call_index = tracker.mock_calls.index(call.session.commit(),
drop_call_index)
+ assert rollback_call_index < drop_call_index < commit_call_index
+
+ def test_do_delete_propagates_original_error_when_rollback_fails(self):
+ session = MagicMock(spec=Session)
+ session.get_bind.return_value.dialect.name = "mysql"
+ session.connection.return_value = object()
+ session.scalars.return_value.one.side_effect = [1, 0]
+
+ metadata, source_table, target_table, query =
_build_do_delete_test_objects()
+ delete_failure = IntegrityError("DELETE FROM dag_version", {},
Exception("fk violation"))
+ session.execute.side_effect = [None, None, delete_failure]
+ session.rollback.side_effect = OperationalError("ROLLBACK", {},
Exception("connection lost"))
+
+ with (
+ patch("airflow.utils.db_cleanup.reflect_tables",
return_value=metadata),
+ patch("airflow.utils.db_cleanup.timezone.utcnow",
return_value=_delete_test_timestamp()),
+ patch.object(target_table, "drop") as drop_mock,
+ ):
+ with pytest.raises(IntegrityError) as exc_info:
+ _do_delete(
+ query=query,
+ orm_model=source_table,
+ skip_archive=True,
+ session=session,
+ batch_size=None,
+ )
+
+ assert exc_info.value is delete_failure
+ session.rollback.assert_called_once_with()
+ drop_mock.assert_called_once_with(bind=session.connection.return_value)
+
+ @pytest.mark.parametrize(
+ ("skip_archive", "expected_commit_count"),
+ [pytest.param(True, 3, id="skip_archive"), pytest.param(False, 2,
id="keep_archive")],
+ )
+ def test_do_delete_success_does_not_call_rollback(self, skip_archive,
expected_commit_count):
+ session = MagicMock(spec=Session)
+ session.get_bind.return_value.dialect.name = "mysql"
+ session.connection.return_value = object()
+ session.scalars.return_value.one.side_effect = [1, 0]
+ session.execute.side_effect = [None, None, None]
+
+ metadata, source_table, target_table, query =
_build_do_delete_test_objects()
+
+ with (
+ patch("airflow.utils.db_cleanup.reflect_tables",
return_value=metadata),
+ patch("airflow.utils.db_cleanup.timezone.utcnow",
return_value=_delete_test_timestamp()),
+ patch.object(target_table, "drop") as drop_mock,
+ ):
+ _do_delete(
+ query=query,
+ orm_model=source_table,
+ skip_archive=skip_archive,
+ session=session,
+ batch_size=None,
+ )
+
+ session.rollback.assert_not_called()
+ assert session.commit.call_count == expected_commit_count
+ if skip_archive:
+ session.connection.assert_called_once_with()
+
drop_mock.assert_called_once_with(bind=session.connection.return_value)
+ else:
+ session.connection.assert_not_called()
+ drop_mock.assert_not_called()
+
+ def test_do_delete_original_error_survives_archive_drop_failure(self):
+ """On the failure path, a drop/commit error in the finally block must
not
+ replace the original delete error (nailo2c review, #66296)."""
+ session = MagicMock(spec=Session)
+ session.get_bind.return_value.dialect.name = "mysql"
+ session.connection.return_value = object()
+ session.scalars.return_value.one.side_effect = [1, 0]
+
+ metadata, source_table, target_table, query =
_build_do_delete_test_objects()
+ delete_failure = IntegrityError("DELETE FROM dag_version", {},
Exception("fk violation"))
+ session.execute.side_effect = [None, None, delete_failure]
+ drop_failure = OperationalError("DROP TABLE", {}, Exception("server
has gone away"))
+
+ with (
+ patch("airflow.utils.db_cleanup.reflect_tables",
return_value=metadata),
+ patch("airflow.utils.db_cleanup.timezone.utcnow",
return_value=_delete_test_timestamp()),
+ patch.object(target_table, "drop", side_effect=drop_failure) as
drop_mock,
+ ):
+ with pytest.raises(IntegrityError) as exc_info:
+ _do_delete(
+ query=query,
+ orm_model=source_table,
+ skip_archive=True,
+ session=session,
+ batch_size=None,
+ )
+
+ assert exc_info.value is delete_failure
+ drop_mock.assert_called_once_with(bind=session.connection.return_value)
+
+ def test_do_delete_success_propagates_archive_drop_error(self):
+ """On the success path, a drop/commit failure is a real error and must
+ still surface (the failure-path guard must not swallow it)."""
+ session = MagicMock(spec=Session)
+ session.get_bind.return_value.dialect.name = "mysql"
+ session.connection.return_value = object()
+ session.scalars.return_value.one.side_effect = [1, 0]
+ session.execute.side_effect = [None, None, None]
+
+ metadata, source_table, target_table, query =
_build_do_delete_test_objects()
+ drop_failure = OperationalError("DROP TABLE", {}, Exception("disk
full"))
+
+ with (
+ patch("airflow.utils.db_cleanup.reflect_tables",
return_value=metadata),
+ patch("airflow.utils.db_cleanup.timezone.utcnow",
return_value=_delete_test_timestamp()),
+ patch.object(target_table, "drop", side_effect=drop_failure),
+ ):
+ with pytest.raises(OperationalError) as exc_info:
+ _do_delete(
+ query=query,
+ orm_model=source_table,
+ skip_archive=True,
+ session=session,
+ batch_size=None,
+ )
+
+ assert exc_info.value is drop_failure
+ session.rollback.assert_not_called()
+
@patch("airflow.utils.db.reflect_tables")
def test_skip_archive_failure_will_remove_table(self, reflect_tables_mock):
"""
@@ -591,6 +756,67 @@ class TestDBCleanup:
archived_table_names = _get_archived_table_names(["dag_run"], session)
assert len(archived_table_names) == 0
+ @pytest.mark.backend("mysql")
+ def test_db_clean_does_not_deadlock_on_fk_violation_mysql(self):
+ clean_before_date = _create_pinned_dag_version_cleanup_data(
+ base_date=pendulum.DateTime(2022, 1, 1,
tzinfo=pendulum.timezone("UTC"))
+ )
+ timeout_seconds = 15
+ finished = threading.Event()
+ result: dict[str, BaseException] = {}
+
+ def cleanup_in_thread():
+ try:
+ with create_session() as session:
+ _cleanup_table(
+ **_dag_version_config_without_row_exclusion(),
+ clean_before_timestamp=clean_before_date,
+ dry_run=False,
+ session=session,
+ table_names=["dag_version"],
+ skip_archive=True,
+ )
+ except BaseException as exc:
+ result["exception"] = exc
+ finally:
+ finished.set()
+
+ worker = threading.Thread(target=cleanup_in_thread, daemon=True)
+ start_time = time.monotonic()
+ worker.start()
+
+ assert finished.wait(timeout_seconds), (
+ "dag_version cleanup timed out after "
+ f"{time.monotonic() - start_time:.2f} seconds waiting for the FK
violation to surface"
+ )
+
+ worker.join(timeout=1)
+ assert not worker.is_alive()
+ assert isinstance(result.get("exception"), IntegrityError)
+
+ with create_session() as session:
+ assert _get_dag_version_archive_table_names(session=session) == []
+
+ @pytest.mark.backend("postgres")
+ def test_db_clean_failure_path_does_not_break_postgres(self):
+ clean_before_date = _create_pinned_dag_version_cleanup_data(
+ base_date=pendulum.DateTime(2022, 1, 1,
tzinfo=pendulum.timezone("UTC"))
+ )
+
+ with pytest.raises(IntegrityError):
+ with create_session() as session:
+ _cleanup_table(
+ **_dag_version_config_without_row_exclusion(),
+ clean_before_timestamp=clean_before_date,
+ dry_run=False,
+ session=session,
+ table_names=["dag_version"],
+ skip_archive=True,
+ )
+
+ with create_session() as session:
+ assert _get_dag_version_archive_table_names(session=session) == []
+
def test_no_models_missing(self):
"""
1. Verify that for all tables in `airflow.models`, we either have them
enabled in db cleanup,
@@ -1115,3 +1341,89 @@ class TestConnectionTestRequestCleanup:
assert seeded[state] in survivors, f"{state} row should NOT be
cleaned up"
for state in ("success", "failed"):
assert seeded[state] not in survivors, f"{state} row should be
cleaned up"
+
+
+def _delete_test_timestamp():
+ return pendulum.DateTime(2024, 1, 2, 3, 4, 5,
tzinfo=pendulum.timezone("UTC"))
+
+
+def _build_do_delete_test_objects():
+ metadata = MetaData()
+ source_table = Table("dag_version", metadata, Column("id", Integer,
primary_key=True))
+ target_table = Table(
+ f"{ARCHIVE_TABLE_PREFIX}{source_table.name}__20240102030405",
+ metadata,
+ Column("id", Integer, primary_key=True),
+ )
+ query = select(literal(1).label("id"))
+ return metadata, source_table, target_table, query
+
+
+def _create_pinned_dag_version_cleanup_data(*, base_date):
+ bundle_name = f"testing-{uuid4()}"
+ dag_id = f"test-dag_{uuid4()}"
+
+ with create_session() as session:
+ session.add(DagBundleModel(name=bundle_name))
+ session.flush()
+
+ dag = DAG(dag_id=dag_id)
+ session.add(DagModel(dag_id=dag_id, bundle_name=bundle_name))
+
+ old_dag_version = DagVersion(
+ dag_id=dag_id,
+ version_number=1,
+ bundle_name=bundle_name,
+ created_at=base_date,
+ last_updated=base_date,
+ )
+ new_dag_version = DagVersion(
+ dag_id=dag_id,
+ version_number=2,
+ bundle_name=bundle_name,
+ created_at=base_date.add(minutes=1),
+ last_updated=base_date.add(minutes=1),
+ )
+ session.add_all([old_dag_version, new_dag_version])
+ session.flush()
+
+ dag_run = DagRun(
+ dag_id,
+ run_id="run-1",
+ run_type=DagRunType.MANUAL,
+ start_date=base_date,
+ )
+ task = PythonOperator(task_id="dummy-task", dag=dag,
python_callable=print)
+ ti = create_task_instance(task, run_id=dag_run.run_id,
dag_version_id=old_dag_version.id)
+ ti.dag_id = dag_id
+ ti.start_date = base_date
+
+ session.add_all([dag_run, ti])
+ session.commit()
+
+ return base_date.add(days=1)
+
+
+def _get_dag_version_archive_table_names(*, session):
+ return [
+ table_name
+ for table_name in _get_archived_table_names(["dag_version"], session)
+ if table_name.startswith(f"{ARCHIVE_TABLE_PREFIX}dag_version__")
+ ]
+
+
+def _dag_version_config_without_row_exclusion():
+ """Live dag_version cleanup config with row-exclusion filters disabled.
+
+ The backend regression tests above pin ``_do_delete``'s
rollback-before-drop
+ behaviour, which requires the
+ DELETE to actually hit the ``task_instance.dag_version_id`` FK violation.
+ The live config may exclude FK-pinned rows from the deletion query (see
+ PR #68339), which would turn these tests into no-ops -- so strip any such
+ filters from a copy of the config.
+ """
+ config = dict(config_dict["dag_version"].__dict__)
+ for key in ("extra_filters", "skip_if_referenced"):
+ if key in config:
+ config[key] = None
+ return config