-
Notifications
You must be signed in to change notification settings - Fork 29.3k
[SPARK-28042][K8S] Support using volume mount as local storage #24879
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from 8 commits
6fca505
ef66c87
5610fe4
f071c2c
45c8bc5
9392dad
6e5fcf6
2abb8e9
c29dd7a
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -18,7 +18,8 @@ package org.apache.spark.deploy.k8s.features | |
|
|
||
| import java.util.UUID | ||
|
|
||
| import io.fabric8.kubernetes.api.model.{ContainerBuilder, HasMetadata, PodBuilder, VolumeBuilder, VolumeMountBuilder} | ||
| import io.fabric8.kubernetes.api.model._ | ||
| import scala.collection.JavaConverters._ | ||
|
|
||
| import org.apache.spark.deploy.k8s.{KubernetesConf, SparkPod} | ||
| import org.apache.spark.deploy.k8s.Config._ | ||
|
|
@@ -28,36 +29,45 @@ private[spark] class LocalDirsFeatureStep( | |
| defaultLocalDir: String = s"/var/data/spark-${UUID.randomUUID}") | ||
| extends KubernetesFeatureConfigStep { | ||
|
|
||
| // Cannot use Utils.getConfiguredLocalDirs because that will default to the Java system | ||
| // property - we want to instead default to mounting an emptydir volume that doesn't already | ||
| // exist in the image. | ||
| // We could make utils.getConfiguredLocalDirs opinionated about Kubernetes, as it is already | ||
| // a bit opinionated about YARN and Mesos. | ||
| private val resolvedLocalDirs = Option(conf.sparkConf.getenv("SPARK_LOCAL_DIRS")) | ||
| .orElse(conf.getOption("spark.local.dir")) | ||
| .getOrElse(defaultLocalDir) | ||
| .split(",") | ||
| private val useLocalDirTmpFs = conf.get(KUBERNETES_LOCAL_DIRS_TMPFS) | ||
|
|
||
| override def configurePod(pod: SparkPod): SparkPod = { | ||
| val localDirVolumes = resolvedLocalDirs | ||
| .zipWithIndex | ||
| .map { case (localDir, index) => | ||
| new VolumeBuilder() | ||
| .withName(s"spark-local-dir-${index + 1}") | ||
| .withNewEmptyDir() | ||
| .withMedium(if (useLocalDirTmpFs) "Memory" else null) | ||
| .endEmptyDir() | ||
| .build() | ||
| } | ||
| val localDirVolumeMounts = localDirVolumes | ||
| .zip(resolvedLocalDirs) | ||
| .map { case (localDirVolume, localDirPath) => | ||
| new VolumeMountBuilder() | ||
| .withName(localDirVolume.getName) | ||
| .withMountPath(localDirPath) | ||
| .build() | ||
| } | ||
| var localDirs = findLocalDirVolumeMount(pod) | ||
| var localDirVolumes : Seq[Volume] = Seq() | ||
| var localDirVolumeMounts : Seq[VolumeMount] = Seq() | ||
|
|
||
| if (localDirs.isEmpty) { | ||
| // Cannot use Utils.getConfiguredLocalDirs because that will default to the Java system | ||
| // property - we want to instead default to mounting an emptydir volume that doesn't already | ||
| // exist in the image. | ||
| // We could make utils.getConfiguredLocalDirs opinionated about Kubernetes, as it is already | ||
| // a bit opinionated about YARN and Mesos. | ||
| val resolvedLocalDirs = Option(conf.sparkConf.getenv("SPARK_LOCAL_DIRS")) | ||
| .orElse(conf.getOption("spark.local.dir")) | ||
| .getOrElse(defaultLocalDir) | ||
| .split(",") | ||
| localDirs = resolvedLocalDirs.toSeq | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Get rid of |
||
| localDirVolumes = resolvedLocalDirs | ||
| .zipWithIndex | ||
| .map { case (_, index) => | ||
| new VolumeBuilder() | ||
| .withName(s"spark-local-dir-${index + 1}") | ||
| .withNewEmptyDir() | ||
| .withMedium(if (useLocalDirTmpFs) "Memory" else null) | ||
| .endEmptyDir() | ||
| .build() | ||
| } | ||
|
|
||
| localDirVolumeMounts = localDirVolumes | ||
| .zip(resolvedLocalDirs) | ||
| .map { case (localDirVolume, localDirPath) => | ||
| new VolumeMountBuilder() | ||
| .withName(localDirVolume.getName) | ||
| .withMountPath(localDirPath) | ||
| .build() | ||
| } | ||
| } | ||
|
|
||
| val podWithLocalDirVolumes = new PodBuilder(pod.pod) | ||
| .editSpec() | ||
| .addToVolumes(localDirVolumes: _*) | ||
|
|
@@ -66,10 +76,19 @@ private[spark] class LocalDirsFeatureStep( | |
| val containerWithLocalDirVolumeMounts = new ContainerBuilder(pod.container) | ||
| .addNewEnv() | ||
| .withName("SPARK_LOCAL_DIRS") | ||
| .withValue(resolvedLocalDirs.mkString(",")) | ||
| .withValue(localDirs.mkString(",")) | ||
| .endEnv() | ||
| .addToVolumeMounts(localDirVolumeMounts: _*) | ||
| .build() | ||
| SparkPod(podWithLocalDirVolumes, containerWithLocalDirVolumeMounts) | ||
| } | ||
|
|
||
| def findLocalDirVolumeMount(pod: SparkPod): Seq[String] = { | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Mount => Mounts |
||
| val localDirVolumes = pod.pod.getSpec.getVolumes.asScala | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Hmm... do you need to list the volumes at all? Can't you just look at the mounts (since they already have the volume name)? |
||
| .filter(_.getName.startsWith("spark-local-dir-")) | ||
|
|
||
| localDirVolumes.flatMap { volume => | ||
| pod.container.getVolumeMounts.asScala | ||
| .filter(_.getName == volume.getName).map(_.getMountPath) } | ||
| } | ||
| } | ||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -43,12 +43,12 @@ private[spark] class KubernetesDriverBuilder { | |
| new DriverServiceFeatureStep(conf), | ||
| new MountSecretsFeatureStep(conf), | ||
| new EnvSecretsFeatureStep(conf), | ||
| new LocalDirsFeatureStep(conf), | ||
| new MountVolumesFeatureStep(conf), | ||
| new DriverCommandFeatureStep(conf), | ||
| new HadoopConfDriverFeatureStep(conf), | ||
| new KerberosConfDriverFeatureStep(conf), | ||
| new PodTemplateConfigMapStep(conf)) | ||
| new PodTemplateConfigMapStep(conf), | ||
| new LocalDirsFeatureStep(conf)) | ||
|
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. why not move this right after
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. I think some volume mounts maybe setup in PodTemplateConfigMapStep, so I move to last to ensure that.
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Yep, since the intent is for users to provide custom local dir volumes via either mount volumes or pod templates it needs to appear after both
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Either is fine, just wanted to point out that the pod template is initialized before any step is executed (i.e. |
||
|
|
||
| val spec = KubernetesDriverSpec( | ||
| initialPod, | ||
|
|
||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Hmm, style checker should have complained about this (scala imports should be separate from others).