Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 2 additions & 2 deletions docs/usage/properties/instance/user/bulk_import.md
Original file line number Diff line number Diff line change
Expand Up @@ -109,8 +109,8 @@ Note that on EMR, the total resource allocation must align with the instance typ
| sleeper.bulk.import.eks.spark.default.parallelism | (EKS mode only) The default parallelism for the Spark job. Used to set spark.default.parallelism. Should scale with the total cores across the EKS cluster, which may differ from EMR.<br>See https://spark.apache.org/docs/latest/configuration.html. | 290 | false |
| sleeper.bulk.import.eks.spark.sql.shuffle.partitions | (EKS mode only) The number of partitions used in a Spark SQL/dataframe shuffle operation. Used to set spark.sql.shuffle.partitions.<br>See https://spark.apache.org/docs/latest/configuration.html. | 290 | false |
| sleeper.bulk.import.eks.spark.dynamic.allocation.enabled | (EKS mode only) Whether Spark should use dynamic allocation to scale resources up and down. Used to set spark.dynamicAllocation.enabled. Kubernetes support for dynamic allocation is more limited than YARN's; consider leaving this disabled on EKS.<br>See https://spark.apache.org/docs/latest/configuration.html. | false | false |
| sleeper.bulk.import.eks.spark.executor.extra.java.options | JVM options passed to the executors. Used to set spark.executor.extraJavaOptions.<br>See https://spark.apache.org/docs/latest/configuration.html. | -XX:+UseG1GC -XX:+UnlockDiagnosticVMOptions -XX:+G1SummarizeConcMark -XX:InitiatingHeapOccupancyPercent=35 -verbose:gc -XX:+PrintGCDetails -XX:+PrintGCDateStamps -XX:OnOutOfMemoryError='kill -9 %p' | false |
| sleeper.bulk.import.eks.spark.driver.extra.java.options | JVM options passed to the driver. Used to set spark.driver.extraJavaOptions.<br>See https://spark.apache.org/docs/latest/configuration.html. | -XX:+UseG1GC -XX:+UnlockDiagnosticVMOptions -XX:+G1SummarizeConcMark -XX:InitiatingHeapOccupancyPercent=35 -verbose:gc -XX:+PrintGCDetails -XX:+PrintGCDateStamps -XX:OnOutOfMemoryError='kill -9 %p' | false |
| sleeper.bulk.import.eks.spark.executor.extra.java.options | JVM options passed to the executors. Used to set spark.executor.extraJavaOptions.<br>See https://spark.apache.org/docs/latest/configuration.html. | -XX:+UnlockDiagnosticVMOptions -XX:+G1SummarizeConcMark -XX:InitiatingHeapOccupancyPercent=35 -verbose:gc -XX:+PrintGCDetails -XX:+PrintGCDateStamps -XX:OnOutOfMemoryError='kill -9 %p' | false |
| sleeper.bulk.import.eks.spark.driver.extra.java.options | JVM options passed to the driver. Used to set spark.driver.extraJavaOptions.<br>See https://spark.apache.org/docs/latest/configuration.html. | -XX:+UnlockDiagnosticVMOptions -XX:+G1SummarizeConcMark -XX:InitiatingHeapOccupancyPercent=35 -verbose:gc -XX:+PrintGCDetails -XX:+PrintGCDateStamps -XX:OnOutOfMemoryError='kill -9 %p' | false |
| sleeper.bulk.import.eks.spark.executor.heartbeat.interval | The interval between heartbeats from executors to the driver. Used to set spark.executor.heartbeatInterval.<br>See https://spark.apache.org/docs/latest/configuration.html. | 60s | false |
| sleeper.bulk.import.eks.spark.network.timeout | The default timeout for network interactions in Spark. Used to set spark.network.timeout.<br>See https://spark.apache.org/docs/latest/configuration.html. | 800s | false |
| sleeper.bulk.import.eks.spark.memory.fraction | The fraction of heap space used for execution and storage. Used to set spark.memory.fraction.<br>See https://spark.apache.org/docs/latest/configuration.html. | 0.80 | false |
Expand Down
4 changes: 2 additions & 2 deletions example/full/instance.properties
Original file line number Diff line number Diff line change
Expand Up @@ -1247,12 +1247,12 @@
# JVM options passed to the executors. Used to set spark.executor.extraJavaOptions.
# See https://spark.apache.org/docs/latest/configuration.html.
# (default value shown below, uncomment to set a value)
# sleeper.bulk.import.eks.spark.executor.extra.java.options=-XX:+UseG1GC -XX:+UnlockDiagnosticVMOptions -XX:+G1SummarizeConcMark -XX:InitiatingHeapOccupancyPercent=35 -verbose:gc -XX:+PrintGCDetails -XX:+PrintGCDateStamps -XX:OnOutOfMemoryError='kill -9 %p'
# sleeper.bulk.import.eks.spark.executor.extra.java.options=-XX:+UnlockDiagnosticVMOptions -XX:+G1SummarizeConcMark -XX:InitiatingHeapOccupancyPercent=35 -verbose:gc -XX:+PrintGCDetails -XX:+PrintGCDateStamps -XX:OnOutOfMemoryError='kill -9 %p'

