potiuk commented on code in PR #73044:
URL: https://github.com/apache/airflow/pull/73044#discussion_r4188865524


##########
providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/triggers/pod.py:
##########
@@ -407,10 +408,25 @@ async def _wait_for_container_completion(self) -> 
TriggerEvent:
             now = datetime.datetime.now(tz=datetime.timezone.utc)
             if time_get_more_logs and now >= time_get_more_logs:
                 if self.get_logs and self.logging_interval:
-                    self.last_log_time = await 
self.pod_manager.fetch_container_logs_before_current_sec(
-                        pod, container_name=self.base_container_name, 
since_time=self.last_log_time
-                    )
+                    # Advance before fetching so a failed read waits a full 
interval rather than
+                    # being retried on the next poll.
                     time_get_more_logs = now + 
datetime.timedelta(seconds=self.logging_interval)
+                    try:
+                        self.last_log_time = await 
self.pod_manager.fetch_container_logs_before_current_sec(
+                            pod,
+                            container_name=self.base_container_name,
+                            since_time=self.last_log_time,
+                        )
+                    except ClientError as e:

Review Comment:
   A side effect worth noting: lines are now logged as they stream, so if the 
read fails partway, the lines already logged stay in the task log, 
`last_log_time` doesn't advance, and the next interval reads them again. 
Duplicates are a better failure mode than losing logs, but could you add a 
comment, or advance `last_log_time` to the timestamp of the last emitted line?



##########
providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/utils/pod_manager.py:
##########
@@ -1234,29 +1234,34 @@ async def fetch_container_logs_before_current_sec(
         """
         Asynchronously read the log file of the specified pod.
 
-        This method streams logs from the base container, skipping log lines 
from the current second to prevent duplicate entries on subsequent reads. It is 
designed to handle long-running containers and gracefully suppresses transient 
interruptions.
+        Lines are consumed one at a time as they arrive, so a large backlog is 
never held in memory.
+        Log lines from the current second are skipped to prevent duplicate 
entries on the next read.
 
         :param pod: The pod specification to monitor.
         :param container_name: The name of the container within the pod.
         :param since_time: The timestamp from which to start reading logs.
-        :return: The timestamp to use for the next log read, representing the 
start of the current second. Returns None if an exception occurred.
+        :return: The timestamp to use for the next log read, representing the 
start of the current second.
         """
         now = pendulum.now()
-        logs = await self._hook.read_logs(
+        log_lines = self._hook.stream_logs(
             name=pod.metadata.name,
             namespace=pod.metadata.namespace,
             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)
+        # ``aclosing`` so the early break below finalises the generator, and 
with it the HTTP
+        # response, deterministically rather than leaving it to the loop's 
asyncgen finaliser.
+        async with aclosing(log_lines) as lines:

Review Comment:
   **Blocking until a logging-path review:** #69661 moved the parse-and-log 
step into a thread so it couldn't block the triggerer loop. This moves it back 
onto the loop, yielding every 256 lines. In Airflow 3, trigger logs go through 
the supervisor channel, so the per-line cost of `self.log.log(...)` is part of 
the stall, not just the parsing. Was the 41 ms stall in the benchmark measured 
with the triggerer's real log handler? If emitting stays on the loop, its cost 
should be measured with that handler; the other option is to keep streaming but 
hand batches of lines to a thread.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to