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(

Reply via email to