# JVM options passed to the driver. Used to set spark.driver.extraJavaOptions.
# See https://spark.apache.org/docs/latest/configuration.html.
# (default value shown below, uncomment to set a value)
# sleeper.bulk.import.eks.spark.driver.extra.java.options=-XX:+UseG1GC -XX:+UnlockDiagnosticVMOptions -XX:+G1SummarizeConcMark -XX:InitiatingHeapOccupancyPercent=35 -verbose:gc -XX:+PrintGCDetails -XX:+PrintGCDateStamps -XX:OnOutOfMemoryError='kill -9 %p'
# sleeper.bulk.import.eks.spark.driver.extra.java.options=-XX:+UnlockDiagnosticVMOptions -XX:+G1SummarizeConcMark -XX:InitiatingHeapOccupancyPercent=35 -verbose:gc -XX:+PrintGCDetails -XX:+PrintGCDateStamps -XX:OnOutOfMemoryError='kill -9 %p'

# The interval between heartbeats from executors to the driver. Used to set
# spark.executor.heartbeatInterval.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -214,6 +214,7 @@ public static Map<String, String> getSparkConfigurationForEKSFromInstancePropert
sparkConf.put("spark.kubernetes.namespace", instanceProperties.get(BULK_IMPORT_EKS_NAMESPACE));
sparkConf.put("spark.kubernetes.authenticate.driver.serviceAccountName", "spark");
sparkConf.put("spark.kubernetes.executor.podTemplateFile", "/tmp/executor-template.yaml");
sparkConf.put("spark.kubernetes.executor.podTemplateContainerName", "spark-kubernetes-executor");

// Ensure pods aren't disrupted/evicted before they finish
sparkConf.put("spark.kubernetes.driver.annotation.karpenter.sh/do-not-disrupt", "true");
Expand Down
15 changes: 4 additions & 11 deletions java/bulk-import/bulk-import-eks/docker/eks/Dockerfile
Original file line number Diff line number Diff line change
Expand Up @@ -12,23 +12,16 @@
# See the License for the specific language governing permissions and
# limitations under the License.

ARG BASE_IMAGE=apache/spark:3.5.6-scala2.12-java17-ubuntu
ARG BASE_IMAGE=public.ecr.aws/emr-on-eks/spark/emr-7.12.0:latest
FROM ${BASE_IMAGE}

