This is an automated email from the ASF dual-hosted git repository.
hussein-awala 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 94a677e8f87 Keep callback triggers assigned to healthy triggerers
(#72848)
94a677e8f87 is described below
commit 94a677e8f877eea59589aa77022f106a3edcf340
Author: Hussein Awala <[email protected]>
AuthorDate: Mon Oct 5 11:50:07 2026 +0200
Keep callback triggers assigned to healthy triggerers (#72848)
Callback triggers could be repeatedly claimed by another healthy triggerer,
causing duplicate callback execution and consuming capacity intended for
waiting triggers.
---
airflow-core/src/airflow/models/trigger.py | 1 +
airflow-core/tests/unit/models/test_trigger.py | 45 ++++++++++++++++++++++++++
2 files changed, 46 insertions(+)
diff --git a/airflow-core/src/airflow/models/trigger.py
b/airflow-core/src/airflow/models/trigger.py
index 6a3e704fe79..aae270756be 100644
--- a/airflow-core/src/airflow/models/trigger.py
+++ b/airflow-core/src/airflow/models/trigger.py
@@ -462,6 +462,7 @@ class Trigger(Base):
# Callback triggers
select(cls.id)
.join(Callback, isouter=False)
+ .where(or_(cls.triggerer_id.is_(None),
cls.triggerer_id.not_in(alive_triggerer_ids)))
.order_by(Callback.priority_weight.desc(), cls.created_date),
# Task Instance triggers
select(cls.id)
diff --git a/airflow-core/tests/unit/models/test_trigger.py
b/airflow-core/tests/unit/models/test_trigger.py
index feb5c7aba6b..80961acb4d9 100644
--- a/airflow-core/tests/unit/models/test_trigger.py
+++ b/airflow-core/tests/unit/models/test_trigger.py
@@ -667,6 +667,51 @@ def test_assign_unassigned(session, create_triggerer,
create_trigger, use_queues
)
[email protected]("queue", [None, "callbacks"])
+def test_assign_unassigned_callbacks_preserves_healthy_owners(session,
create_triggerer, time_machine, queue):
+ now = timezone.datetime(2026, 1, 1)
+ time_machine.move_to(now, tick=False)
+ queues = {queue} if queue else None
+ healthy_owner = create_triggerer(session, State.RUNNING,
latest_heartbeat=now)
+ claiming_triggerer = create_triggerer(session, State.RUNNING,
latest_heartbeat=now)
+ stale_owner = create_triggerer(
+ session, State.RUNNING, latest_heartbeat=now -
datetime.timedelta(seconds=31)
+ )
+ finished_owner = create_triggerer(session, State.SUCCESS,
latest_heartbeat=now, end_date=now)
+ session.flush()
+
+ expected_owners = {}
+ for owner, priority in (
+ (healthy_owner, 10),
+ (claiming_triggerer, 10),
+ (stale_owner, 1),
+ (finished_owner, 1),
+ (None, 1),
+ ):
+ callback = TriggererCallback(
+ callback_def=AsyncCallback("asyncio.sleep", kwargs={"delay": 60},
queue=queue),
+ priority_weight=priority,
+ )
+ callback.queue(session=session)
+ callback.trigger.triggerer_id = owner.id if owner else None
+ session.add(callback)
+ session.flush()
+ expected_owners[callback.trigger.id] = (
+ healthy_owner.id if owner is healthy_owner else
claiming_triggerer.id
+ )
+ session.commit()
+
+ for triggerer in (claiming_triggerer, healthy_owner):
+ Trigger.assign_unassigned(
+ triggerer.id,
+ capacity=4,
+ health_check_threshold=30,
+ queues=queues,
+ )
+ session.expire_all()
+ assert dict(session.execute(select(Trigger.id,
Trigger.triggerer_id)).all()) == expected_owners
+
+
@pytest.mark.need_serialized_dag
@conf_vars({("triggerer", "queues_enabled"): "True"})
def test_assign_unassigned_with_qeueus(session, create_triggerer,
create_trigger) -> None: