This is an automated email from the ASF dual-hosted git repository.

jason810496 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 19cc828aed8 Add running_pod_log_lines config option to 
KubernetesExecutor (#69301)
19cc828aed8 is described below

commit 19cc828aed8359057743d2bc9431f734a6bb2006
Author: Jason(Zhe-You) Liu <[email protected]>
AuthorDate: Mon Jul 20 16:43:18 2026 +0800

    Add running_pod_log_lines config option to KubernetesExecutor (#69301)
    
    * Resolve conflict in KubernetesExecutor.__init__ and drop redundant comment
    
    The running_pod_log_lines config addition conflicted with the
    pod-launch-failure requeue state added to __init__ after this branch
    diverged; keep both initializations. Also drops a comment that only
    restated the line below it, per review feedback.
    
    * Clarify fallback reference in KubernetesExecutor running_pod_log_lines
    
    self.RUNNING_POD_LOG_LINES looked like a self-referential fallback since
    it isn't set on the instance until this assignment completes; referencing
    the class attribute directly makes clear it falls back to the class default.
---
 providers/cncf/kubernetes/provider.yaml                   |  9 +++++++++
 .../cncf/kubernetes/executors/kubernetes_executor.py      |  8 ++++++++
 .../providers/cncf/kubernetes/get_provider_info.py        |  7 +++++++
 .../cncf/kubernetes/executors/test_kubernetes_executor.py | 15 +++++++++++++++
 4 files changed, 39 insertions(+)

diff --git a/providers/cncf/kubernetes/provider.yaml 
b/providers/cncf/kubernetes/provider.yaml
index 047f6b410c0..9e211647790 100644
--- a/providers/cncf/kubernetes/provider.yaml
+++ b/providers/cncf/kubernetes/provider.yaml
@@ -287,6 +287,15 @@ config:
         type: boolean
         example: ~
         default: "False"
+      running_pod_log_lines:
+        description: |
+          Number of lines read from the end of a running task's pod log when 
the task log
+          is served through the kube API, e.g. when viewing logs of a running 
task in the UI.
+          The value must be greater than 0.
+        version_added: 10.20.0
+        type: integer
+        example: ~
+        default: "100"
       pod_template_file:
         description: |
           Path to the YAML pod file that forms the basis for 
KubernetesExecutor workers.
diff --git 
a/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py
 
b/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py
index ba0355504b1..67616ea5989 100644
--- 
a/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py
+++ 
b/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py
@@ -154,6 +154,14 @@ class KubernetesExecutor(BaseExecutor):
         # instead of requeuing. The orphaned task instance itself is still 
recovered by the
         # scheduler's adopt_or_reset_orphaned_tasks(), which re-queues it with 
a fresh attempt.
         self.pod_launch_attempts: dict[TaskInstanceKey, _PodLaunchAttempt] = {}
+        self.RUNNING_POD_LOG_LINES = self.conf.getint(
+            "kubernetes_executor", "running_pod_log_lines", 
fallback=KubernetesExecutor.RUNNING_POD_LOG_LINES
+        )
+        if self.RUNNING_POD_LOG_LINES <= 0:
+            raise ValueError(
+                "The [kubernetes_executor] running_pod_log_lines configuration 
must be greater than 0, "
+                f"got {self.RUNNING_POD_LOG_LINES}."
+            )
         self.completed: dict[tuple[str, str], KubernetesResults] = {}
         self.create_pods_after: datetime | None = None
 
diff --git 
a/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/get_provider_info.py
 
b/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/get_provider_info.py
index 34c9857a315..df11e510921 100644
--- 
a/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/get_provider_info.py
+++ 
b/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/get_provider_info.py
@@ -158,6 +158,13 @@ def get_provider_info():
                         "example": None,
                         "default": "False",
                     },
+                    "running_pod_log_lines": {
+                        "description": "Number of lines read from the end of a 
running task's pod log when the task log\nis served through the kube API, e.g. 
when viewing logs of a running task in the UI.\nThe value must be greater than 
0.\n",
+                        "version_added": "10.20.0",
+                        "type": "integer",
+                        "example": None,
+                        "default": "100",
+                    },
                     "pod_template_file": {
                         "description": "Path to the YAML pod file that forms 
the basis for KubernetesExecutor workers.\n",
                         "version_added": None,
diff --git 
a/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/executors/test_kubernetes_executor.py
 
b/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/executors/test_kubernetes_executor.py
index e5716e30ffc..a71101dc555 100644
--- 
a/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/executors/test_kubernetes_executor.py
+++ 
b/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/executors/test_kubernetes_executor.py
@@ -641,6 +641,21 @@ class TestAirflowKubernetesScheduler:
         assert kube_executor.RUNNING_POD_LOG_LINES == 100
         assert kube_executor_2.RUNNING_POD_LOG_LINES == 200
 
+    @conf_vars({("kubernetes_executor", "running_pod_log_lines"): "500"})
+    def test_running_pod_log_lines_from_config(self):
+        kube_executor = KubernetesExecutor()
+
+        assert kube_executor.RUNNING_POD_LOG_LINES == 500
+        assert KubernetesExecutor.RUNNING_POD_LOG_LINES == 100
+
+    @pytest.mark.parametrize("invalid_value", ["0", "-1"])
+    def test_running_pod_log_lines_invalid_config(self, invalid_value):
+        with conf_vars({("kubernetes_executor", "running_pod_log_lines"): 
invalid_value}):
+            with pytest.raises(
+                ValueError, match="running_pod_log_lines configuration must be 
greater than 0"
+            ):
+                KubernetesExecutor()
+
 
 class TestKubernetesExecutor:
     """

Reply via email to