-
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 3 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 collection.JavaConverters._ | ||
| import io.fabric8.kubernetes.api.model.{ContainerBuilder, PodBuilder, Volume, VolumeBuilder, VolumeMount, VolumeMountBuilder} | ||
|
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. at this point, just use the wildcard instead of a really huge line.
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 see, Thanks. |
||
|
|
||
| import org.apache.spark.deploy.k8s.{KubernetesConf, SparkPod} | ||
| import org.apache.spark.deploy.k8s.Config._ | ||
|
|
@@ -43,33 +44,52 @@ private[spark] class LocalDirsFeatureStep( | |
| 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 name = s"spark-local-dir-${index + 1}" | ||
| findVolume(pod, name) match { | ||
| case Some(volume) => volume | ||
|
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. integration tests should exercise both cases (found the volume or created it)
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 agree with you, there is an existing unit test for creating a local directory via emptyDir, this PR adds a unit test for the case of found volume. Both cases should already be covered. Is that ok? |
||
| case None => | ||
| new VolumeBuilder() | ||
| .withName(s"spark-local-dir-${index + 1}") | ||
|
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. What is the scope of visibility on volume names "spark-local-dir-1", "spark-local-dir-2" ...
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. It should be visible for global scope. Since the task distributed should run after setting up executor, so it will not collide with volumes which are built later.
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. volume names in this context are scoped to the pod spec |
||
| .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() | ||
| findVolumeMount(pod, localDirVolume.getName, localDirPath) match { | ||
| case Some(volumeMount) => volumeMount | ||
| case None => | ||
| new VolumeMountBuilder() | ||
| .withName(localDirVolume.getName) | ||
| .withMountPath(localDirPath) | ||
| .build() | ||
| } | ||
| } | ||
| val podWithLocalDirVolumes = new PodBuilder(pod.pod) | ||
| .editSpec() | ||
| .addToVolumes(localDirVolumes: _*) | ||
| .addToVolumes(localDirVolumes.filter(v => findVolume(pod, v.getName).isEmpty): _*) | ||
| .endSpec() | ||
| .build() | ||
| val containerWithLocalDirVolumeMounts = new ContainerBuilder(pod.container) | ||
| .addNewEnv() | ||
| .withName("SPARK_LOCAL_DIRS") | ||
| .withValue(resolvedLocalDirs.mkString(",")) | ||
| .endEnv() | ||
| .addToVolumeMounts(localDirVolumeMounts: _*) | ||
| .addToVolumeMounts(localDirVolumeMounts | ||
| .filter(m => findVolumeMount(pod, m.getName, m.getMountPath).isEmpty): _*) | ||
| .build() | ||
| SparkPod(podWithLocalDirVolumes, containerWithLocalDirVolumeMounts) | ||
| } | ||
|
|
||
| def findVolume(pod: SparkPod, name: String): Option[Volume] = { | ||
| pod.pod.getSpec.getVolumes.asScala.find(v => v.getName.equals(name)) | ||
| } | ||
|
|
||
| def findVolumeMount(pod: SparkPod, name: String, path: String): Option[VolumeMount] = { | ||
| pod.container.getVolumeMounts.asScala | ||
| .find(m => m.getName.equals(name) && m.getMountPath.equals(path)) | ||
| } | ||
| } | ||
| 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, | ||
|
|
||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -19,9 +19,8 @@ package org.apache.spark.deploy.k8s.features | |
| import io.fabric8.kubernetes.api.model.{EnvVarBuilder, VolumeBuilder, VolumeMountBuilder} | ||
|
|
||
| import org.apache.spark.{SparkConf, SparkFunSuite} | ||
| import org.apache.spark.deploy.k8s.{KubernetesTestConf, SparkPod} | ||
| import org.apache.spark.deploy.k8s.{KubernetesHostPathVolumeConf, KubernetesTestConf, KubernetesVolumeSpec, SparkPod} | ||
|
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. Use the wildcard when the import line gets too long. |
||
| import org.apache.spark.deploy.k8s.Config._ | ||
| import org.apache.spark.deploy.k8s.submit.JavaMainAppResource | ||
| import org.apache.spark.util.SparkConfWithEnv | ||
|
|
||
| class LocalDirsFeatureStepSuite extends SparkFunSuite { | ||
|
|
@@ -116,4 +115,28 @@ class LocalDirsFeatureStepSuite extends SparkFunSuite { | |
| .withValue(defaultLocalDir) | ||
| .build()) | ||
| } | ||
|
|
||
| test("local dir on mounted volume") { | ||
|
dongjoon-hyun marked this conversation as resolved.
|
||
| val volumeConf = KubernetesVolumeSpec( | ||
| "spark-local-dir-1", | ||
| "/tmp", | ||
| "", | ||
| false, | ||
| KubernetesHostPathVolumeConf("/hostPath/tmp") | ||
| ) | ||
| val mountVolumeConf = KubernetesTestConf.createDriverConf(volumes = Seq(volumeConf)) | ||
| val mountVolumeStep = new MountVolumesFeatureStep(mountVolumeConf) | ||
| val configuredPod = mountVolumeStep.configurePod(SparkPod.initialPod()) | ||
|
|
||
| val sparkConf = new SparkConfWithEnv(Map("SPARK_LOCAL_DIRS" -> "/tmp")) | ||
| val localDirConf = KubernetesTestConf.createDriverConf(sparkConf = sparkConf) | ||
| val localDirStep = new LocalDirsFeatureStep(localDirConf, defaultLocalDir) | ||
| val newConfiguredPod = localDirStep.configurePod(configuredPod) | ||
|
|
||
| assert(newConfiguredPod.pod.getSpec.getVolumes.size() === 1) | ||
| assert(newConfiguredPod.pod.getSpec.getVolumes.get(0).getHostPath.getPath === "/hostPath/tmp") | ||
| assert(newConfiguredPod.container.getVolumeMounts.size() === 1) | ||
| assert(newConfiguredPod.container.getVolumeMounts.get(0).getMountPath === "/tmp") | ||
| assert(newConfiguredPod.container.getVolumeMounts.get(0).getName === "spark-local-dir-1") | ||
| } | ||
| } | ||
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.
import scala....