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


##########
airflow-core/src/airflow/jobs/scheduler_job_runner.py:
##########
@@ -1469,16 +1468,14 @@ def process_executor_events(
         `dag.test` execute DAGs with no scheduler, therefore it needs to 
handle the events pushed by the
         executors as well.
         """
-        ti_event_keys: dict[tuple[str, str, str, int], list[TaskInstanceKey]] 
= defaultdict(list)
-        event_buffer = executor.get_event_buffer()
+        event_buffer = executor.get_event_buffer_with_task_ids()
         num_events = len(event_buffer)
-        tis_with_right_state: list[TaskInstanceKey] = []
+        tis_with_right_state: list[UUID] = []
         callback_keys_with_events: list[CallbackKey] = []
 
         # Report execution - handle both task and callback events
         for key, (state, _) in event_buffer.items():
-            if isinstance(key, TaskInstanceKey):
-                ti_event_keys[key.primary].append(key)
+            if isinstance(key, UUID):
                 cls.logger().info("Received executor event with state %s for 
task instance %s", state, key)

Review Comment:
   This now prints only the UUID. For events that match a row the later 
"TaskInstance Finished" line fills in the rest, but for one that doesn't 
there's nothing left that names the task, since the legacy coordinates get 
pruned in the same drain. The "Ignoring executor event for a different attempt" 
line also went away with nothing replacing it, so an event for a retired 
attempt now misses the `TI.id.in_` select and vanishes silently.
   
   Could `get_event_buffer_with_task_ids` hand back the coordinates it already 
has, and the scheduler log whatever UUIDs are still in `event_buffer` after the 
loop? That leftover set also includes rows skipped by `skip_locked`, so the 
message should allow for both.



##########
airflow-core/tests/unit/executors/test_base_executor.py:
##########
@@ -185,6 +218,161 @@ def test_fail_and_success():
     assert len(executor.get_event_buffer()) == 3
 
 
[email protected]
+def task_workload():
+    return workloads.ExecuteTask(
+        ti=workloads.TaskInstanceDTO(
+            id=uuid4(),
+            dag_version_id=uuid4(),
+            dag_id="dag",
+            task_id="task",
+            run_id="run",
+            try_number=1,
+            pool_slots=1,
+            priority_weight=1,
+            queue="default",
+        ),
+        dag_rel_path="dag.py",
+        bundle_info=BundleInfo(name="bundle"),
+        token="",
+        log_path=None,
+    )
+
+
[email protected]("override_queue", [False, True])
[email protected]("terminal_state", [TaskInstanceState.SUCCESS, 
TaskInstanceState.FAILED])
[email protected]("next_try_number", [1, 2])
+def test_legacy_executor_dispatch_and_completion_keep_submitted_identity(
+    task_workload, override_queue, terminal_state, next_try_number
+):
+    class LegacyExecutor(BaseExecutor):
+        def _process_workloads(self, items):
+            for workload in items:
+                key = workload.ti.key
+                del self.executor_queues[WorkloadType.EXECUTE_TASK][key]
+                self.running.add(key)
+                self.event_buffer[key] = (TaskInstanceState.QUEUED, 
"remote-id")
+
+    class OverrideExecutor(LegacyExecutor):
+        def queue_workload(self, workload, session):
+            self.executor_queues[workload.type][workload.key] = workload
+
+    executor = (OverrideExecutor if override_queue else LegacyExecutor)()
+    submitted_id, submitted_key = task_workload.ti.id, task_workload.ti.key
+    executor.queue_workload(task_workload, session=None)
+    assert executor.has_task(task_workload.ti)
+    executor.heartbeat()

Review Comment:
   Nothing in this test drains events while the workload is still queued, so it 
passes with `task_queue` removed from the pruning loop in `get_event_buffer`. 
That clause is what keeps the snapshot alive for a legacy executor sitting at 
its parallelism cap. An `assert executor.get_event_buffer_with_task_ids() == 
{}` just before `heartbeat()` would pin it: without the clause, that drain 
prunes the snapshot and the queued event after `heartbeat()` gets discarded.



##########
airflow-core/src/airflow/executors/workloads/types.py:
##########
@@ -27,7 +28,7 @@
 from airflow.utils.state import CallbackState, TaskInstanceState
 
 # Type aliases for workload keys and states (used by executor layer)
-WorkloadKey: TypeAlias = TaskInstanceKey | CallbackKey | ConnectionTestKey
+WorkloadKey: TypeAlias = UUID | TaskInstanceKey | CallbackKey | 
ConnectionTestKey

Review Comment:
   `CallbackKey` and `ConnectionTestKey` both wrap their id in a frozen 
dataclass so a bare value can't pass the wrong `isinstance` check (the 
`CallbackKey` docstring says as much). A bare `UUID` for tasks breaks that 
pattern, and once the provider PRs start emitting it, it becomes a 
provider-facing contract that's hard to change. Would a small typed 
`TaskInstanceUUIDKey` be worth it before then? Nothing collides today.



##########
airflow-core/docs/core-concepts/executor/index.rst:
##########
@@ -336,6 +336,18 @@ The ``BaseExecutor`` class interface contains a set of 
attributes that Airflow c
 * ``is_single_threaded``: Whether or not the executor is single threaded. This 
is particularly relevant to what database backends are supported. Single 
threaded executors can run with any backend, including SQLite.
 * ``is_production``: Whether or not the executor should be used for production 
purposes. A UI message is displayed to users when they are using a 
non-production ready executor.
 * ``serve_logs``: Whether or not the executor supports serving logs, see 
:doc:`/administration-and-deployment/logging-monitoring/logging-tasks`.
+* ``supports_task_instance_uuid``: Whether task submission, adoption, result 
decoding and cleanup all use task-instance UUIDs. Defaults to ``False`` to 
support existing provider releases. ``LocalExecutor`` enables this capability.
+
+UUID-capable executors use ``get_workload_key(workload)`` and 
``get_task_key(ti)`` for executor bookkeeping.
+``ExecuteTask.key`` retains its coordinate-based value for existing providers. 
Providers supporting older
+Airflow releases must retain their older-core paths when adopting these 
helpers.

Review Comment:
   Could this name the version? `get_task_key`, `get_workload_key` and 
`supports_task_instance_uuid` don't exist on 3.3.x, so a provider calling the 
helpers there gets an `AttributeError`. Pointing at `AIRFLOW_V_3_4_PLUS` (or 
adding a `versionadded:: 3.4.0`) would save each provider PR from rediscovering 
that.



##########
airflow-core/src/airflow/executors/base_executor.py:
##########
@@ -382,10 +413,29 @@ def start(self):  # pragma: no cover
 
     def log_task_event(self, *, event: str, extra: str, ti_key: WorkloadKey):
         """Add an event to the log table."""
-        if not isinstance(ti_key, TaskInstanceKey):
+        if not isinstance(ti_key, (UUID, TaskInstanceKey)):
             self.log.debug("Skipping log_task_event for callback key %s 
(event=%s)", ti_key, event)
             return
-        self._task_event_logs.append(Log(event=event, task_instance=ti_key, 
extra=extra))
+        coordinates = ti_key if isinstance(ti_key, TaskInstanceKey) else 
self._task_coordinates.get(ti_key)
+        if coordinates is None:
+            extra = f"Task instance {ti_key}: {extra}"
+        self._task_event_logs.append(Log(event=event, 
task_instance=coordinates, extra=extra))
+
+    def register_task(self, ti: TaskInstance | TaskInstanceDTO) -> None:

Review Comment:
   Are providers meant to call `register_task` and 
`get_event_buffer_with_task_ids`? In core the only callers are the 
`__init_subclass__` wrappers, the base `queue_workload` and 
`process_executor_events`, but the Kubernetes test in this PR already reaches 
for `register_task` behind `hasattr`, which is how provider code will start 
depending on it too. If they're plumbing for the legacy bridge, a `_` prefix 
keeps them removable when the bridge goes. If providers are meant to use them, 
they probably want a line in the executor docs next to `get_task_key`.



##########
airflow-core/src/airflow/executors/base_executor.py:
##########
@@ -248,6 +253,31 @@ def jwt_generator(self) -> JWTGenerator:
 
     def __init_subclass__(cls, **kwargs: Any) -> None:
         super().__init_subclass__(**kwargs)
+        if queue_workload := cls.__dict__.get("queue_workload"):
+
+            @wraps(queue_workload)
+            def register_queued_task(self, workload, *args, **kwargs):
+                if isinstance(workload, ExecuteTask):
+                    self.register_task(workload.ti)
+                return queue_workload(self, workload, *args, **kwargs)
+
+            cls.queue_workload = register_queued_task  # type: 
ignore[method-assign]
+        if try_adopt := cls.__dict__.get("try_adopt_task_instances"):
+
+            @wraps(try_adopt)
+            def register_adopted_tasks(self, tis, *args, **kwargs):
+                keys = [(ti.id, self.get_task_key(ti)) for ti in tis]
+                for ti in tis:
+                    self.register_task(ti)
+                rejected = try_adopt(self, tis, *args, **kwargs)
+                rejected_ids = {ti.id for ti in rejected}
+                for task_id, key in keys:
+                    state, _ = self.event_buffer.get(key, (None, None))
+                    if task_id not in rejected_ids and state not in 
State.finished:
+                        self.running.add(key)

Review Comment:
   This also puts adopted ECS, Batch and Lambda tasks into `running`, which 
their own `try_adopt_task_instances` never did, so after a failover they now 
count against `parallelism` until they finish. I think that's right, and the 
pruning loop needs the key there anyway, but it's a slot-accounting change for 
those executors on a core-only upgrade and probably worth a line in the 
description.



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