ENV PATH="$PATH:/opt/spark/bin"
USER root
RUN apt-get update && apt-get upgrade -y && rm -rf /var/lib/apt/lists/*
RUN rm /opt/spark/jars/*
RUN mkdir /opt/spark/workdir
USER spark

# Replace Spark jars with versions managed by Sleeper (the version of Spark must match)
COPY ./spark/* /opt/spark/jars
COPY ./bulk-import-runner.jar /opt/spark/workdir
RUN mkdir -p /opt/spark/workdir && chown hadoop:hadoop /opt/spark/workdir
COPY --chown=hadoop:hadoop ./bulk-import-runner.jar /opt/spark/workdir/

USER root
COPY ./sleeper-entrypoint.sh /opt/sleeper-entrypoint.sh
RUN chmod +x /opt/sleeper-entrypoint.sh
USER spark

USER hadoop:hadoop
ENTRYPOINT ["/opt/sleeper-entrypoint.sh"]
Original file line number Diff line number Diff line change
Expand Up @@ -19,4 +19,4 @@ if [ -n "${EXECUTOR_POD_TEMPLATE:-}" ]; then
printf '%s' "$EXECUTOR_POD_TEMPLATE" > /tmp/executor-template.yaml
fi

exec /opt/entrypoint.sh "$@"
exec /usr/bin/entrypoint.sh "$@"
42 changes: 0 additions & 42 deletions java/bulk-import/bulk-import-eks/pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -24,46 +24,4 @@
<modelVersion>4.0.0</modelVersion>

<artifactId>bulk-import-eks</artifactId>

<dependencies>
<dependency>
<groupId>org.apache.spark</groupId>
<artifactId>spark-sql_${scala.version}</artifactId>
<scope>runtime</scope>
</dependency>
<dependency>
<groupId>org.apache.spark</groupId>
<artifactId>spark-kubernetes_${scala.version}</artifactId>
<scope>runtime</scope>
</dependency>
<dependency>
<groupId>org.apache.arrow</groupId>
<artifactId>arrow-memory-unsafe</artifactId>
<scope>runtime</scope>
</dependency>
</dependencies>

<build>
<plugins>
<plugin>
<artifactId>maven-dependency-plugin</artifactId>
<executions>
<execution>
<!-- The Spark Docker image uses an old version of Hadoop, which is not compatible with the -->
<!-- version used in EMR. -->
<!-- This outputs the jars required to run Spark in EKS, to completely replace the jars -->
<!-- in the Spark Docker image with versions managed by Sleeper. -->
<!-- The version of Spark must match between the base image and Sleeper. -->
<goals>
<goal>copy-dependencies</goal>
</goals>
<configuration>
<outputDirectory>${project.build.directory}/spark</outputDirectory>
<includeScope>runtime</includeScope>
</configuration>
</execution>
</executions>
</plugin>
</plugins>
</build>
</project>
Original file line number Diff line number Diff line change
Expand Up @@ -37,7 +37,6 @@
*/
public class StateMachinePlatformExecutor implements PlatformExecutor {
private static final String SPARK_IMAGE_JAR_LOCATION = "local:///opt/spark/workdir/bulk-import-runner.jar";
private static final String SPARK_IMAGE_JAVA_HOME = "/opt/java/openjdk";
private static final String NATIVE_IMAGE_JAR_LOCATION = "local:///opt/spark/workdir/bulk-import-runner.jar";
private static final String NATIVE_IMAGE_LOG4J_LOCATION = "file:///opt/spark/workdir/log4j.properties";
private static final String NATIVE_IMAGE_JAVA_HOME = "/usr/lib/jvm/java-11-amazon-corretto";
Expand Down Expand Up @@ -92,7 +91,6 @@ private List<String> constructArgs(BulkImportArguments arguments, String taskId)
baseSparkConfig.put("spark.executor.extraJavaOptions", "-Dlog4j.configuration=" + NATIVE_IMAGE_LOG4J_LOCATION);
jarLocation = NATIVE_IMAGE_JAR_LOCATION;
} else {
baseSparkConfig.put("spark.executorEnv.JAVA_HOME", SPARK_IMAGE_JAVA_HOME);
jarLocation = SPARK_IMAGE_JAR_LOCATION;
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -12,7 +12,7 @@
"--conf",
"spark.driver.cores=5",
"--conf",
"spark.driver.extraJavaOptions=-XX:+UseG1GC -XX:+UnlockDiagnosticVMOptions -XX:+G1SummarizeConcMark -XX:InitiatingHeapOccupancyPercent=35 -verbose:gc -XX:+PrintGCDetails -XX:+PrintGCDateStamps -XX:OnOutOfMemoryError='kill -9 %p'",
"spark.driver.extraJavaOptions=-XX:+UnlockDiagnosticVMOptions -XX:+G1SummarizeConcMark -XX:InitiatingHeapOccupancyPercent=35 -verbose:gc -XX:+PrintGCDetails -XX:+PrintGCDateStamps -XX:OnOutOfMemoryError='kill -9 %p'",
"--conf",
"spark.driver.memory=12g",
"--conf",
Expand All @@ -22,7 +22,7 @@
"--conf",
"spark.executor.cores=5",
"--conf",
"spark.executor.extraJavaOptions=-XX:+UseG1GC -XX:+UnlockDiagnosticVMOptions -XX:+G1SummarizeConcMark -XX:InitiatingHeapOccupancyPercent=35 -verbose:gc -XX:+PrintGCDetails -XX:+PrintGCDateStamps -XX:OnOutOfMemoryError='kill -9 %p'",
"spark.executor.extraJavaOptions=-XX:+UnlockDiagnosticVMOptions -XX:+G1SummarizeConcMark -XX:InitiatingHeapOccupancyPercent=35 -verbose:gc -XX:+PrintGCDetails -XX:+PrintGCDateStamps -XX:OnOutOfMemoryError='kill -9 %p'",
"--conf",
"spark.executor.heartbeatInterval=60s",
"--conf",
Expand All @@ -32,8 +32,6 @@
"--conf",
"spark.executor.memoryOverhead=1g",
"--conf",
"spark.executorEnv.JAVA_HOME=/opt/java/openjdk",
"--conf",
"spark.hadoop.fs.s3a.aws.credentials.provider=software.amazon.awssdk.auth.credentials.WebIdentityTokenFileCredentialsProvider",
"--conf",
"spark.hadoop.fs.s3a.block.size=32M",
Expand Down Expand Up @@ -62,6 +60,8 @@
"--conf",
"spark.kubernetes.executor.podNamePrefix=job-test-job",
"--conf",
"spark.kubernetes.executor.podTemplateContainerName=spark-kubernetes-executor",
"--conf",
"spark.kubernetes.executor.podTemplateFile=/tmp/executor-template.yaml",
"--conf",
"spark.kubernetes.namespace=eks-namespace",
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -144,7 +144,7 @@ public interface EKSProperty {
.description("JVM options passed to the executors. Used to set spark.executor.extraJavaOptions.\n" +
"See https://spark.apache.org/docs/latest/configuration.html.")
.defaultValue(
"-XX:+UseG1GC -XX:+UnlockDiagnosticVMOptions -XX:+G1SummarizeConcMark -XX:InitiatingHeapOccupancyPercent=35 -verbose:gc -XX:+PrintGCDetails -XX:+PrintGCDateStamps -XX:OnOutOfMemoryError='kill -9 %p'")
"-XX:+UnlockDiagnosticVMOptions -XX:+G1SummarizeConcMark -XX:InitiatingHeapOccupancyPercent=35 -verbose:gc -XX:+PrintGCDetails -XX:+PrintGCDateStamps -XX:OnOutOfMemoryError='kill -9 %p'")
.propertyGroup(InstancePropertyGroup.BULK_IMPORT).build();
UserDefinedInstanceProperty BULK_IMPORT_EKS_SPARK_DRIVER_EXTRA_JAVA_OPTIONS = Index.propertyBuilder("sleeper.bulk.import.eks.spark.driver.extra.java.options")
.description("JVM options passed to the driver. Used to set spark.driver.extraJavaOptions.\n" +
Expand Down
4 changes: 2 additions & 2 deletions scripts/templates/instanceproperties.template
Original file line number Diff line number Diff line change
Expand Up @@ -1251,12 +1251,12 @@
# JVM options passed to the executors. Used to set spark.executor.extraJavaOptions.
# See https://spark.apache.org/docs/latest/configuration.html.
# (default value shown below, uncomment to set a value)
# sleeper.bulk.import.eks.spark.executor.extra.java.options=-XX:+UseG1GC -XX:+UnlockDiagnosticVMOptions -XX:+G1SummarizeConcMark -XX:InitiatingHeapOccupancyPercent=35 -verbose:gc -XX:+PrintGCDetails -XX:+PrintGCDateStamps -XX:OnOutOfMemoryError='kill -9 %p'
# sleeper.bulk.import.eks.spark.executor.extra.java.options=-XX:+UnlockDiagnosticVMOptions -XX:+G1SummarizeConcMark -XX:InitiatingHeapOccupancyPercent=35 -verbose:gc -XX:+PrintGCDetails -XX:+PrintGCDateStamps -XX:OnOutOfMemoryError='kill -9 %p'

# JVM options passed to the driver. Used to set spark.driver.extraJavaOptions.
# See https://spark.apache.org/docs/latest/configuration.html.
# (default value shown below, uncomment to set a value)
# sleeper.bulk.import.eks.spark.driver.extra.java.options=-XX:+UseG1GC -XX:+UnlockDiagnosticVMOptions -XX:+G1SummarizeConcMark -XX:InitiatingHeapOccupancyPercent=35 -verbose:gc -XX:+PrintGCDetails -XX:+PrintGCDateStamps -XX:OnOutOfMemoryError='kill -9 %p'
# sleeper.bulk.import.eks.spark.driver.extra.java.options=-XX:+UnlockDiagnosticVMOptions -XX:+G1SummarizeConcMark -XX:InitiatingHeapOccupancyPercent=35 -verbose:gc -XX:+PrintGCDetails -XX:+PrintGCDateStamps -XX:OnOutOfMemoryError='kill -9 %p'

# The interval between heartbeats from executors to the driver. Used to set
# spark.executor.heartbeatInterval.
Expand Down
Loading