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]
