This is an automated email from the ASF dual-hosted git repository.
potiuk pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/airflow.git
The following commit(s) were added to refs/heads/main by this push:
new ee5eefcb9fd Fix KubernetesPodOperator discarding successful XCom when
sidecar kill fails (#72068)
ee5eefcb9fd is described below
commit ee5eefcb9fd53477b79bd4d6116a53c459d4f11f
Author: Jeremy Schoemaker <[email protected]>
AuthorDate: Mon Sep 21 09:30:07 2026 -0500
Fix KubernetesPodOperator discarding successful XCom when sidecar kill
fails (#72068)
closes: #71369
---
.../providers/cncf/kubernetes/operators/job.py | 6 ++-
.../providers/cncf/kubernetes/operators/pod.py | 12 ++++--
.../providers/cncf/kubernetes/utils/pod_manager.py | 24 +++++++++--
.../unit/cncf/kubernetes/operators/test_job.py | 47 +++++++++++++++++++++-
.../unit/cncf/kubernetes/utils/test_pod_manager.py | 41 +++++++++++++++++++
5 files changed, 122 insertions(+), 8 deletions(-)
diff --git
a/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/operators/job.py
b/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/operators/job.py
index e888159c739..21fabd19c01 100644
---
a/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/operators/job.py
+++
b/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/operators/job.py
@@ -238,7 +238,11 @@ class KubernetesJobOperator(KubernetesPodOperator):
pod=pod, container_name=self.base_container_name
)
self.pod_manager.await_xcom_sidecar_container_start(pod=pod)
- xcom_result.append(self.extract_xcom(pod=pod))
+ # `wait_until_job_complete` below polls without a
timeout, and the
+ # Job can never reach a terminal state while the xcom
sidecar is
+ # still looping. So unlike KubernetesPodOperator, this
path must
+ # fail loudly when the sidecar cannot be killed rather
than hang.
+ xcom_result.append(self.extract_xcom(pod=pod,
ignore_kill_failure=False))
self.job = self.hook.wait_until_job_complete(
job_name=self.job.metadata.name,
namespace=self.job.metadata.namespace,
diff --git
a/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/operators/pod.py
b/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/operators/pod.py
index 2436fb7bb5d..5178222c647 100644
---
a/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/operators/pod.py
+++
b/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/operators/pod.py
@@ -810,9 +810,15 @@ class KubernetesPodOperator(BaseOperator):
self._read_pod_events(pod, reraise=False)
raise
- def extract_xcom(self, pod: k8s.V1Pod) -> dict[Any, Any] | None:
- """Retrieve xcom value and kill xcom sidecar container."""
- result = self.pod_manager.extract_xcom(pod)
+ def extract_xcom(self, pod: k8s.V1Pod, *, ignore_kill_failure: bool =
True) -> dict[Any, Any] | None:
+ """
+ Retrieve xcom value and kill xcom sidecar container.
+
+ :param pod: the pod to read the XCom result from.
+ :param ignore_kill_failure: when True (the default), a failure to kill
the sidecar
+ container does not discard the XCom value that was already read.
+ """
+ result = self.pod_manager.extract_xcom(pod,
ignore_kill_failure=ignore_kill_failure)
if isinstance(result, str) and result.rstrip() == EMPTY_XCOM_RESULT:
self.log.info("xcom result file is empty.")
return None
diff --git
a/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/utils/pod_manager.py
b/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/utils/pod_manager.py
index 0914de4e90a..96588fe8334 100644
---
a/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/utils/pod_manager.py
+++
b/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/utils/pod_manager.py
@@ -974,8 +974,17 @@ class PodManager(LoggingMixin):
)
time.sleep(1)
- def extract_xcom(self, pod: V1Pod) -> str:
- """Retrieve XCom value and kill xcom sidecar container."""
+ def extract_xcom(self, pod: V1Pod, *, ignore_kill_failure: bool = True) ->
str:
+ """
+ Retrieve XCom value and kill xcom sidecar container.
+
+ :param pod: the pod to read the XCom result from.
+ :param ignore_kill_failure: when True (the default), a failure to kill
the sidecar
+ container is logged as a warning and the successfully read XCom
value is still
+ returned. Set it to False when the caller cannot tolerate a
sidecar that keeps
+ running, for example ``KubernetesJobOperator``, whose Job can
never reach a
+ terminal state while the sidecar is alive.
+ """
# make sure that xcom sidecar container is still running
if not self.container_is_running(pod,
PodDefaults.SIDECAR_CONTAINER_NAME):
raise XComRetrievalError(
@@ -986,7 +995,16 @@ class PodManager(LoggingMixin):
result = self.extract_xcom_json(pod)
return result
finally:
- self.extract_xcom_kill(pod)
+ try:
+ self.extract_xcom_kill(pod)
+ except (PodCommandException, ApiException) as e:
+ if not ignore_kill_failure:
+ raise
+ self.log.warning(
+ "Failed to kill xcom sidecar container in pod %s, leaving
it running: %s",
+ pod.metadata.name,
+ e,
+ )
@generic_api_retry
def extract_xcom_json(self, pod: V1Pod) -> str:
diff --git
a/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/operators/test_job.py
b/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/operators/test_job.py
index 60536cbb439..6f2cdd7d425 100644
--- a/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/operators/test_job.py
+++ b/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/operators/test_job.py
@@ -36,7 +36,7 @@ from airflow.providers.cncf.kubernetes.operators.job import (
KubernetesPatchJobOperator,
)
from airflow.providers.cncf.kubernetes.triggers.job import KubernetesJobTrigger
-from airflow.providers.cncf.kubernetes.utils.pod_manager import PodManager
+from airflow.providers.cncf.kubernetes.utils.pod_manager import
PodCommandException, PodManager
from airflow.providers.common.compat.sdk import AirflowException,
TaskDeferred, timezone
from airflow.utils.session import create_session
from airflow.utils.types import DagRunType
@@ -867,6 +867,51 @@ class TestKubernetesJobOperator:
with pytest.raises(AirflowProviderDeprecationWarning):
assert op.pod == mock_pods_expected[0]
+ @patch(f"{POD_MANAGER_CLASS}.extract_xcom_kill")
+ @patch(f"{POD_MANAGER_CLASS}.extract_xcom_json", return_value='{"a":
"true"}')
+ @patch(f"{POD_MANAGER_CLASS}.container_is_running", return_value=True)
+ @patch(f"{POD_MANAGER_CLASS}.await_xcom_sidecar_container_start")
+ @patch(f"{POD_MANAGER_CLASS}.await_container_completion")
+ @patch(JOB_OPERATORS_PATH.format("KubernetesJobOperator.get_pods"))
+
@patch(JOB_OPERATORS_PATH.format("KubernetesJobOperator.build_job_request_obj"))
+ @patch(JOB_OPERATORS_PATH.format("KubernetesJobOperator.create_job"))
+ @patch(f"{HOOK_CLASS}.wait_until_job_complete")
+ def test_xcom_sidecar_kill_failure_fails_job_instead_of_hanging(
+ self,
+ mock_wait_until_job_complete,
+ mock_create_job,
+ mock_build_job_request_obj,
+ mock_get_pods,
+ mock_await_container_completion,
+ mock_await_xcom_sidecar_container_start,
+ mock_container_is_running,
+ mock_extract_xcom_json,
+ mock_extract_xcom_kill,
+ ):
+ """
+ The Job path must not swallow a sidecar kill failure.
+
+ ``wait_until_job_complete`` polls without a timeout and the Job can
never reach a
+ terminal state while the xcom sidecar keeps looping, so the task would
hang. It
+ has to fail loudly before the polling starts instead.
+ """
+ mock_get_pods.return_value = [mock.MagicMock()]
+ mock_extract_xcom_kill.side_effect = PodCommandException(
+ "Command failed with stderr: Permission denied"
+ )
+
+ op = KubernetesJobOperator(
+ task_id="test_task_id",
+ wait_until_job_complete=True,
+ job_poll_interval=POLL_INTERVAL,
+ do_xcom_push=True,
+ )
+
+ with pytest.raises(PodCommandException, match="Permission denied"):
+ op.execute(context=dict(ti=mock.MagicMock()))
+
+ mock_wait_until_job_complete.assert_not_called()
+
@pytest.mark.parametrize("do_xcom_push", [True, False])
@pytest.mark.parametrize("get_logs", [True, False])
@patch(JOB_OPERATORS_PATH.format("KubernetesJobOperator._write_logs"))
diff --git
a/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/utils/test_pod_manager.py
b/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/utils/test_pod_manager.py
index f4c21683762..45e4b20a84f 100644
---
a/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/utils/test_pod_manager.py
+++
b/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/utils/test_pod_manager.py
@@ -32,6 +32,7 @@ from urllib3.exceptions import HTTPError as BaseHTTPError
from airflow.providers.cncf.kubernetes.exceptions import KubernetesApiError
from airflow.providers.cncf.kubernetes.utils.pod_manager import (
AsyncPodManager,
+ PodCommandException,
PodLogsConsumer,
PodManager,
PodPhase,
@@ -1225,6 +1226,46 @@ class TestPodManager:
assert ret == xcom_json
assert mock_exec_xcom_kill.call_count == 1
+
@mock.patch("airflow.providers.cncf.kubernetes.utils.pod_manager.kubernetes_stream")
+
@mock.patch("airflow.providers.cncf.kubernetes.utils.pod_manager.PodManager.extract_xcom_kill")
+ @mock.patch(
+
"airflow.providers.cncf.kubernetes.utils.pod_manager.PodManager.container_is_running",
+ return_value=True,
+ )
+ def test_extract_xcom_returns_result_when_sidecar_kill_fails(
+ self, mock_container_is_running, mock_exec_xcom_kill,
mock_kubernetes_stream
+ ):
+ """A failure to kill the sidecar must not discard the XCom value
already read."""
+ xcom_json = """{"a": "true"}"""
+ mock_client = MagicMock()
+ mock_client.peek_stderr.return_value = ""
+ mock_client.read_all.return_value = xcom_json
+ mock_kubernetes_stream.return_value = mock_client
+ mock_exec_xcom_kill.side_effect = PodCommandException("Command failed
with stderr: Permission denied")
+ ret = self.pod_manager.extract_xcom(pod=MagicMock())
+ assert ret == xcom_json
+ assert mock_exec_xcom_kill.call_count == 1
+
+
@mock.patch("airflow.providers.cncf.kubernetes.utils.pod_manager.kubernetes_stream")
+
@mock.patch("airflow.providers.cncf.kubernetes.utils.pod_manager.PodManager.extract_xcom_kill")
+ @mock.patch(
+
"airflow.providers.cncf.kubernetes.utils.pod_manager.PodManager.container_is_running",
+ return_value=True,
+ )
+ def test_extract_xcom_reraises_kill_failure_when_not_ignored(
+ self, mock_container_is_running, mock_exec_xcom_kill,
mock_kubernetes_stream
+ ):
+ """With ignore_kill_failure=False the kill failure still propagates to
the caller."""
+ xcom_json = """{"a": "true"}"""
+ mock_client = MagicMock()
+ mock_client.peek_stderr.return_value = ""
+ mock_client.read_all.return_value = xcom_json
+ mock_kubernetes_stream.return_value = mock_client
+ mock_exec_xcom_kill.side_effect = PodCommandException("Command failed
with stderr: Permission denied")
+ with pytest.raises(PodCommandException, match="Permission denied"):
+ self.pod_manager.extract_xcom(pod=MagicMock(),
ignore_kill_failure=False)
+ assert mock_exec_xcom_kill.call_count == 1
+
@mock.patch("airflow.providers.cncf.kubernetes.utils.pod_manager.kubernetes_stream")
@mock.patch("airflow.providers.cncf.kubernetes.utils.pod_manager.PodManager.extract_xcom_kill")
@mock.patch(