diff --git a/providers/src/airflow/providers/cncf/kubernetes/utils/pod_manager.py b/providers/src/airflow/providers/cncf/kubernetes/utils/pod_manager.py index e123ef0b0d284..bda798cc928fa 100644 --- a/providers/src/airflow/providers/cncf/kubernetes/utils/pod_manager.py +++ b/providers/src/airflow/providers/cncf/kubernetes/utils/pod_manager.py @@ -451,7 +451,10 @@ def consume_logs(*, since_time: DateTime | None = None) -> tuple[DateTime | None line=line, client=self._client, mode=ExecutionMode.SYNC ) if message_to_log is not None: - self.log.info("[%s] %s", container_name, message_to_log) + if is_log_group_marker(message_to_log): + print(message_to_log) + else: + self.log.info("[%s] %s", container_name, message_to_log) last_captured_timestamp = message_timestamp message_to_log = message message_timestamp = line_timestamp @@ -467,7 +470,10 @@ def consume_logs(*, since_time: DateTime | None = None) -> tuple[DateTime | None line=line, client=self._client, mode=ExecutionMode.SYNC ) if message_to_log is not None: - self.log.info("[%s] %s", container_name, message_to_log) + if is_log_group_marker(message_to_log): + print(message_to_log) + else: + self.log.info("[%s] %s", container_name, message_to_log) last_captured_timestamp = message_timestamp except TimeoutError as e: # in case of timeout, increment return time by 2 seconds to avoid @@ -820,3 +826,8 @@ class OnFinishAction(str, enum.Enum): KEEP_POD = "keep_pod" DELETE_POD = "delete_pod" DELETE_SUCCEEDED_POD = "delete_succeeded_pod" + + +def is_log_group_marker(line: str) -> bool: + """Check if the line is a log group marker like `::group::` or `::endgroup::`.""" + return line.startswith("::group::") or line.startswith("::endgroup::")