sunchao commented on code in PR #57350:
URL: https://github.com/apache/spark/pull/57350#discussion_r3706593855


##########
resource-managers/kubernetes/core/src/main/scala/org/apache/spark/deploy/k8s/features/DriverUIServiceFeatureStep.scala:
##########
@@ -0,0 +1,108 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *    http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.spark.deploy.k8s.features
+
+import scala.jdk.CollectionConverters._
+
+import io.fabric8.kubernetes.api.model.{HasMetadata, ServiceBuilder}
+
+import org.apache.spark.deploy.k8s.{KubernetesDriverConf, SparkPod}
+import org.apache.spark.deploy.k8s.Config.{
+  KUBERNETES_DRIVER_UI_SERVICE_ENABLED,
+  KUBERNETES_DRIVER_UI_SERVICE_NAME,
+  KUBERNETES_DRIVER_UI_SERVICE_TYPE
+}
+import org.apache.spark.deploy.k8s.Constants._
+import org.apache.spark.internal.{config, Logging}
+
+/**
+ * Optionally provisions a dedicated Kubernetes Service exposing only the 
Spark driver's Web UI
+ * port.
+ *
+ * At creation time, the Service's `targetPort` is set to the configured 
`spark.ui.port`
+ * (a placeholder when the user has requested a random port). Once the 
driver's Jetty server
+ * has bound, `K8sDriverUIServicePatcher` updates the Service's `targetPort` 
to reflect the
+ * actual bound port.
+ */
+private[spark] class DriverUIServiceFeatureStep(kubernetesConf: 
KubernetesDriverConf)
+  extends KubernetesFeatureConfigStep with Logging {
+  import DriverUIServiceFeatureStep._
+
+  private val enabled = 
kubernetesConf.get(KUBERNETES_DRIVER_UI_SERVICE_ENABLED)
+  private lazy val serviceType = 
kubernetesConf.get(KUBERNETES_DRIVER_UI_SERVICE_TYPE)
+  private lazy val configuredUIPort = kubernetesConf.get(config.UI.UI_PORT)
+
+  /**
+   * Port value used when building the Service. When the user has requested a 
random UI port
+   * (`spark.ui.port=0`), the actual port is only known after the driver's 
Jetty server binds,
+   * so we substitute the default UI port (typically 4040) purely as a 
placeholder to satisfy
+   * Kubernetes' Service port validation (must be > 0). After the driver JVM 
starts,
+   * [[org.apache.spark.scheduler.cluster.k8s.K8sDriverUIServicePatcher]] 
updates the Service's
+   * `targetPort` to the real bound port.
+   */
+  private lazy val servicePort: Int = if (configuredUIPort == 0) {

Review Comment:
   [P1] This remains reproducible on 
`44b952a19a85a35e9ecad237dde9cfb604183e9b`. 
`BasicDriverFeatureStep.scala:100,117-120` still copies `spark.ui.port=0` into 
the driver's `containerPort`, and `DriverServiceFeatureStep.scala:52,89-93` 
independently emits Service `port=0` and `targetPort=0`. The proposed 
`spark.kubernetes.executor.useDriverPodIP=true` plus excluding 
`DriverServiceFeatureStep` only removes the second invalid resource; 
`BasicDriverFeatureStep` still makes the Pod invalid before the runtime patcher 
can execute. Because this PR explicitly promises `spark.ui.port=0`, both 
existing declarations need to be omitted/replaced, with a full 
`KubernetesDriverBuilder` regression test.



##########
resource-managers/kubernetes/core/src/main/scala/org/apache/spark/scheduler/cluster/k8s/KubernetesClusterSchedulerBackend.scala:
##########
@@ -120,6 +121,27 @@ private[spark] class KubernetesClusterSchedulerBackend(
     if (!conf.get(KUBERNETES_EXECUTOR_DISABLE_CONFIGMAP)) {
       setUpExecutorConfigMap(podAllocator.driverPod)
     }
+    maybePatchDriverUIServiceTargetPort()

Review Comment:
   [P2] This still fails on `44b952a19a85a35e9ecad237dde9cfb604183e9b`, and the 
latest endpointless-Service change makes it affect even a normal fixed port 
with no bind collision. For `spark.kubernetes.driver.master=local[*]`, 
`KubernetesClusterManager.scala:56-72` returns `LocalSchedulerBackend`, but 
only `KubernetesClusterSchedulerBackend.start()` invokes the patcher. 
`DriverUIServiceFeatureStep` now creates every enabled UI Service without a 
selector, so `spark.ui.port=4040` plus 
`spark.kubernetes.driver.ui.service.enabled=true` leaves the Service 
permanently endpointless. Setting the feature to false avoids the configuration 
rather than handling the explicitly enabled one; please patch from a shared 
driver lifecycle or reject/disable this incompatible combination.



##########
resource-managers/kubernetes/core/src/main/scala/org/apache/spark/deploy/k8s/features/DriverUIServiceFeatureStep.scala:
##########
@@ -0,0 +1,167 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *    http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.spark.deploy.k8s.features
+
+import scala.jdk.CollectionConverters._
+
+import io.fabric8.kubernetes.api.model.{HasMetadata, ServiceBuilder}
+
+import org.apache.spark.deploy.k8s.{KubernetesDriverConf, SparkPod}
+import org.apache.spark.deploy.k8s.Config.{
+  KUBERNETES_DRIVER_SERVICE_IP_FAMILIES,
+  KUBERNETES_DRIVER_SERVICE_IP_FAMILY_POLICY,
+  KUBERNETES_DRIVER_UI_SERVICE_ENABLED,
+  KUBERNETES_DRIVER_UI_SERVICE_NAME,
+  KUBERNETES_DRIVER_UI_SERVICE_TYPE
+}
+import org.apache.spark.deploy.k8s.Constants._
+import org.apache.spark.internal.{config, Logging}
+
+/**
+ * Optionally provisions a dedicated Kubernetes Service exposing only the 
Spark driver's Web UI
+ * port.
+ *
+ * When the user requests a random UI port (`spark.ui.port=0`), the actual 
port is unknown at
+ * submission time, so the Service's `targetPort` uses a placeholder. To avoid 
routing traffic
+ * to the wrong endpoint during that window (in `hostNetwork` mode a pod 
endpoint is the node IP,
+ * so the placeholder `targetPort` could reach an unrelated driver co-located 
on the same node),
+ * the Service is created *without* a selector, leaving it endpointless. Once 
the driver's Jetty
+ * server has bound, `K8sDriverUIServicePatcher` patches the selector and the 
actual `targetPort`
+ * together. When the port is fixed up front, the Service is created with its 
selector and final
+ * `targetPort` immediately.
+ */
+private[spark] class DriverUIServiceFeatureStep(kubernetesConf: 
KubernetesDriverConf)
+  extends KubernetesFeatureConfigStep with Logging {
+  import DriverUIServiceFeatureStep._
+
+  private val enabled = 
kubernetesConf.get(KUBERNETES_DRIVER_UI_SERVICE_ENABLED)
+  private lazy val serviceType = 
kubernetesConf.get(KUBERNETES_DRIVER_UI_SERVICE_TYPE)
+  private lazy val configuredUIPort = kubernetesConf.get(config.UI.UI_PORT)
+
+  /**
+   * Port value used when building the Service. When the user has requested a 
random UI port
+   * (`spark.ui.port=0`), the actual port is only known after the driver's 
Jetty server binds,
+   * so we substitute the default UI port (typically 4040) purely as a 
placeholder to satisfy
+   * Kubernetes' Service port validation (must be > 0). After the driver JVM 
starts,
+   * [[org.apache.spark.scheduler.cluster.k8s.K8sDriverUIServicePatcher]] 
updates the Service's
+   * `targetPort` to the real bound port.
+   */
+  private lazy val servicePort: Int = if (configuredUIPort == 0) {
+    config.UI.UI_PORT.defaultValue.get
+  } else {
+    configuredUIPort
+  }
+
+  private lazy val serviceName: String = 
kubernetesConf.get(KUBERNETES_DRIVER_UI_SERVICE_NAME)

Review Comment:
   [P1] This collision surface is introduced by the new arbitrary 
`spark.kubernetes.driver.ui.service.name` override. The existing mandatory 
driver Service derives its name from the application-specific randomized 
`resourceNamePrefix`; existing driver Pods use create-only semantics. By 
contrast, two otherwise distinct applications can now both set `shared-ui`: 
`KubernetesUtils.addOwnerReference` replaces the controller owner UID, and 
`KubernetesClientApplication.scala:175-177` uses 
`forceConflicts().serverSideApply()` on that same Service; the second 
application's runtime patch then replaces the selector. The first application 
silently loses its Service, and the Service can be garbage-collected with the 
wrong Pod. Choosing this application's own mandatory driver-Service name also 
makes both feature steps emit the same Service identity. Please reject name 
collisions or verify existing ownership before applying; unique defaults do not 
protect explicitly configured names.



-- 
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]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to