kaxil commented on code in PR #73917:
URL: https://github.com/apache/airflow/pull/73917#discussion_r4138736201
##########
providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py:
##########
@@ -1146,21 +1177,31 @@ def adopt_launched_task(
self.log.error("attempting to adopt taskinstance which was not
specified by database: %s", ti_key)
return
- new_worker_id_label =
self._make_safe_label_value(self.scheduler_job_id)
+ key: UUID | TaskInstanceKey = ti_key
+ metadata = {"labels": {"airflow-worker":
self._make_safe_label_value(self.scheduler_job_id)}}
+ if self.supports_task_instance_uuid:
+ task_id = task_instance_id_from_pod(pod)
Review Comment:
I don't think the pre-upgrade case works against a real cluster.
`_list_pods` requests `PartialObjectMetadataList`, so the pods passed in here
have no spec (the dynamic client's `ResourceField` returns None for a missing
attribute). For any pod without the annotation, which is every pod launched
before the upgrade, `task_instance_id_from_pod` returns None and the pod is
refused. The scheduler then resets the TI with a fresh id, and the old pod
keeps running under the dead scheduler's `airflow-worker` label, where no
watcher or adoption pass will see it again. The UUID branch of `revoke_task`
has the same gap, since it also lists through `_list_pods`.
`test_adoption_uses_pod_identity_and_patches_uuid[annotated=False]` passes
because it builds a full `V1Pod` with a spec. Could unannotated candidates be
fetched with `read_namespaced_pod` before extracting the id? That would also
make the "pre-upgrade task pods can be adopted" line in
`kubernetes_executor.rst` hold.
##########
providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor_utils.py:
##########
@@ -172,6 +176,13 @@ def _run(
"run_id": annotations.get("run_id"),
"try_number": annotations["try_number"],
}
+ if TASK_INSTANCE_ID_ANNOTATION in annotations:
+ task_instance_related_annotations[TASK_INSTANCE_ID_ANNOTATION]
= annotations[
+ TASK_INSTANCE_ID_ANNOTATION
+ ]
+ elif hasattr(BaseExecutor, "get_task_key"):
Review Comment:
The watcher here and `process_watcher_task` decide UUID mode from
`hasattr(BaseExecutor, "get_task_key")`, while the executor uses
`supports_task_instance_uuid`. When the flag is False on a UUID-capable core,
the watcher recovers a UUID from the pod args, `_change_state` doesn't find it
in `running`, and the result is dropped. That is the state the new autouse
`coordinate_key_contract` fixture puts all of `TestKubernetesExecutor` in, so
on main the existing adoption, revoke and `_change_state` tests only run the
coordinate branches. The UUID branches added in `_iter_tis_to_flush`,
`_adopt_completed_pods`, `cleanup_stuck_queued_tasks` and the `"completed"`
exemption in `_change_state` aren't exercised anywhere. Could
`AirflowKubernetesScheduler` and the watcher take the executor's flag instead,
so those tests can run with the real value?
##########
providers/celery/tests/unit/celery/executors/test_celery_executor.py:
##########
@@ -1669,3 +1673,220 @@ def test_team_app_name_includes_team_name(self):
team_conf = ExecutorConf(team_name="team_beta")
celery_app = celery_executor_utils.create_celery_app(team_conf)
assert "team_beta" in celery_app.main
+
+
[email protected]
+def identity_workload():
+ if not AIRFLOW_V_3_0_PLUS:
+ pytest.skip("Task workloads require Airflow 3+")
+ return workloads.ExecuteTask(
+ ti=workloads.TaskInstance(
+ id=uuid4(),
+ dag_version_id=uuid4(),
+ dag_id="dag",
+ task_id="task",
+ run_id="run",
+ try_number=1,
+ map_index=-1,
+ pool_slots=1,
+ priority_weight=1,
+ queue="default",
+ external_executor_id="celery-a",
Review Comment:
`workloads.TaskInstance` only has `external_executor_id` from 3.2 (it isn't
in 3.0.6 or 3.1.8), and pydantic silently drops the unknown kwarg. So on the
3.0 and 3.1 compat jobs `test_celery_adoption_keeps_original_task_identity`
fails with `AttributeError` on `ti.external_executor_id` in
`try_adopt_task_instances`. Gate that test on `AIRFLOW_V_3_2_PLUS`?
##########
providers/amazon/src/airflow/providers/amazon/aws/executors/ecs/ecs_executor.py:
##########
@@ -669,7 +670,7 @@ def try_adopt_task_instances(self, tis:
Sequence[TaskInstance]) -> Sequence[Task
self.active_workers.add_task(
task,
- ti.key,
+ self.get_task_key(ti) if
self.supports_task_instance_uuid else ti.key,
Review Comment:
`prepare_db_for_next_try` gives the TI a new id but keeps
`external_executor_id`, and ECS doesn't pre-assign one at queue time. If a
retry is still queued (for example in `pending_workloads` after a capacity
failure) when the scheduler dies, adoption describes attempt 1's task ARN and
this binds its result to attempt 2's UUID. Lambda and K8s check identity at
this point; ECS and Batch can't. It happened with coordinate keys too, so
probably a core follow-up: clear `external_executor_id` in
`prepare_db_for_next_try`, as the clear path already does.
##########
providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py:
##########
@@ -1103,11 +1114,31 @@ def revoke_task(self, *, ti: TaskInstance):
assert self.kube_client
assert self.kube_scheduler
- self.running.discard(ti.key)
+ key = self.get_task_key(ti) if self.supports_task_instance_uuid else
ti.key
+ self.running.discard(key)
+ self.pod_launch_attempts.pop(key, None)
+ self.task_publish_retries.pop(key, None)
if AIRFLOW_V_3_4_PLUS:
- self.executor_queues[WorkloadType.EXECUTE_TASK].pop(ti.key, None)
+ self.executor_queues[WorkloadType.EXECUTE_TASK].pop(key, None)
else:
- self.queued_tasks.pop(ti.key, None)
+ self.queued_tasks.pop(key, None)
+ if self.supports_task_instance_uuid:
+ selector = PodGenerator.build_selector_for_k8s_executor_pod(
+ dag_id=ti.dag_id,
+ task_id=ti.task_id,
+ run_id=ti.run_id,
+ try_number=ti.try_number,
+ map_index=ti.map_index,
+ include_version=False,
+ )
+ for pod in self._list_pods({"label_selector": selector}):
Review Comment:
The coordinate path this replaces goes through
`get_pod_combined_search_str_to_pod_map`, which merges
`kube_client_request_args` into the list call. This one doesn't, so a
configured `_request_timeout` no longer applies to revoke.
##########
providers/celery/tests/unit/celery/executors/test_celery_executor.py:
##########
@@ -1669,3 +1673,220 @@ def test_team_app_name_includes_team_name(self):
team_conf = ExecutorConf(team_name="team_beta")
celery_app = celery_executor_utils.create_celery_app(team_conf)
assert "team_beta" in celery_app.main
+
+
[email protected]
+def identity_workload():
+ if not AIRFLOW_V_3_0_PLUS:
+ pytest.skip("Task workloads require Airflow 3+")
+ return workloads.ExecuteTask(
+ ti=workloads.TaskInstance(
+ id=uuid4(),
+ dag_version_id=uuid4(),
+ dag_id="dag",
+ task_id="task",
+ run_id="run",
+ try_number=1,
+ map_index=-1,
+ pool_slots=1,
+ priority_weight=1,
+ queue="default",
+ external_executor_id="celery-a",
+ ),
+ dag_rel_path="dag.py",
+ bundle_info=workloads.BundleInfo(name="bundle"),
+ token="",
+ log_path=None,
+ )
+
+
[email protected]
+def identity_executor(mocker):
+ app = Celery("identity-test", broker="memory://",
backend="cache+memory://")
+ mocker.patch.object(celery_executor_utils, "create_celery_app",
autospec=True, return_value=app)
+ executor = CeleryExecutor()
+ yield executor
+ app.close()
+
+
[email protected]("legacy", [False, True])
[email protected]("state", ["SUCCESS", "FAILURE"])
+def test_task_identity_survives_celery_dispatch_and_completion(
+ identity_executor, identity_workload, legacy, state, monkeypatch, mocker
+):
+ executor, workload = identity_executor, identity_workload
+ native = hasattr(BaseExecutor, "get_task_key") and not legacy
+ if legacy and hasattr(BaseExecutor, "get_task_key"):
+ monkeypatch.setattr(CeleryExecutor, "supports_task_instance_uuid",
False)
+ key = workload.ti.id if native else workload.ti.key
+ result = executor.celery_app.AsyncResult("celery-a")
+ sender = mocker.patch.object(
+ executor,
+ "_send_workloads_to_celery",
+ autospec=True,
+ side_effect=lambda items: [(item[0], None, result) for item in items],
+ )
+ executor.queue_workload(workload, session=None)
+ executor._process_workloads([workload])
+
+ assert executor.running == {key}
+ assert executor.workloads == {key: result}
+ assert not executor.queued_tasks
+ assert executor.get_event_buffer() == {key: (State.QUEUED, "celery-a")}
+ assert sender.call_args.args[0][0][0] == key
+
+ mocker.patch.object(
+ executor.bulk_state_fetcher, "get_many", autospec=True,
return_value={"celery-a": (state, None)}
+ )
+ executor.sync()
+ assert not executor.running
+ assert not executor.workloads
+ assert executor.get_event_buffer() == {key: (State.SUCCESS if state ==
"SUCCESS" else State.FAILED, None)}
+
+
[email protected](not hasattr(BaseExecutor, "get_task_key"),
reason="Requires UUID executor contract")
+def test_same_coordinate_tasks_keep_distinct_celery_results(identity_executor,
identity_workload, mocker):
+ first = identity_workload
+ second = first.model_copy(
+ update={"ti": first.ti.model_copy(update={"id": uuid4(),
"external_executor_id": "celery-b"})}
+ )
+ executor = identity_executor
+ results = {
+ "celery-a": executor.celery_app.AsyncResult("celery-a"),
+ "celery-b": executor.celery_app.AsyncResult("celery-b"),
+ }
+ mocker.patch.object(
+ executor,
+ "_send_workloads_to_celery",
+ autospec=True,
+ side_effect=lambda items: [
+ (item[0], None, results[item[1].ti.external_executor_id]) for item
in items
+ ],
+ )
+ executor.queue_workload(first, session=None)
+ executor.queue_workload(second, session=None)
+ executor._process_workloads([first, second])
+ executor.get_event_buffer()
+ mocker.patch.object(
+ executor.bulk_state_fetcher,
+ "get_many",
+ autospec=True,
+ return_value={"celery-a": ("SUCCESS", None), "celery-b": ("PENDING",
None)},
+ )
+ executor.sync()
+ assert executor.get_event_buffer() == {first.ti.id: (State.SUCCESS, None)}
+ assert executor.running == {second.ti.id}
+ assert executor.workloads == {second.ti.id: results["celery-b"]}
+
+
[email protected]("state", ["PENDING", "SUCCESS", "FAILURE"])
[email protected]("legacy", [False, True])
+def test_celery_adoption_keeps_original_task_identity(
+ identity_executor, identity_workload, state, legacy, monkeypatch, mocker
+):
+ executor, ti = identity_executor, identity_workload.ti
+ native = hasattr(BaseExecutor, "get_task_key") and not legacy
+ if legacy and hasattr(BaseExecutor, "get_task_key"):
+ monkeypatch.setattr(CeleryExecutor, "supports_task_instance_uuid",
False)
+ key = ti.id if native else ti.key
+ result = executor.celery_app.AsyncResult("celery-a")
+ mocker.patch("celery.result.AsyncResult", autospec=True,
return_value=result)
+ mocker.patch.object(
+ executor.bulk_state_fetcher, "get_many", autospec=True,
return_value={"celery-a": (state, None)}
+ )
+ assert executor.try_adopt_task_instances([ti]) == []
+ ti.id = uuid4()
+ if state == "PENDING":
+ assert executor.running == {key}
+ assert executor.workloads == {key: result}
+ else:
+ assert not executor.running
+ assert not executor.workloads
+ assert executor.get_event_buffer() == {
+ key: (State.SUCCESS if state == "SUCCESS" else State.FAILED, None)
+ }
+
+
[email protected]("legacy", [False, True])
+def test_celery_revoke_removes_only_target_attempt(
+ identity_executor, identity_workload, legacy, monkeypatch, mocker
+):
+ executor, workload = identity_executor, identity_workload
+ native = hasattr(BaseExecutor, "get_task_key") and not legacy
+ if legacy and hasattr(BaseExecutor, "get_task_key"):
+ monkeypatch.setattr(CeleryExecutor, "supports_task_instance_uuid",
False)
+ key = workload.ti.id if native else workload.ti.key
+ executor.queue_workload(workload, session=None)
+ executor.running.add(key)
+ executor.workloads[key] = executor.celery_app.AsyncResult("celery-a")
+ other_key = uuid4() if native else workload.ti.key.with_try_number(2)
+ executor.running.add(other_key)
+ executor.workloads[other_key] = executor.celery_app.AsyncResult("celery-b")
+ revoke = mocker.patch.object(executor.celery_app.control, "revoke",
autospec=True)
+ executor.revoke_task(ti=workload.ti)
+ assert executor.running == {other_key}
+ assert set(executor.workloads) == {other_key}
+ assert not executor.queued_tasks
+ revoke.assert_called_once_with("celery-a")
+
+
[email protected](not hasattr(BaseExecutor, "get_task_key"),
reason="Requires UUID executor contract")
+def test_uuid_publish_timeout_keeps_task_queued_for_retry(identity_executor,
identity_workload, mocker):
+ executor, workload = identity_executor, identity_workload
+ key = workload.ti.id
+ executor.queue_workload(workload, session=None)
+ failure =
celery_executor_utils.ExceptionWithTraceback(AirflowTaskTimeout(), "traceback")
+ mocker.patch.object(
+ executor, "_send_workloads_to_celery", autospec=True,
return_value=[(key, None, failure)]
+ )
+ executor._process_workloads([workload])
+ assert executor.workload_publish_retries[key] == 1
+ assert executor.has_task(workload.ti)
+ assert executor.get_event_buffer() == {}
+
+
[email protected](not AIRFLOW_V_3_2_PLUS, reason="Callback workloads require
Airflow 3.2+")
Review Comment:
This needs `AIRFLOW_V_3_3_PLUS`. On 3.2.2 `ExecuteCallback` has no `key`
property (only `CallbackDTO` does), and neither workload has
`success_state`/`failure_state`, so it fails on the 3.2 compat job.
##########
providers/amazon/docs/executors/general.rst:
##########
@@ -363,3 +363,12 @@ The Airflow DB needs to be initialized before it can be
used and a user needs to
airflow users create --username admin --password admin --firstname <your
first name> --lastname <your last name> --email <your email> --role Admin
.. END INIT_DB
+
+Task-instance identity
Review Comment:
This lands after `.. END INIT_DB`, outside any BEGIN/END pair, and
`general.rst` isn't in a toctree. The ECS, Batch and Lambda pages only pull in
fragments with `start-after`, so this section never renders. It needs its own
marker pair and an include on each executor page.
--
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]