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 a6699a40ac2 Prevent deferrable KubernetesPodOperator log parsing from
blocking the triggerer event loop (#69661)
a6699a40ac2 is described below
commit a6699a40ac254c4b34b476c694ffde4c2ec8f92a
Author: Jorge Rocamora <[email protected]>
AuthorDate: Fri Jul 31 21:21:37 2026 +0200
Prevent deferrable KubernetesPodOperator log parsing from blocking the
triggerer event loop (#69661)
---
.../providers/cncf/kubernetes/hooks/kubernetes.py | 9 +++++---
.../providers/cncf/kubernetes/utils/pod_manager.py | 6 +++++-
.../unit/cncf/kubernetes/hooks/test_kubernetes.py | 24 ++++++++++++++++++++++
.../unit/cncf/kubernetes/utils/test_pod_manager.py | 20 ++++++++++++++++++
4 files changed, 55 insertions(+), 4 deletions(-)
diff --git
a/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/hooks/kubernetes.py
b/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/hooks/kubernetes.py
index 8ebf00c69e4..bccf6093ea0 100644
---
a/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/hooks/kubernetes.py
+++
b/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/hooks/kubernetes.py
@@ -814,6 +814,10 @@ def _get_bool(val) -> bool | None:
return None
+def _split_log_bytes(raw_bytes: bytes) -> list[str]:
+ return raw_bytes.decode("utf-8", errors="replace").splitlines()
+
+
class AsyncKubernetesHook(KubernetesHook):
"""Hook to use Kubernetes SDK asynchronously."""
@@ -1111,9 +1115,8 @@ class AsyncKubernetesHook(KubernetesHook):
raw_resp: ClientResponse = await
v1_api.read_namespaced_pod_log(**kwargs) # type: ignore #
_preload_content=False makes returning ClientResponse instead of str!
raw_bytes = await raw_resp.read()
- logs = raw_bytes.decode("utf-8", errors="replace")
- logs_list: list[str] = logs.splitlines()
- return logs_list
+ # CPU-bound decode/split, offloaded so it can't block the
triggerer event loop.
+ return await asyncio.to_thread(_split_log_bytes, raw_bytes)
except HTTPError as e:
raise KubernetesApiError from e
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 ec7d339f377..926e59aab1e 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
@@ -1231,6 +1231,11 @@ class AsyncPodManager(LoggingMixin):
container_name=container_name,
since_seconds=(math.ceil((now - since_time).total_seconds()) if
since_time else None),
)
+ # CPU-bound per-line parse/emit, offloaded so it can't block the
triggerer event loop.
+ await asyncio.to_thread(self._emit_container_logs, logs, now,
container_name)
+ return now # Return the current time as the last log time to ensure
logs from the current second are read in the next fetch.
+
+ def _emit_container_logs(self, logs: list[str], now: DateTime,
container_name: str) -> None:
message_to_log = None
try:
now_seconds = now.replace(microsecond=0)
@@ -1266,4 +1271,3 @@ class AsyncPodManager(LoggingMixin):
else:
level = _parse_log_level(message_to_log)
self.log.log(level, "[%s] %s", container_name,
message_to_log)
- return now # Return the current time as the last log time to ensure
logs from the current second are read in the next fetch.
diff --git
a/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/hooks/test_kubernetes.py
b/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/hooks/test_kubernetes.py
index d22dc067868..7105cfad845 100644
---
a/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/hooks/test_kubernetes.py
+++
b/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/hooks/test_kubernetes.py
@@ -39,6 +39,7 @@ from airflow.models import Connection
from airflow.providers.cncf.kubernetes.hooks.kubernetes import (
AsyncKubernetesHook,
KubernetesHook,
+ _split_log_bytes,
_TimeoutAsyncK8sApiClient,
_TimeoutK8sApiClient,
)
@@ -1857,6 +1858,29 @@ class TestAsyncKubernetesHook:
lib_method.assert_called_once()
assert lib_method.call_args.kwargs.get("_preload_content") is False
+ @pytest.mark.asyncio
+ @mock.patch("asyncio.to_thread", new_callable=mock.AsyncMock)
+ @mock.patch(KUBE_API.format("read_namespaced_pod_log"))
+ async def test_read_logs_decodes_off_the_event_loop(self, lib_method,
mock_to_thread, kube_config_loader):
+ """The CPU-bound decode/splitlines is offloaded to a worker thread,
not run on the loop."""
+ raw_bytes = b"2023-01-11 Some string logs..."
+ mock_raw_resp = mock.AsyncMock()
+ mock_raw_resp.read = mock.AsyncMock(return_value=raw_bytes)
+ lib_method.return_value = self.mock_await_result(mock_raw_resp)
+ mock_to_thread.return_value = ["decoded line"]
+
+ hook = AsyncKubernetesHook(
+ conn_id=None,
+ in_cluster=False,
+ config_file=None,
+ cluster_context=None,
+ )
+
+ logs = await hook.read_logs(name=POD_NAME, namespace=NAMESPACE,
container_name=CONTAINER_NAME)
+
+ assert logs == ["decoded line"]
+ mock_to_thread.assert_awaited_once_with(_split_log_bytes, raw_bytes)
+
@pytest.mark.asyncio
@mock.patch(KUBE_BATCH_API.format("read_namespaced_job_status"))
async def test_get_job_status(self, lib_method, kube_config_loader):
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 97f2aff6056..f4c21683762 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
@@ -1788,6 +1788,26 @@ class TestAsyncPodManager:
pod=pod, container_name=container_name, since_time=since_time
)
+ @pytest.mark.asyncio
+ @mock.patch("asyncio.to_thread", new_callable=mock.AsyncMock)
+ async def
test_fetch_container_logs_offloads_parse_off_the_event_loop(self,
mock_to_thread):
+ """The CPU-bound per-line parse/emit loop is offloaded to a worker
thread, not run on the loop."""
+ now = pendulum.datetime(2024, 1, 1, 12, 0, 0)
+ pod = mock.MagicMock()
+ container_name = "base"
+ log_lines = [f"{now.subtract(seconds=2).to_iso8601_string()} hello"]
+ self.mock_async_hook.read_logs.return_value = log_lines
+
+ with
mock.patch("airflow.providers.cncf.kubernetes.utils.pod_manager.pendulum.now",
return_value=now):
+ result = await
self.async_pod_manager.fetch_container_logs_before_current_sec(
+ pod=pod, container_name=container_name,
since_time=now.subtract(minutes=1)
+ )
+
+ assert result == now
+ mock_to_thread.assert_awaited_once_with(
+ self.async_pod_manager._emit_container_logs, log_lines, now,
container_name
+ )
+
class TestPodLogsConsumer:
@pytest.mark.parametrize(