Skip to content

Deferrable KubernetesPodOperator marks task SUCCESS without pushing XCom when the pod is gone at re-entry (regression from #66716, provider 10.17.1+) #73117

Description

@aroina

Apache Airflow Provider(s)

cncf-kubernetes

Versions of Apache Airflow Providers

apache-airflow-providers-cncf-kubernetes==10.19.0 (observed in production).
The code path is unchanged in 10.20.0, 10.21.0, 10.21.1 and on main as of today.
Introduced in 10.17.1 by #66716 (merge commit bda472d).
Not present in 10.17.0 and earlier.

Apache Airflow version

3.2.2 (task-sdk 1.2.2)

Operating System

Debian 12 (official apache/airflow:slim-3.2.2-python3.12 image)

Deployment

Official Apache Airflow Helm Chart

Deployment details

  • CeleryExecutor, on-premise Kubernetes 1.3x.
  • KubernetesPodOperator with deferrable=True, do_xcom_push=True, on_finish_action=delete_pod, and active_deadline_seconds set on the pod.
  • Dynamic task mapping: a mapped download task (about 750 map indexes per DAG run) whose return_value XCom is consumed by a mapped process task through a Jinja template.

What happened

Since upgrading to a provider that contains #66716, deferrable KubernetesPodOperator tasks are intermittently marked SUCCESS while pushing no return_value XCom. The downstream task then fails at templating time with:

TypeError: the JSON object must be str, bytes or bytearray, not NoneType

Nothing in the upstream task indicates a problem: state is success, the log ends with a warning and no error, and retries of the downstream task can never recover because the upstream XCom does not exist. In one affected day, 65 of 746 mapped instances of the upstream task were "successful" without a return_value; every one of them produced a downstream failure that had to be repaired by clearing the upstream task by hand.

Sequence, taken from the task log of one instance (timestamps trimmed):

21:57:58  Building pod ... / Pausing task as DEFERRED.
21:58:04  Pod has reached Running phase before launch timeout
21:59:15  Trigger fired event ... TriggerEvent<{'status': 'success', ...}>
          (main container finished; xcom sidecar keeps the pod alive waiting for the worker)
          -- worker does not resume the task for ~13 minutes (scheduler backlog) --
          -- pod hits spec.activeDeadlineSeconds, kubelet kills the sidecar, pod -> Failed/DeadlineExceeded --
          -- pod object is then removed by cluster pod garbage collection --
22:12:01  [warning] Pod <ns>/<pod> not found after resuming from deferral — already GC'd.  (pod.py:996)
22:12:01  ::group::Post Execute  ->  task state SUCCESS, no return_value XCom

The responsible code is the new 404 handling in trigger_reentry (provider 10.19.0, operators/pod.py lines 989–1007):

        try:
            self.pod = self.hook.get_pod(pod_name, pod_namespace)
        except ApiException as e:
            if e.status != 404:
                raise
            self.log.warning(
                "Pod %s/%s not found after resuming from deferral — already GC'd.", ...
            )
            if event["status"] == "success":
                # Trigger already observed the pod completed successfully;
                # logs/XCom are unrecoverable but the task itself succeeded.
                return
            raise PodNotFoundException(...) from e

The comment acknowledges that the XCom is unrecoverable, yet the method returns normally. For an operator with do_xcom_push=True, "the task itself succeeded" is not true: the contract of the task is to produce a value, and the value is lost. Before #66716 the same situation raised ApiException(404) out of trigger_reentry, the task failed, and the retry re-ran the pod and produced a correct XCom. The fix for the crash (#66715) is legitimate, but it turned a loud failure into silent data loss for XCom-producing tasks.

Note that if the pod still exists with a dead sidecar, PodManager.extract_xcom raises XComRetrievalError and the task fails as expected. The silent path only opens when the pod object is gone, which is exactly the case #66716 targets.

What you think should happen instead

When the pod cannot be found at re-entry and do_xcom_push is True, the task must fail (retryable), because its result cannot be produced. The silent-success shortcut should only apply when no XCom is expected, for example:

            if event["status"] == "success" and not self.do_xcom_push:
                return
            raise PodNotFoundException(
                f"Pod {pod_namespace}/{pod_name} not found after resuming from deferral"
                + (" — XCom cannot be retrieved" if self.do_xcom_push else "")
            ) from e

Alternatively, the shortcut could be made opt-in (for example succeed_if_pod_gone=True), keeping the pre-10.17.1 behaviour as the default.

How to reproduce

  1. Cluster where terminated pods are removed (any pod GC, or delete the pod by hand in step 4).
  2. Deferrable KubernetesPodOperator with do_xcom_push=True, an image that writes /airflow/xcom/return.json and exits 0 within a few seconds, and active_deadline_seconds=60. Put the task in a pool.
  3. Downstream task using {{ ti.xcom_pull(task_ids='upstream') }} with json.loads or equivalent.
  4. Trigger the DAG. Once the trigger has fired (Trigger fired event ... status: success), prevent the worker from resuming the task for more than 60 s (for example set the pool slots to 0 while the task is in scheduled state). The kubelet kills the sidecar at the deadline; delete the Failed pod or let pod GC remove it.
  5. Reopen the pool. The upstream task logs not found after resuming from deferral — already GC'd and ends in success with no return_value. The downstream task fails with TypeError ... NoneType.

Provider 10.17.0 or earlier: step 5 ends with the upstream task in failed/up_for_retry and the retry succeeds with a correct XCom.

Anything else

  • The problem is intermittent and load dependent: it appears whenever the delay between the trigger firing and the worker resuming the task exceeds the pod's activeDeadlineSeconds (or any other reason the completed pod disappears). Under scheduler backlog we measured resume delays of up to 58 minutes.
  • Related but distinct: KubernetesPodOperator task is marked SUCCESS when the task process receives SIGTERM mid-execute (pod deleted by on_kill, cleanup() short-circuits on _killed) #73006 (task marked SUCCESS after SIGTERM via on_kill / _killed short-circuit). Both are "false success" paths in KubernetesPodOperator; this one needs no signal.
  • Workaround used on our side: subclass override of trigger_reentry that raises PodNotFoundException when do_xcom_push is set and self.pod is still None after super().trigger_reentry() returns.

Are you willing to submit PR?

  • Yes, I am willing to submit a PR, once maintainers agree on the approach (fail when do_xcom_push, or opt-in shortcut).

Code of Conduct

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    Type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions