This is an automated email from the ASF dual-hosted git repository.
ferruzzi 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 6f7e635d3d9 Include task_reschedule in task_instance dependent tables
for db clean (#71990)
6f7e635d3d9 is described below
commit 6f7e635d3d97d8fcae417bbbe034bb01637549e3
Author: gaurav kumar pandey <[email protected]>
AuthorDate: Wed Sep 23 03:11:45 2026 +0530
Include task_reschedule in task_instance dependent tables for db clean
(#71990)
In Airflow 3, task_reschedule rows reference task_instance via ti_id
foreign key with CASCADE delete. When cleaning task_instance rows via airflow
db clean, task_reschedule rows must be cleaned and archived before
task_instance rows are removed.
- Add task_reschedule to task_instance.dependent_tables in
airflow.utils.db_cleanup.
- Add test coverage verifying task_reschedule cleanup and archiving when
task_instance is cleaned.
---
airflow-core/newsfragments/71990.bugfix.rst | 1 +
airflow-core/src/airflow/utils/db_cleanup.py | 2 +-
airflow-core/tests/unit/utils/test_db_cleanup.py | 121 +++++++++++++++++++++++
3 files changed, 123 insertions(+), 1 deletion(-)
diff --git a/airflow-core/newsfragments/71990.bugfix.rst
b/airflow-core/newsfragments/71990.bugfix.rst
new file mode 100644
index 00000000000..819b6ab5953
--- /dev/null
+++ b/airflow-core/newsfragments/71990.bugfix.rst
@@ -0,0 +1 @@
+Fix ``airflow db clean`` to include ``task_reschedule`` in ``task_instance``
dependent tables so task reschedules are cleaned and archived prior to task
instances.
diff --git a/airflow-core/src/airflow/utils/db_cleanup.py
b/airflow-core/src/airflow/utils/db_cleanup.py
index 9c56814632b..ad41b482c44 100644
--- a/airflow-core/src/airflow/utils/db_cleanup.py
+++ b/airflow-core/src/airflow/utils/db_cleanup.py
@@ -258,7 +258,7 @@ config_list: list[_TableConfig] = [
_TableConfig(
table_name="task_instance",
recency_column_name="start_date",
- dependent_tables=["task_instance_history", "xcom"],
+ dependent_tables=["task_instance_history", "xcom", "task_reschedule"],
dag_id_column_name="dag_id",
),
_TableConfig(
diff --git a/airflow-core/tests/unit/utils/test_db_cleanup.py
b/airflow-core/tests/unit/utils/test_db_cleanup.py
index da3e7bae4a2..b55ef1ace1b 100644
--- a/airflow-core/tests/unit/utils/test_db_cleanup.py
+++ b/airflow-core/tests/unit/utils/test_db_cleanup.py
@@ -52,6 +52,7 @@ from airflow.models.dagbundle import DagBundleModel
from airflow.models.deadline import Deadline
from airflow.models.serialized_dag import SerializedDagModel
from airflow.models.task_state_store import TaskStateStoreModel
+from airflow.models.taskreschedule import TaskReschedule
from airflow.providers.standard.operators.python import PythonOperator
from airflow.sdk.definitions.callback import AsyncCallback
from airflow.serialization.serialized_objects import LazyDeserializedDAG
@@ -2002,6 +2003,126 @@ class TestSchemaQualifiedTableConfig:
assert config.orm_model.schema is None
[email protected]_test
+class TestTaskRescheduleCleanup:
+ @pytest.fixture(autouse=True)
+ def clear_airflow_tables(self):
+ drop_tables_with_prefix("_airflow_")
+ yield
+ drop_tables_with_prefix("_airflow_")
+
+ def test_cleanup_task_reschedule(self):
+ base_date = pendulum.DateTime(2023, 1, 1,
tzinfo=pendulum.timezone("UTC"))
+ bundle_name = "testing"
+ with create_session() as session:
+ session.add(DagBundleModel(name=bundle_name))
+ session.flush()
+
+ dag_id = f"test-tr-cleanup_{uuid4()}"
+ dag = DAG(dag_id=dag_id)
+ dm = DagModel(dag_id=dag_id, bundle_name=bundle_name)
+ session.add(dm)
+ SerializedDagModel.write_dag(LazyDeserializedDAG.from_dag(dag),
bundle_name=bundle_name)
+ dag_version = DagVersion.get_latest_version(dag.dag_id)
+
+ dag_run = DagRun(
+ dag.dag_id,
+ run_id="run_1",
+ run_type=DagRunType.SCHEDULED,
+ start_date=base_date,
+ )
+ ti = create_task_instance(
+ PythonOperator(task_id="dummy-task", python_callable=print),
+ run_id=dag_run.run_id,
+ dag_version_id=dag_version.id,
+ )
+ ti.dag_id = dag.dag_id
+ ti.start_date = base_date
+ session.add(dag_run)
+ session.add(ti)
+ session.flush()
+
+ tr_old = TaskReschedule(
+ ti_id=ti.id,
+ start_date=base_date,
+ end_date=base_date.add(minutes=1),
+ reschedule_date=base_date.add(minutes=5),
+ )
+ tr_new = TaskReschedule(
+ ti_id=ti.id,
+ start_date=base_date.add(days=10),
+ end_date=base_date.add(days=10, minutes=1),
+ reschedule_date=base_date.add(days=10, minutes=5),
+ )
+ session.add_all([tr_old, tr_new])
+ session.commit()
+
+ run_cleanup(
+ clean_before_timestamp=base_date.add(days=5),
+ table_names=["task_reschedule"],
+ dry_run=False,
+ confirm=False,
+ session=session,
+ )
+
+ remaining = session.scalars(select(TaskReschedule)).all()
+ assert len(remaining) == 1
+ assert remaining[0].id == tr_new.id
+
+ def test_cleanup_task_reschedule_as_task_instance_dependent(self):
+ base_date = pendulum.DateTime(2023, 1, 1,
tzinfo=pendulum.timezone("UTC"))
+ bundle_name = "testing"
+ with create_session() as session:
+ session.add(DagBundleModel(name=bundle_name))
+ session.flush()
+
+ dag_id = f"test-tr-cascade_{uuid4()}"
+ dag = DAG(dag_id=dag_id)
+ dm = DagModel(dag_id=dag_id, bundle_name=bundle_name)
+ session.add(dm)
+ SerializedDagModel.write_dag(LazyDeserializedDAG.from_dag(dag),
bundle_name=bundle_name)
+ dag_version = DagVersion.get_latest_version(dag.dag_id)
+
+ dag_run = DagRun(
+ dag.dag_id,
+ run_id="run_1",
+ run_type=DagRunType.SCHEDULED,
+ start_date=base_date,
+ )
+ ti = create_task_instance(
+ PythonOperator(task_id="dummy-task", python_callable=print),
+ run_id=dag_run.run_id,
+ dag_version_id=dag_version.id,
+ )
+ ti.dag_id = dag.dag_id
+ ti.start_date = base_date
+ session.add(dag_run)
+ session.add(ti)
+ session.flush()
+
+ tr = TaskReschedule(
+ ti_id=ti.id,
+ start_date=base_date,
+ end_date=base_date.add(minutes=1),
+ reschedule_date=base_date.add(minutes=5),
+ )
+ session.add(tr)
+ session.commit()
+
+ run_cleanup(
+ clean_before_timestamp=base_date.add(days=5),
+ table_names=["task_instance"],
+ dry_run=False,
+ confirm=False,
+ session=session,
+ )
+
+ assert session.scalar(select(func.count(TaskInstance.id))) == 0
+ assert session.scalar(select(func.count(TaskReschedule.id))) == 0
+ archives = _get_archived_table_names(["task_reschedule"], session)
+ assert len(archives) == 1
+
+
@pytest.mark.backend("postgres")
class TestSchemaQualifiedTableCleanupIntegration:
"""