apache / apache/airflow

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

Open
#73,117 1 comment 1 reaction 0 assignees View on GitHub
Dominant language
Python
Stars
46.9k
Forks
17.8k
Avg merge
2d 9h
Merged PRs (30d)
472

Description

### 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 bda472db1be86872534a17a603ad2ccad04814a1).
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 / 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):

```python
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:

```python
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: #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?

- [x] 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

- [x] I agree to follow this project's [Code of Conduct](https://github.com/apache/airflow/blob/main/CODE_OF_CONDUCT.md)

Contributor guide

Open the contributing guide

Research direction

Start in operators/pod.py at trigger_reentry, especially the 404 handling around lines 989–1007, and inspect the existing KubernetesPodOperator behavior for deferrable tasks. Reproduce the case where a successfully completed pod is gone before re-entry with do_xcom_push enabled. Done means the task fails retryably instead of reporting success without its return_value XCom, while the no-XCom path retains its intended behavior.

Written by the indexing model from the issue text.

Assessment

Tech stack
kubernetes, python
Domain
backend, devops
Issue type
Bug
Difficulty
3/5
Estimated time
1-2 days
Activity status
Active
Clarity
Clearly specified
Newbie friendliness
76/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.