kaxil commented on code in PR #70961:
URL: https://github.com/apache/airflow/pull/70961#discussion_r4125259456


##########
airflow-core/src/airflow/models/trigger.py:
##########
@@ -231,42 +231,54 @@ def fetch_trigger_ids_with_non_task_associations(cls, *, 
session: Session = NEW_
         return set(session.scalars(query))
 
     @classmethod
-    @provide_session
-    def clean_unused(cls, *, session: Session = NEW_SESSION) -> None:
+    def clean_unused(cls) -> None:
         """
-        Delete all triggers that have no tasks dependent on them and are not 
associated to an asset.
+        Delete triggers that have no dependent tasks, assets, or callbacks in 
bounded transactions.
 
         Triggers have a one-to-many relationship to task instances, so we need 
to clean those up first.
         Afterward we can drop the triggers not referenced by anyone.
         """
-        # Update all task instances with trigger IDs that are not DEFERRED to 
remove them
-        for attempt in run_with_db_retries():
-            with attempt:
-                session.execute(
-                    update(TaskInstance)
-                    .where(
-                        TaskInstance.state != TaskInstanceState.DEFERRED, 
TaskInstance.trigger_id.is_not(None)
-                    )
-                    .values(trigger_id=None)
-                )
-
-        # Get all triggers that have no task instances, assets, or callbacks 
depending on them and delete them
-        ids = select(cls.id).where(
-            ~cls.assets.any(),
-            ~cls.callback.has(),
-            ~cls.task_instance.has(),
-        )
+        batch_size = conf.getint("triggerer", 
"unreferenced_triggers_cleanup_batch_size", fallback=500)
+        if batch_size <= 0:
+            raise ValueError("[triggerer] 
unreferenced_triggers_cleanup_batch_size must be at least 1")
+
+        clear_task_instance_references = True
+        while True:

Review Comment:
   Thanks, the re-check in the DELETE looks right and the new race test covers 
all three referrers. One thing left: this loop still drains the whole backlog 
inside a single `clean_unused()` call, and `run_once` only reaches 
`heartbeat()` after it returns, so the loop stall #68243 describes is 
unchanged. Batching bounds how long each transaction holds locks, not how long 
the supervisor is away from the heartbeat and the runner comms.
   
   Could this do one batch per call instead (or cap the batches per call)? 
`clean_unused` already runs on every supervisor loop, so the backlog still 
drains, and it lets you drop the `while True` and the 
`clear_task_instance_references` flag. If you want to keep the full drain, the 
description should stop claiming it fixes the heartbeat delay.
   
   On Postgres, the time per batch also hinges on the `callback.trigger_id` 
index you scoped out. That column has an FK to `trigger.id` but no index, so 
every deleted trigger row runs the FK check as a sequential scan of `callback`. 
In a local Postgres 16 run with ~500k callback rows, one 500-row batch took 
about 5,600 ms, and about 6 ms once `callback(trigger_id)` was indexed. So even 
a single batch isn't bounded without it. Fine as a follow-up, but it's the part 
that actually bounds the time.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to