This is an automated email from the ASF dual-hosted git repository.
ashb 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 79de1d48039 Fix `refresh_from_task` not applied on retries (#65932)
79de1d48039 is described below
commit 79de1d4803941e7ba4a6361529cc7079bb36373e
Author: Jorge Rocamora <[email protected]>
AuthorDate: Sat Oct 3 11:19:13 2026 +0200
Fix `refresh_from_task` not applied on retries (#65932)
task_instance_mutation_hook and dynamic PriorityWeightStrategy
subclasses are not re-applied on natural retries: refresh_from_task
is missing from the UP_FOR_RETRY → SCHEDULED transition.
This adds it to DagRun.schedule_tis against the new computed
try_number. First attempts and UP_FOR_RESCHEDULE are unaffected,
per-TI UPDATEs only fire when field values actually change.
---
.../cluster-policies.rst | 11 ++--
airflow-core/src/airflow/models/dagrun.py | 4 ++
airflow-core/tests/unit/models/test_dagrun.py | 68 ++++++++++++++++++++++
3 files changed, 79 insertions(+), 4 deletions(-)
diff --git
a/airflow-core/docs/administration-and-deployment/cluster-policies.rst
b/airflow-core/docs/administration-and-deployment/cluster-policies.rst
index 1a09f69916b..c08a6bebfa0 100644
--- a/airflow-core/docs/administration-and-deployment/cluster-policies.rst
+++ b/airflow-core/docs/administration-and-deployment/cluster-policies.rst
@@ -41,10 +41,13 @@ There are three main types of cluster policy:
``task_instance_mutation_hook`` applies not to a task but to the instance of
a task that
relates to a particular DagRun. It is executed scheduler-side while task
instances are created or
reconciled (not in the Dag file processor, and not on the worker). The
policy is only applied to the
- currently executed run (i.e. instance) of that task. The ``dag_run``
argument lets the policy route on
- run configuration (``dag_run.conf``); it may be ``None`` in early
task-instance construction, and a hook
- that only declares ``task_instance`` keeps working unchanged. Note that
``dag_run.conf`` is only populated
- for manually triggered or API-triggered runs; scheduled runs carry an empty
``conf``.
+ currently executed run (i.e. instance) of that task. The policy can run more
than once for the same task
+ instance, for example when it is created or reconciled, so implementations
should be idempotent. It also
+ runs when a retry is scheduled, with ``task_instance.try_number`` set to the
attempt about to run. The
+ ``dag_run`` argument lets the policy route on run configuration
(``dag_run.conf``); it may be ``None`` in
+ early task-instance construction, and a hook that only declares
``task_instance`` keeps working
+ unchanged. Note that ``dag_run.conf`` is only populated for manually
triggered or API-triggered runs;
+ scheduled runs carry an empty ``conf``.
.. warning::
diff --git a/airflow-core/src/airflow/models/dagrun.py
b/airflow-core/src/airflow/models/dagrun.py
index c2b165d49cc..fc26cbccd16 100644
--- a/airflow-core/src/airflow/models/dagrun.py
+++ b/airflow-core/src/airflow/models/dagrun.py
@@ -2257,6 +2257,10 @@ class DagRun(Base, LoggingMixin):
debug_try_number_check = self.log.isEnabledFor(logging.DEBUG)
expected_try_number_by_ti_id: dict[UUID, tuple[int, int, str | None]]
= {}
for ti in schedulable_tis:
+ if ti.state == TaskInstanceState.UP_FOR_RETRY:
+ if TYPE_CHECKING:
+ assert ti.task
+ ti.refresh_from_task(ti.task, dag_run=self)
if not ti.is_schedulable:
empty_ti_ids.append(ti.id)
# The defer_task method will check "start_trigger_args" to see
whether the operator
diff --git a/airflow-core/tests/unit/models/test_dagrun.py
b/airflow-core/tests/unit/models/test_dagrun.py
index 07a3151181b..e897a59f35e 100644
--- a/airflow-core/tests/unit/models/test_dagrun.py
+++ b/airflow-core/tests/unit/models/test_dagrun.py
@@ -95,6 +95,7 @@ from tests_common.test_utils.mapping import
expand_mapped_task, push_mapped_leng
from tests_common.test_utils.mock_operators import MockOperator
from tests_common.test_utils.taskinstance import create_task_instance,
run_task_instance
from unit.models import DEFAULT_DATE as _DEFAULT_DATE
+from unit.plugins.priority_weight_strategy import DecreasingPriorityStrategy,
TestPriorityWeightStrategyPlugin
if TYPE_CHECKING:
from airflow.serialization.definitions.dag import SerializedDAG
@@ -2832,6 +2833,73 @@ def
test_schedule_tis_up_for_reschedule_does_not_increment_try_number(dag_maker,
assert refreshed_ti.try_number == 3
[email protected]_plugin_manager(plugins=[TestPriorityWeightStrategyPlugin])
+def test_schedule_tis_refreshes_task_instance_only_on_retry(dag_maker,
session):
+ with dag_maker(session=session) as dag:
+ for task_id in ("first_attempt", "retry", "reschedule"):
+ BashOperator(task_id=task_id, bash_command="echo 1",
weight_rule=DecreasingPriorityStrategy())
+
+ dr = dag_maker.create_dagrun(session=session)
+ tis = {ti.task_id: ti for ti in dr.get_task_instances(session=session)}
+ for task_id, state, try_number in (
+ ("first_attempt", None, 0),
+ ("retry", TaskInstanceState.UP_FOR_RETRY, 2),
+ ("reschedule", TaskInstanceState.UP_FOR_RESCHEDULE, 3),
+ ):
+ tis[task_id].refresh_from_task(dag.get_task(task_id))
+ tis[task_id].state = state
+ tis[task_id].try_number = try_number
+ session.commit()
+
+ hook_calls = []
+
+ def route_to_retry_queue(task_instance, dag_run=None):
+ hook_calls.append((task_instance.task_id, task_instance.try_number,
dag_run))
+ task_instance.queue = "retry_queue"
+
+ with _registered_mutation_hook(route_to_retry_queue):
+ assert dr.schedule_tis(tis.values(), session=session) == 3
+ session.commit()
+
+ assert hook_calls == [("retry", 2, dr)]
+ session.expire_all()
+ retry_ti = dr.get_task_instance("retry", session=session)
+ assert retry_ti.state == TaskInstanceState.SCHEDULED
+ assert retry_ti.try_number == 2
+ assert retry_ti.queue == "retry_queue"
+ assert retry_ti.priority_weight == 2
+
+
+def test_schedule_tis_refreshes_a_retry_that_defers_from_trigger(dag_maker,
session):
+ with dag_maker(session=session):
+ task = MockOperator(task_id="task")
+ task.start_from_trigger = True
+ task.start_trigger_args = StartTriggerArgs(
+ trigger_cls="airflow.triggers.testing.SuccessTrigger",
+ next_method="execute_complete",
+ )
+
+ dr = dag_maker.create_dagrun(session=session)
+ ti = dag_maker.create_ti("task", dag_run=dr)
+ ti.state = TaskInstanceState.UP_FOR_RETRY
+ ti.try_number = 2
+ session.commit()
+
+ hook_calls = []
+
+ def route_to_retry_queue(task_instance, dag_run=None):
+ hook_calls.append((task_instance.task_id, task_instance.try_number))
+ task_instance.queue = "retry_queue"
+
+ with _registered_mutation_hook(route_to_retry_queue):
+ dr.schedule_tis((ti,), session=session)
+ session.commit()
+
+ assert hook_calls == [("task", 2)]
+ session.expire_all()
+ assert (ti.state, ti.queue) == (TaskInstanceState.DEFERRED, "retry_queue")
+
+
def test_schedule_tis_empty_operator_is_noop_if_ti_already_running(dag_maker,
session):
with dag_maker(session=session) as dag:
EmptyOperator(task_id="empty_task")