wolvery commented on code in PR #69762:
URL: https://github.com/apache/airflow/pull/69762#discussion_r3629206631


##########
providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py:
##########
@@ -391,6 +391,94 @@ def _process_workloads(self, workloads: 
Sequence[workloads.All]) -> None:
             self.execute_async(key=key, command=command, queue=queue, 
executor_config=executor_config)
             self.running.add(key)
 
+    def _should_create_pod_for_job(self, task: KubernetesJob) -> bool:
+        """
+        Check whether an executor job still represents the current queued task 
instance.
+
+        The scheduler creates an ``ExecuteTask`` workload while the task 
instance is queued, but the
+        Kubernetes pod may be created much later, for example after API-server 
throttling or quota
+        failures. In an HA scheduler deployment, the task instance may have 
been retried, cleared, or
+        otherwise replaced before this executor gets another chance to create 
the pod. Revalidating the
+        immutable task instance id and launch ownership here prevents an 
obsolete workload from creating
+        a stale worker pod.
+        """
+        try:
+            from airflow.executors.workloads import ExecuteTask
+        except ImportError:
+            # Compatibility with older Airflow versions tested by provider 
compatibility jobs.
+            return True
+
+        if not task.command or not isinstance(task.command[0], ExecuteTask):
+            return True
+
+        return self._should_create_pod_for_execute_task(task, task.command[0])
+
+    @provide_session
+    def _should_create_pod_for_execute_task(
+        self,
+        task: KubernetesJob,
+        workload: Any,
+        *,
+        session: Session = NEW_SESSION,
+    ) -> bool:
+        """Check that an ``ExecuteTask`` workload still owns the queued task 
instance row."""
+        from airflow.models.taskinstance import TaskInstance
+
+        workload_ti = workload.ti
+        try:
+            scheduler_job_id = int(self.scheduler_job_id) if 
self.scheduler_job_id is not None else None
+        except ValueError:
+            self.log.debug(
+                "Skipping stale Kubernetes workload check because 
scheduler_job_id %r is not numeric",
+                self.scheduler_job_id,
+            )
+            return True
+
+        ti = session.execute(
+            select(
+                TaskInstance.id,
+                TaskInstance.state,
+                TaskInstance.try_number,
+                TaskInstance.queued_by_job_id,
+            ).where(cast(TaskInstance.id, String) == str(workload_ti.id))
+        ).one_or_none()
+        if ti is None:
+            self.log.info(
+                "Dropping stale Kubernetes workload for %s because task 
instance id %s no longer exists",
+                task.key,
+                workload_ti.id,
+            )
+            return False
+
+        _, state, try_number, queued_by_job_id = ti
+        if (
+            state == TaskInstanceState.QUEUED
+            and try_number == workload_ti.try_number
+            and queued_by_job_id == scheduler_job_id
+        ):
+            return True
+
+        self.log.info(
+            "Dropping stale Kubernetes workload for %s because current task 
instance state does not "
+            "match the queued workload. task_instance_id=%s, state=%s, 
try_number=%s, "
+            "queued_by_job_id=%s, workload_try_number=%s, scheduler_job_id=%s",
+            task.key,
+            workload_ti.id,
+            state,
+            try_number,
+            queued_by_job_id,
+            workload_ti.try_number,
+            scheduler_job_id,
+        )
+        return False
+
+    def _discard_stale_pod_creation_task(self, task: KubernetesJob) -> None:
+        """Remove executor bookkeeping for a stale job that will not create a 
pod."""
+        self.running.discard(task.key)
+        if self.event_buffer.get(task.key) == (TaskInstanceState.QUEUED, 
self.scheduler_job_id):
+            self.event_buffer.pop(task.key, None)
+        Stats.incr("kubernetes_executor.stale_workload_dropped")

Review Comment:
   I think the logs might be enough at this moment



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