Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
21 commits
Select commit Hold shift + click to select a range
270c221
7607 Pass key into BulkImportJobWriterToS3
patchwork01 Jul 17, 2026
b378792
7607 Pass S3 object key as an argument to Spark
patchwork01 Jul 17, 2026
c858cb2
7607 Inline job ID in BulkImportJobLoaderFromS3IT
patchwork01 Jul 17, 2026
e6311a6
7607 Split EKS bulk import system test
patchwork01 Jul 17, 2026
bb800fe
7607 Update system-test-suites.md
patchwork01 Jul 17, 2026
ebba103
7607 Add missing brackets in BulkImportArguments
patchwork01 Jul 20, 2026
191bfea
7607 Fix GenerateSystemTestSuiteDocumentationIT
patchwork01 Jul 20, 2026
508cba0
Merge branch 'develop' into 7607-explicit-bulk-import-s3-path
patchwork01 Jul 20, 2026
ca0a3d6
Merge branch 'develop' into 7607-explicit-bulk-import-s3-path
patchwork01 Jul 20, 2026
495eb2f
7607 Fix EMR Serverless arguments
patchwork01 Jul 20, 2026
55f9e93
7607 Refactor entry point arguments
patchwork01 Jul 20, 2026
d9f9deb
7607 Move EksAutoBulkImportST to Slow3 suite
patchwork01 Jul 20, 2026
01f49e7
7607 Swap IngestPerformanceST and CompactionPerformanceST between suites
patchwork01 Jul 20, 2026
b36613c
7607 Fix GenerateSystemTestSuiteDocumentationIT
patchwork01 Jul 20, 2026
67be31a
Merge branch 'develop' into 7607-explicit-bulk-import-s3-path
patchwork01 Aug 3, 2026
e93e224
Merge branch 'develop' into 7607-explicit-bulk-import-s3-path
patchwork01 Aug 3, 2026
7c2caa0
7607 Fix usage check of BulkImportJobDriver
patchwork01 Aug 4, 2026
96631af
Merge branch 'develop' into 7607-explicit-bulk-import-s3-path
patchwork01 Aug 4, 2026
b56f60e
7607 Adjust comment on EksBulkImportPerformanceST
patchwork01 Aug 4, 2026
09535d6
7607 Fix DirectEmrServerlessDriver
patchwork01 Aug 4, 2026
531e205
7607 Deploy EKS system test instances with EKS enabled
patchwork01 Aug 4, 2026
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
18 changes: 9 additions & 9 deletions docs/development/system-test-suites.md
Original file line number Diff line number Diff line change
Expand Up @@ -6,15 +6,15 @@ it takes to complete the nightly system tests.

| Slow1 | Slow2 | Slow3 |
|----------------------------|------------------------------------|---------------------------------|
| AutoDeleteS3ObjectsST | CompactionOnEC2ST | MultipleTablesST |
| AutoStopEcsTaskST | ECSStateStoreCommitterST | StateStoreCommitterThroughputST |
| CompactionCreationST | ECSStateStoreCommitterThroughputST |
| EmrPersistentBulkImportST | EksBulkImportST |
| AutoDeleteS3ObjectsST | CompactionOnEC2ST | EksAutoBulkImportST |
| AutoStopEcsTaskST | ECSStateStoreCommitterST | MultipleTablesST |
| CompactionCreationST | ECSStateStoreCommitterThroughputST | StateStoreCommitterThroughputST |
| EmrPersistentBulkImportST | EksFargateBulkImportST |
| OptionalFeaturesDisabledST |
| RedeployOptionalStacksST |

