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

Reply via email to