| Expensive1 | Expensive2 | Expensive3 |
|-----------------------------------|----------------------------|-----------------------|
| CompactionDataFusionPerformanceST | CompactionPerformanceST | IngestPerformanceST |
| CompactionVeryLargeST | EksBulkImportPerformanceST | ParallelCompactionsST |
| | EmrBulkImportPerformanceST |
| Expensive1 | Expensive2 | Expensive3 |
|-----------------------------------|----------------------------|-------------------------|
| CompactionDataFusionPerformanceST | EksBulkImportPerformanceST | CompactionPerformanceST |
| CompactionVeryLargeST | EmrBulkImportPerformanceST | ParallelCompactionsST |
| | IngestPerformanceST |
Original file line number Diff line number Diff line change
Expand Up @@ -113,9 +113,9 @@ public BulkImportJobDriver(
* {@link org.apache.spark.SparkException}, which extends {@link Exception} but is thrown unchecked at runtime
* because Spark is implemented in Scala, which does not enforce checked exceptions.
*
* @param job the bulk import job to run
* @param jobRunId the ID of this run of the job
* @param taskId the ID of the task running the job
* @param job the bulk import job to run
* @param jobRunId the ID of this run of the job
* @param taskId the ID of the task running the job
* @throws Exception if the job fails for any reason, including Spark failures
*/
public void run(BulkImportJob job, String jobRunId, String taskId) throws Exception {
Expand Down Expand Up @@ -228,15 +228,16 @@ public static void start(String[] args, BulkImportJobRunner runner) {
}

private static void startOrThrow(String[] args, BulkImportJobRunner runner) throws Exception {
if (args.length != 5) {
throw new IllegalArgumentException("Expected 5 arguments:" +
" <config bucket name> <bulk import job ID> <bulk import task ID> <bulk import job run ID> <bulk import mode>");
if (args.length != 6) {
throw new IllegalArgumentException("Expected 6 arguments:" +
" <config bucket name> <bulk import job ID> <bulk import task ID> <bulk import job run ID> <S3 object key of job definition file> <bulk import mode>");
}
String configBucket = args[0];
String jobId = args[1];
String taskId = args[2];
String jobRunId = args[3];
String bulkImportMode = args[4];
String jobFileObjectKey = args[4];
String bulkImportMode = args[5];

LOGGER.info("Starting bulk import job driver");
LOGGER.info("Config bucket: {}", configBucket);
Expand All @@ -261,7 +262,7 @@ private static void startOrThrow(String[] args, BulkImportJobRunner runner) thro
throw e;
}

BulkImportJob bulkImportJob = BulkImportJobLoaderFromS3.loadJob(instanceProperties, jobId, jobRunId, s3Client);
BulkImportJob bulkImportJob = BulkImportJobLoaderFromS3.loadJob(instanceProperties, jobFileObjectKey, s3Client);

TablePropertiesProvider tablePropertiesProvider = S3TableProperties.createProvider(instanceProperties, s3Client, dynamoClient);
StateStoreProvider stateStoreProvider = StateStoreFactory.createProvider(instanceProperties, s3Client, dynamoClient);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -36,16 +36,15 @@ private BulkImportJobLoaderFromS3() {

}

public static BulkImportJob loadJob(InstanceProperties instanceProperties, String jobId, String jobRunId, S3Client s3Client) {
public static BulkImportJob loadJob(InstanceProperties instanceProperties, String objectKey, S3Client s3Client) {
String bulkImportBucket = instanceProperties.get(BULK_IMPORT_BUCKET);
if (null == bulkImportBucket) {
throw new RuntimeException("sleeper.bulk.import.bucket was not set. Has one of the bulk import stacks been deployed?");
}
String jsonJobKey = "bulk_import/" + jobId + "-" + jobRunId + ".json";
LOGGER.info("Loading bulk import job from key {} in bulk import bucket {}", jsonJobKey, bulkImportBucket);
LOGGER.info("Loading bulk import job from key {} in bulk import bucket {}", objectKey, bulkImportBucket);
String jsonJob = s3Client.getObjectAsBytes(GetObjectRequest.builder()
.bucket(bulkImportBucket)
.key(jsonJobKey)
.key(objectKey)
.build()).asUtf8String();
try {
return new BulkImportJobSerDe().fromJson(jsonJob);
Expand All @@ -55,7 +54,7 @@ public static BulkImportJob loadJob(InstanceProperties instanceProperties, Strin
} finally {
s3Client.deleteObject(DeleteObjectRequest.builder()
.bucket(bulkImportBucket)
.key(jsonJobKey)
.key(objectKey)
.build());
}
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -43,20 +43,19 @@ void setUp() {
@Test
void shouldLoadBulkImportJobFromS3() {
// Given
String jobRunId = "load-run";
String jobId = "load-job-id";
String objectKey = "folder/test.json";

BulkImportJob bulkImportJob = BulkImportJob.builder()
.id(jobId)
.id("load-job-id")
.tableId("test-table-id")
.files(List.of("/load-job.parquet"))
.build();

BulkImportJobWriterToS3 bulkImportJobWriterToS3 = new BulkImportJobWriterToS3(instanceProperties, s3Client);
bulkImportJobWriterToS3.writeJobToBulkImportBucket(bulkImportJob, jobRunId);
bulkImportJobWriterToS3.writeJobToBulkImportBucket(bulkImportJob, objectKey);

// When / Then
assertThat(BulkImportJobLoaderFromS3.loadJob(instanceProperties, jobId, jobRunId, s3Client))
assertThat(BulkImportJobLoaderFromS3.loadJob(instanceProperties, objectKey, s3Client))
.isEqualTo(bulkImportJob);
// And the file is deleted after it is loaded
assertThat(listObjectKeys(instanceProperties.get(BULK_IMPORT_BUCKET))).isEmpty();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -30,18 +30,21 @@

import static java.util.Map.entry;
import static sleeper.core.properties.instance.BulkImportProperty.BULK_IMPORT_CLASS_NAME;
import static sleeper.core.properties.instance.CdkDefinedInstanceProperty.BULK_IMPORT_EMR_SERVERLESS_CLUSTER_NAME;
import static sleeper.core.properties.instance.CdkDefinedInstanceProperty.CONFIG_BUCKET;

public class BulkImportArguments {

private final InstanceProperties instanceProperties;
private final BulkImportJob bulkImportJob;
private final String jobFileObjectKey;
private final String jobRunId;

private BulkImportArguments(Builder builder) {
instanceProperties = builder.instanceProperties;
bulkImportJob = builder.bulkImportJob;
jobRunId = builder.jobRunId;
instanceProperties = Objects.requireNonNull(builder.instanceProperties, "instanceProperties must not be null");
bulkImportJob = Objects.requireNonNull(builder.bulkImportJob, "bulkImportJob must not be null");
jobFileObjectKey = Objects.requireNonNull(builder.jobFileObjectKey, "jobFileObjectKey must not be null");
jobRunId = Objects.requireNonNull(builder.jobRunId, "jobRunId must not be null");
}

public static Builder builder() {
Expand All @@ -61,16 +64,23 @@ public List<String> sparkSubmitCommandForEKSCluster(String taskId, String jarLoc
}

private List<String> sparkSubmitCommandForCluster(String taskId, String jarLocation, Map<String, String> baseSparkConfig, String bulkImportMode) {
String configBucket = instanceProperties.get(CONFIG_BUCKET);
String jobId = bulkImportJob.getId();
return Stream.of(
Stream.of("spark-submit", "--deploy-mode", "cluster"),
sparkSubmitParameters(baseSparkConfig),
Stream.of(jarLocation, configBucket, jobId, taskId, jobRunId, bulkImportMode))
Stream.of(jarLocation),
Stream.of(entryPointArguments(taskId, bulkImportMode)))
.flatMap(partialArgs -> partialArgs)
.collect(Collectors.toUnmodifiableList());
}

public String[] entryPointArgumentsForServerless() {
return entryPointArguments(instanceProperties.get(BULK_IMPORT_EMR_SERVERLESS_CLUSTER_NAME) + "-EMRS", "EMR");
}

private String[] entryPointArguments(String taskId, String bulkImportMode) {
return new String[]{instanceProperties.get(CONFIG_BUCKET), bulkImportJob.getId(), taskId, jobRunId, jobFileObjectKey, bulkImportMode};
}

public String sparkSubmitParametersForServerless() {
return sparkSubmitParameters(
SparkConfigurationUtils.getSparkServerlessConfigurationFromInstanceProperties(
Expand Down Expand Up @@ -123,9 +133,33 @@ public String getJobRunId() {
return jobRunId;
}

@Override
public String toString() {
return "BulkImportArguments{instanceProperties=" + instanceProperties + ", bulkImportJob=" + bulkImportJob + ", jobFileObjectKey=" + jobFileObjectKey + ", jobRunId=" + jobRunId + "}";
}

@Override
public int hashCode() {
return Objects.hash(instanceProperties, bulkImportJob, jobFileObjectKey, jobRunId);
}

@Override
public boolean equals(Object obj) {
if (this == obj) {
return true;
}
if (!(obj instanceof BulkImportArguments)) {
return false;
}
BulkImportArguments other = (BulkImportArguments) obj;
return Objects.equals(instanceProperties, other.instanceProperties) && Objects.equals(bulkImportJob, other.bulkImportJob) && Objects.equals(jobFileObjectKey, other.jobFileObjectKey)
&& Objects.equals(jobRunId, other.jobRunId);
}

public static final class Builder {
private InstanceProperties instanceProperties;
private BulkImportJob bulkImportJob;
private String jobFileObjectKey;
private String jobRunId;

private Builder() {
Expand All @@ -141,6 +175,11 @@ public Builder bulkImportJob(BulkImportJob bulkImportJob) {
return this;
}

public Builder jobFileObjectKey(String jobFileObjectKey) {
this.jobFileObjectKey = jobFileObjectKey;
return this;
}

public Builder jobRunId(String jobRunId) {
this.jobRunId = jobRunId;
return this;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -57,6 +57,10 @@ public BulkImportExecutor(
this.validationTimeSupplier = validationTimeSupplier;
}

public static String createJobFileObjectKey(BulkImportJob job, String jobRunId) {
return "bulk_import/" + job.getId() + "-" + jobRunId + ".json";
}

public void runJob(BulkImportJob bulkImportJob) {
runJob(bulkImportJob, UUID.randomUUID().toString());
}
Expand All @@ -75,11 +79,12 @@ public void runJob(BulkImportJob bulkImportJob, String jobRunId) {
.jobRunId(jobRunId).build());
try {
LOGGER.info("Writing job with id {} to JSON file", bulkImportJob.getId());
writeJobToBucket.writeJobToBulkImportBucket(bulkImportJob, jobRunId);
String jobFileObjectKey = createJobFileObjectKey(bulkImportJob, jobRunId);
writeJobToBucket.writeJobToBulkImportBucket(bulkImportJob, jobFileObjectKey);
LOGGER.info("Submitting job with id {}", bulkImportJob.getId());
platformExecutor.runJobOnPlatform(BulkImportArguments.builder()
.instanceProperties(instanceProperties)
.bulkImportJob(bulkImportJob).jobRunId(jobRunId)
.bulkImportJob(bulkImportJob).jobRunId(jobRunId).jobFileObjectKey(jobFileObjectKey)
.build());
LOGGER.info("Successfully submitted job");
} catch (RuntimeException e) {
Expand Down Expand Up @@ -125,6 +130,6 @@ private boolean validateJob(BulkImportJob bulkImportJob) {
@FunctionalInterface
public interface WriteJobToBucket {

void writeJobToBulkImportBucket(BulkImportJob bulkImportJob, String jobRunID);
void writeJobToBulkImportBucket(BulkImportJob bulkImportJob, String objectKey);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -39,18 +39,17 @@ public BulkImportJobWriterToS3(InstanceProperties instanceProperties, S3Client s
}

@Override
public void writeJobToBulkImportBucket(BulkImportJob bulkImportJob, String jobRunID) {
public void writeJobToBulkImportBucket(BulkImportJob bulkImportJob, String objectKey) {
String bulkImportBucket = instanceProperties.get(BULK_IMPORT_BUCKET);
if (null == bulkImportBucket) {
throw new RuntimeException("sleeper.bulk.import.bucket was not set. Has one of the bulk import stacks been deployed?");
}
String key = "bulk_import/" + bulkImportJob.getId() + "-" + jobRunID + ".json";
String bulkImportJobJSON = new BulkImportJobSerDe().toJson(bulkImportJob);
s3Client.putObject(PutObjectRequest.builder()
.bucket(bulkImportBucket)
.key(key)
.key(objectKey)
.build(),
RequestBody.fromString(bulkImportJobJSON));
LOGGER.info("Put object for job {} to key {} in bucket {}", bulkImportJob.getId(), key, bulkImportBucket);
LOGGER.info("Put object for job {} to key {} in bucket {}", bulkImportJob.getId(), objectKey, bulkImportBucket);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -23,14 +23,11 @@
import software.amazon.awssdk.services.emrserverless.model.SparkSubmit;
import software.amazon.awssdk.services.emrserverless.model.StartJobRunRequest;

import sleeper.bulkimport.core.job.BulkImportJob;
import sleeper.core.properties.instance.InstanceProperties;

import static sleeper.core.properties.instance.CdkDefinedInstanceProperty.BULK_IMPORT_BUCKET;
import static sleeper.core.properties.instance.CdkDefinedInstanceProperty.BULK_IMPORT_EMR_SERVERLESS_APPLICATION_ID;
import static sleeper.core.properties.instance.CdkDefinedInstanceProperty.BULK_IMPORT_EMR_SERVERLESS_CLUSTER_NAME;
import static sleeper.core.properties.instance.CdkDefinedInstanceProperty.BULK_IMPORT_EMR_SERVERLESS_CLUSTER_ROLE_ARN;
import static sleeper.core.properties.instance.CdkDefinedInstanceProperty.CONFIG_BUCKET;
import static sleeper.core.properties.instance.CdkDefinedInstanceProperty.VERSION;
import static sleeper.core.properties.instance.CommonProperty.JARS_BUCKET;

Expand All @@ -48,20 +45,14 @@ public EmrServerlessPlatformExecutor(EmrServerlessClient emrClient, InstanceProp

@Override
public void runJobOnPlatform(BulkImportArguments arguments) {
BulkImportJob bulkImportJob = arguments.getBulkImportJob();
String jobName = String.join("-", "job", arguments.getJobRunId());
String applicationName = instanceProperties.get(BULK_IMPORT_EMR_SERVERLESS_CLUSTER_NAME);
String applicationId = instanceProperties.get(BULK_IMPORT_EMR_SERVERLESS_APPLICATION_ID);

StartJobRunRequest job = StartJobRunRequest.builder()
.applicationId(applicationId)
.name(jobName)
.applicationId(instanceProperties.get(BULK_IMPORT_EMR_SERVERLESS_APPLICATION_ID))
.name(String.join("-", "job", arguments.getJobRunId()))
.executionRoleArn(
instanceProperties.get(BULK_IMPORT_EMR_SERVERLESS_CLUSTER_ROLE_ARN))
.jobDriver(JobDriver.builder().sparkSubmit(SparkSubmit.builder()
.entryPoint("s3://" + instanceProperties.get(JARS_BUCKET) + "/bulk-import-runner-" + instanceProperties.get(VERSION) + ".jar")
.entryPointArguments(instanceProperties.get(CONFIG_BUCKET),
bulkImportJob.getId(), applicationName + "-EMRS", arguments.getJobRunId(), "EMR")
.entryPointArguments(arguments.entryPointArgumentsForServerless())
.sparkSubmitParameters(arguments.sparkSubmitParametersForServerless())
.build()).build())
.configurationOverrides(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -47,6 +47,7 @@ void shouldConstructArgs() {
.files(Lists.newArrayList("file1.parquet"))
.build())
.jobRunId("test-run")
.jobFileObjectKey("folder/job.json")
.build();

// When / Then
Expand All @@ -61,6 +62,7 @@ void shouldConstructArgs() {
"my-job",
"test-task",
"test-run",
"folder/job.json",
"EMR");
}
}
Loading
Loading