From 270c22132200d7f185640872b99d5b11267be00e Mon Sep 17 00:00:00 2001 From: patchwork01 <110390516+patchwork01@users.noreply.github.com> Date: Fri, 17 Jul 2026 14:11:17 +0000 Subject: [PATCH 01/16] 7607 Pass key into BulkImportJobWriterToS3 --- .../sleeper/bulkimport/runner/BulkImportJobDriver.java | 9 +++++---- .../bulkimport/runner/BulkImportJobLoaderFromS3.java | 9 ++++----- .../bulkimport/runner/BulkImportJobLoaderFromS3IT.java | 6 +++--- .../starter/executor/BulkImportExecutor.java | 5 +++-- .../starter/executor/BulkImportJobWriterToS3.java | 7 +++---- .../starter/executor/BulkImportExecutorTest.java | 10 +++++----- 6 files changed, 23 insertions(+), 23 deletions(-) diff --git a/java/bulk-import/bulk-import-runner/src/main/java/sleeper/bulkimport/runner/BulkImportJobDriver.java b/java/bulk-import/bulk-import-runner/src/main/java/sleeper/bulkimport/runner/BulkImportJobDriver.java index dfffe8c6987..a7dd9817e02 100644 --- a/java/bulk-import/bulk-import-runner/src/main/java/sleeper/bulkimport/runner/BulkImportJobDriver.java +++ b/java/bulk-import/bulk-import-runner/src/main/java/sleeper/bulkimport/runner/BulkImportJobDriver.java @@ -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 { @@ -261,7 +261,8 @@ private static void startOrThrow(String[] args, BulkImportJobRunner runner) thro throw e; } - BulkImportJob bulkImportJob = BulkImportJobLoaderFromS3.loadJob(instanceProperties, jobId, jobRunId, s3Client); + String key = "bulk_import/" + jobId + "-" + jobRunId + ".json"; + BulkImportJob bulkImportJob = BulkImportJobLoaderFromS3.loadJob(instanceProperties, key, s3Client); TablePropertiesProvider tablePropertiesProvider = S3TableProperties.createProvider(instanceProperties, s3Client, dynamoClient); StateStoreProvider stateStoreProvider = StateStoreFactory.createProvider(instanceProperties, s3Client, dynamoClient); diff --git a/java/bulk-import/bulk-import-runner/src/main/java/sleeper/bulkimport/runner/BulkImportJobLoaderFromS3.java b/java/bulk-import/bulk-import-runner/src/main/java/sleeper/bulkimport/runner/BulkImportJobLoaderFromS3.java index 397292e33a4..90036409261 100644 --- a/java/bulk-import/bulk-import-runner/src/main/java/sleeper/bulkimport/runner/BulkImportJobLoaderFromS3.java +++ b/java/bulk-import/bulk-import-runner/src/main/java/sleeper/bulkimport/runner/BulkImportJobLoaderFromS3.java @@ -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); @@ -55,7 +54,7 @@ public static BulkImportJob loadJob(InstanceProperties instanceProperties, Strin } finally { s3Client.deleteObject(DeleteObjectRequest.builder() .bucket(bulkImportBucket) - .key(jsonJobKey) + .key(objectKey) .build()); } } diff --git a/java/bulk-import/bulk-import-runner/src/test/java/sleeper/bulkimport/runner/BulkImportJobLoaderFromS3IT.java b/java/bulk-import/bulk-import-runner/src/test/java/sleeper/bulkimport/runner/BulkImportJobLoaderFromS3IT.java index f75e941352a..52925f862d0 100644 --- a/java/bulk-import/bulk-import-runner/src/test/java/sleeper/bulkimport/runner/BulkImportJobLoaderFromS3IT.java +++ b/java/bulk-import/bulk-import-runner/src/test/java/sleeper/bulkimport/runner/BulkImportJobLoaderFromS3IT.java @@ -43,7 +43,7 @@ void setUp() { @Test void shouldLoadBulkImportJobFromS3() { // Given - String jobRunId = "load-run"; + String objectKey = "folder/test.json"; String jobId = "load-job-id"; BulkImportJob bulkImportJob = BulkImportJob.builder() @@ -53,10 +53,10 @@ void shouldLoadBulkImportJobFromS3() { .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(); diff --git a/java/bulk-import/bulk-import-starter/src/main/java/sleeper/bulkimport/starter/executor/BulkImportExecutor.java b/java/bulk-import/bulk-import-starter/src/main/java/sleeper/bulkimport/starter/executor/BulkImportExecutor.java index 6988f838aaa..e14cd933ffb 100644 --- a/java/bulk-import/bulk-import-starter/src/main/java/sleeper/bulkimport/starter/executor/BulkImportExecutor.java +++ b/java/bulk-import/bulk-import-starter/src/main/java/sleeper/bulkimport/starter/executor/BulkImportExecutor.java @@ -75,7 +75,8 @@ 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 key = "bulk_import/" + bulkImportJob.getId() + "-" + jobRunId + ".json"; + writeJobToBucket.writeJobToBulkImportBucket(bulkImportJob, key); LOGGER.info("Submitting job with id {}", bulkImportJob.getId()); platformExecutor.runJobOnPlatform(BulkImportArguments.builder() .instanceProperties(instanceProperties) @@ -125,6 +126,6 @@ private boolean validateJob(BulkImportJob bulkImportJob) { @FunctionalInterface public interface WriteJobToBucket { - void writeJobToBulkImportBucket(BulkImportJob bulkImportJob, String jobRunID); + void writeJobToBulkImportBucket(BulkImportJob bulkImportJob, String objectKey); } } diff --git a/java/bulk-import/bulk-import-starter/src/main/java/sleeper/bulkimport/starter/executor/BulkImportJobWriterToS3.java b/java/bulk-import/bulk-import-starter/src/main/java/sleeper/bulkimport/starter/executor/BulkImportJobWriterToS3.java index da46a585d0e..4cef05283bb 100644 --- a/java/bulk-import/bulk-import-starter/src/main/java/sleeper/bulkimport/starter/executor/BulkImportJobWriterToS3.java +++ b/java/bulk-import/bulk-import-starter/src/main/java/sleeper/bulkimport/starter/executor/BulkImportJobWriterToS3.java @@ -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); } } diff --git a/java/bulk-import/bulk-import-starter/src/test/java/sleeper/bulkimport/starter/executor/BulkImportExecutorTest.java b/java/bulk-import/bulk-import-starter/src/test/java/sleeper/bulkimport/starter/executor/BulkImportExecutorTest.java index 8118b58c4b6..5b98f489181 100644 --- a/java/bulk-import/bulk-import-starter/src/test/java/sleeper/bulkimport/starter/executor/BulkImportExecutorTest.java +++ b/java/bulk-import/bulk-import-starter/src/test/java/sleeper/bulkimport/starter/executor/BulkImportExecutorTest.java @@ -66,7 +66,7 @@ class BulkImportExecutorTest { private final String tableId = tableProperties.get(TABLE_ID); private final IngestJobTracker tracker = new InMemoryIngestJobTracker(); private final List jobsInBucket = new ArrayList<>(); - private final List jobRunIdsOfJobsInBucket = new ArrayList<>(); + private final List keysInBucket = new ArrayList<>(); private final List jobsRun = new ArrayList<>(); private final List jobRunIdsOfJobsRun = new ArrayList<>(); @@ -179,7 +179,7 @@ void shouldCallRunOnPlatformIfJobIsValid() { // Then assertThat(jobsInBucket).containsExactly(importJob); - assertThat(jobRunIdsOfJobsInBucket).containsExactly("job-run-id"); + assertThat(keysInBucket).containsExactly("bulk_import/my-job-job-run-id.json"); assertThat(jobsRun).containsExactly(importJob); assertThat(jobRunIdsOfJobsRun).containsExactly("job-run-id"); assertThat(tracker.getAllJobs(tableId)) @@ -202,7 +202,7 @@ void shouldSucceedIfS3ObjectIsADirectoryContainingFiles() { // Then assertThat(jobsInBucket).containsExactly(importJob); - assertThat(jobRunIdsOfJobsInBucket).containsExactly("job-run-id"); + assertThat(keysInBucket).containsExactly("bulk_import/my-job-job-run-id.json"); assertThat(jobsRun).containsExactly(importJob); assertThat(jobRunIdsOfJobsRun).containsExactly("job-run-id"); assertThat(tracker.getAllJobs(tableId)) @@ -337,9 +337,9 @@ public void runJobOnPlatform(BulkImportArguments arguments) { private class RecordWriteJobToBucket implements WriteJobToBucket { @Override - public void writeJobToBulkImportBucket(BulkImportJob bulkImportJob, String jobRunID) { + public void writeJobToBulkImportBucket(BulkImportJob bulkImportJob, String objectKey) { jobsInBucket.add(bulkImportJob); - jobRunIdsOfJobsInBucket.add(jobRunID); + keysInBucket.add(objectKey); } } From b378792451c6b497b8eeb6c712fdb263532bc383 Mon Sep 17 00:00:00 2001 From: patchwork01 <110390516+patchwork01@users.noreply.github.com> Date: Fri, 17 Jul 2026 14:42:10 +0000 Subject: [PATCH 02/16] 7607 Pass S3 object key as an argument to Spark --- .../runner/BulkImportJobDriver.java | 6 +- .../starter/executor/BulkImportArguments.java | 31 +++++++++- .../starter/executor/BulkImportExecutor.java | 6 +- .../executor/BulkImportArgumentsTest.java | 2 + .../executor/BulkImportExecutorTest.java | 62 ++++++++++--------- .../example/emr/runjobflow-request.json | 1 + .../persistent-emr/addjobflow-request.json | 1 + .../step-functions/startexecution-input.json | 1 + 8 files changed, 74 insertions(+), 36 deletions(-) diff --git a/java/bulk-import/bulk-import-runner/src/main/java/sleeper/bulkimport/runner/BulkImportJobDriver.java b/java/bulk-import/bulk-import-runner/src/main/java/sleeper/bulkimport/runner/BulkImportJobDriver.java index a7dd9817e02..ab012789b4d 100644 --- a/java/bulk-import/bulk-import-runner/src/main/java/sleeper/bulkimport/runner/BulkImportJobDriver.java +++ b/java/bulk-import/bulk-import-runner/src/main/java/sleeper/bulkimport/runner/BulkImportJobDriver.java @@ -236,7 +236,8 @@ private static void startOrThrow(String[] args, BulkImportJobRunner runner) thro 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); @@ -261,8 +262,7 @@ private static void startOrThrow(String[] args, BulkImportJobRunner runner) thro throw e; } - String key = "bulk_import/" + jobId + "-" + jobRunId + ".json"; - BulkImportJob bulkImportJob = BulkImportJobLoaderFromS3.loadJob(instanceProperties, key, s3Client); + BulkImportJob bulkImportJob = BulkImportJobLoaderFromS3.loadJob(instanceProperties, jobFileObjectKey, s3Client); TablePropertiesProvider tablePropertiesProvider = S3TableProperties.createProvider(instanceProperties, s3Client, dynamoClient); StateStoreProvider stateStoreProvider = StateStoreFactory.createProvider(instanceProperties, s3Client, dynamoClient); diff --git a/java/bulk-import/bulk-import-starter/src/main/java/sleeper/bulkimport/starter/executor/BulkImportArguments.java b/java/bulk-import/bulk-import-starter/src/main/java/sleeper/bulkimport/starter/executor/BulkImportArguments.java index 6d796e23e09..d6cc3f44bb6 100644 --- a/java/bulk-import/bulk-import-starter/src/main/java/sleeper/bulkimport/starter/executor/BulkImportArguments.java +++ b/java/bulk-import/bulk-import-starter/src/main/java/sleeper/bulkimport/starter/executor/BulkImportArguments.java @@ -36,11 +36,13 @@ 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; + jobFileObjectKey = builder.jobFileObjectKey; jobRunId = builder.jobRunId; } @@ -66,7 +68,7 @@ private List sparkSubmitCommandForCluster(String taskId, String jarLocat return Stream.of( Stream.of("spark-submit", "--deploy-mode", "cluster"), sparkSubmitParameters(baseSparkConfig), - Stream.of(jarLocation, configBucket, jobId, taskId, jobRunId, bulkImportMode)) + Stream.of(jarLocation, configBucket, jobId, taskId, jobRunId, jobFileObjectKey, bulkImportMode)) .flatMap(partialArgs -> partialArgs) .collect(Collectors.toUnmodifiableList()); } @@ -123,9 +125,31 @@ 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() { @@ -141,6 +165,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; diff --git a/java/bulk-import/bulk-import-starter/src/main/java/sleeper/bulkimport/starter/executor/BulkImportExecutor.java b/java/bulk-import/bulk-import-starter/src/main/java/sleeper/bulkimport/starter/executor/BulkImportExecutor.java index e14cd933ffb..3426a2eb13f 100644 --- a/java/bulk-import/bulk-import-starter/src/main/java/sleeper/bulkimport/starter/executor/BulkImportExecutor.java +++ b/java/bulk-import/bulk-import-starter/src/main/java/sleeper/bulkimport/starter/executor/BulkImportExecutor.java @@ -75,12 +75,12 @@ public void runJob(BulkImportJob bulkImportJob, String jobRunId) { .jobRunId(jobRunId).build()); try { LOGGER.info("Writing job with id {} to JSON file", bulkImportJob.getId()); - String key = "bulk_import/" + bulkImportJob.getId() + "-" + jobRunId + ".json"; - writeJobToBucket.writeJobToBulkImportBucket(bulkImportJob, key); + String jobFileObjectKey = "bulk_import/" + bulkImportJob.getId() + "-" + jobRunId + ".json"; + 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) { diff --git a/java/bulk-import/bulk-import-starter/src/test/java/sleeper/bulkimport/starter/executor/BulkImportArgumentsTest.java b/java/bulk-import/bulk-import-starter/src/test/java/sleeper/bulkimport/starter/executor/BulkImportArgumentsTest.java index d04a5da285f..839d9b4b2d8 100644 --- a/java/bulk-import/bulk-import-starter/src/test/java/sleeper/bulkimport/starter/executor/BulkImportArgumentsTest.java +++ b/java/bulk-import/bulk-import-starter/src/test/java/sleeper/bulkimport/starter/executor/BulkImportArgumentsTest.java @@ -47,6 +47,7 @@ void shouldConstructArgs() { .files(Lists.newArrayList("file1.parquet")) .build()) .jobRunId("test-run") + .jobFileObjectKey("folder/job.json") .build(); // When / Then @@ -61,6 +62,7 @@ void shouldConstructArgs() { "my-job", "test-task", "test-run", + "folder/job.json", "EMR"); } } diff --git a/java/bulk-import/bulk-import-starter/src/test/java/sleeper/bulkimport/starter/executor/BulkImportExecutorTest.java b/java/bulk-import/bulk-import-starter/src/test/java/sleeper/bulkimport/starter/executor/BulkImportExecutorTest.java index 5b98f489181..f92d33a6b43 100644 --- a/java/bulk-import/bulk-import-starter/src/test/java/sleeper/bulkimport/starter/executor/BulkImportExecutorTest.java +++ b/java/bulk-import/bulk-import-starter/src/test/java/sleeper/bulkimport/starter/executor/BulkImportExecutorTest.java @@ -38,7 +38,9 @@ import java.time.Instant; import java.util.ArrayList; +import java.util.HashMap; import java.util.List; +import java.util.Map; import java.util.UUID; import java.util.function.Supplier; @@ -65,10 +67,8 @@ class BulkImportExecutorTest { private final String bucketName = UUID.randomUUID().toString(); private final String tableId = tableProperties.get(TABLE_ID); private final IngestJobTracker tracker = new InMemoryIngestJobTracker(); - private final List jobsInBucket = new ArrayList<>(); - private final List keysInBucket = new ArrayList<>(); - private final List jobsRun = new ArrayList<>(); - private final List jobRunIdsOfJobsRun = new ArrayList<>(); + private final Map objectKeyToJobFile = new HashMap<>(); + private final List runJobInvocations = new ArrayList<>(); @Nested @DisplayName("Failing validation") @@ -86,8 +86,8 @@ void shouldFailValidationIfFileListIsEmpty() { executor(atTime(validationTime)).runJob(importJob); // Then - assertThat(jobsInBucket).isEmpty(); - assertThat(jobsRun).isEmpty(); + assertThat(objectKeyToJobFile).isEmpty(); + assertThat(runJobInvocations).isEmpty(); assertThat(tracker.getAllJobs(tableId)) .usingRecursiveFieldByFieldElementComparator(IGNORE_UPDATE_TIMES) .containsExactly(ingestJobStatus(importJob.toIngestJob(), @@ -108,8 +108,8 @@ void shouldFailValidationIfFileListIsNull() { executor(atTime(validationTime)).runJob(importJob); // Then - assertThat(jobsInBucket).isEmpty(); - assertThat(jobsRun).isEmpty(); + assertThat(objectKeyToJobFile).isEmpty(); + assertThat(runJobInvocations).isEmpty(); assertThat(tracker.getAllJobs(tableId)) .usingRecursiveFieldByFieldElementComparator(IGNORE_UPDATE_TIMES) .containsExactly(ingestJobStatus(importJob.toIngestJob(), @@ -130,8 +130,8 @@ void shouldFailValidationIfJobIdContainsMoreThan63Characters() { executor(atTime(validationTime)).runJob(importJob); // Then - assertThat(jobsInBucket).isEmpty(); - assertThat(jobsRun).isEmpty(); + assertThat(objectKeyToJobFile).isEmpty(); + assertThat(runJobInvocations).isEmpty(); assertThat(tracker.getAllJobs(tableId)) .usingRecursiveFieldByFieldElementComparator(IGNORE_UPDATE_TIMES) .containsExactly(ingestJobStatus(importJob.toIngestJob(), @@ -152,8 +152,8 @@ void shouldFailValidationIfJobIdContainsUppercaseLetters() { executor(atTime(validationTime)).runJob(importJob); // Then - assertThat(jobsInBucket).isEmpty(); - assertThat(jobsRun).isEmpty(); + assertThat(objectKeyToJobFile).isEmpty(); + assertThat(runJobInvocations).isEmpty(); assertThat(tracker.getAllJobs(tableId)) .usingRecursiveFieldByFieldElementComparator(IGNORE_UPDATE_TIMES) .containsExactly(ingestJobStatus(importJob.toIngestJob(), @@ -178,10 +178,13 @@ void shouldCallRunOnPlatformIfJobIsValid() { executor(atTime(validationTime)).runJob(importJob, "job-run-id"); // Then - assertThat(jobsInBucket).containsExactly(importJob); - assertThat(keysInBucket).containsExactly("bulk_import/my-job-job-run-id.json"); - assertThat(jobsRun).containsExactly(importJob); - assertThat(jobRunIdsOfJobsRun).containsExactly("job-run-id"); + assertThat(objectKeyToJobFile).isEqualTo(Map.of("bulk_import/my-job-job-run-id.json", importJob)); + assertThat(runJobInvocations).containsExactly(BulkImportArguments.builder() + .instanceProperties(instanceProperties) + .bulkImportJob(importJob) + .jobRunId("job-run-id") + .jobFileObjectKey("bulk_import/my-job-job-run-id.json") + .build()); assertThat(tracker.getAllJobs(tableId)) .usingRecursiveFieldByFieldElementComparator(IGNORE_UPDATE_TIMES) .containsExactly(ingestJobStatus(importJob.toIngestJob(), @@ -189,7 +192,7 @@ void shouldCallRunOnPlatformIfJobIsValid() { } @Test - void shouldSucceedIfS3ObjectIsADirectoryContainingFiles() { + void shouldSucceedIfInputFileIsADirectoryContainingFiles() { // Given BulkImportJob importJob = jobForTable() .id("my-job") @@ -201,10 +204,13 @@ void shouldSucceedIfS3ObjectIsADirectoryContainingFiles() { executor(atTime(validationTime)).runJob(importJob, "job-run-id"); // Then - assertThat(jobsInBucket).containsExactly(importJob); - assertThat(keysInBucket).containsExactly("bulk_import/my-job-job-run-id.json"); - assertThat(jobsRun).containsExactly(importJob); - assertThat(jobRunIdsOfJobsRun).containsExactly("job-run-id"); + assertThat(objectKeyToJobFile.values()).containsExactly(importJob); + assertThat(runJobInvocations).containsExactly(BulkImportArguments.builder() + .instanceProperties(instanceProperties) + .bulkImportJob(importJob) + .jobRunId("job-run-id") + .jobFileObjectKey("bulk_import/my-job-job-run-id.json") + .build()); assertThat(tracker.getAllJobs(tableId)) .usingRecursiveFieldByFieldElementComparator(IGNORE_UPDATE_TIMES) .containsExactly(ingestJobStatus(importJob.toIngestJob(), @@ -217,8 +223,8 @@ void shouldDoNothingWhenJobIsNull() { executor(noTimes()).runJob(null); // Then - assertThat(jobsInBucket).isEmpty(); - assertThat(jobsRun).isEmpty(); + assertThat(objectKeyToJobFile).isEmpty(); + assertThat(runJobInvocations).isEmpty(); } @Test @@ -242,7 +248,7 @@ void shouldFailJobRunWhenWriteToBucketFails() { writeJobToBucketFails(failure), recordPlatformExecutor(), atTimes(validationTime, failureTime)) .runJob(importJob, "some-job-run")) .isSameAs(failure); - assertThat(jobsRun).isEmpty(); + assertThat(runJobInvocations).isEmpty(); assertThat(tracker.getAllJobs(tableId)) .usingRecursiveFieldByFieldElementComparator(IGNORE_UPDATE_TIMES) .containsExactly(ingestJobStatus(importJob.toIngestJob(), @@ -271,7 +277,7 @@ void shouldFailJobRunWhenPlatformExecutorFails() { recordWriteJobToBucket(), platformExecutorFails(failure), atTimes(validationTime, failureTime)) .runJob(importJob, "some-job-run")) .isSameAs(failure); - assertThat(jobsInBucket).contains(importJob); + assertThat(objectKeyToJobFile.values()).contains(importJob); assertThat(tracker.getAllJobs(tableId)) .usingRecursiveFieldByFieldElementComparator(IGNORE_UPDATE_TIMES) .containsExactly(ingestJobStatus(importJob.toIngestJob(), @@ -330,16 +336,14 @@ private Supplier noTimes() { private class RecordPlatformExecutor implements PlatformExecutor { @Override public void runJobOnPlatform(BulkImportArguments arguments) { - jobsRun.add(arguments.getBulkImportJob()); - jobRunIdsOfJobsRun.add(arguments.getJobRunId()); + runJobInvocations.add(arguments); } } private class RecordWriteJobToBucket implements WriteJobToBucket { @Override public void writeJobToBulkImportBucket(BulkImportJob bulkImportJob, String objectKey) { - jobsInBucket.add(bulkImportJob); - keysInBucket.add(objectKey); + objectKeyToJobFile.put(objectKey, bulkImportJob); } } diff --git a/java/bulk-import/bulk-import-starter/src/test/resources/example/emr/runjobflow-request.json b/java/bulk-import/bulk-import-starter/src/test/resources/example/emr/runjobflow-request.json index 8b0e4e38aba..40d01d883be 100644 --- a/java/bulk-import/bulk-import-starter/src/test/resources/example/emr/runjobflow-request.json +++ b/java/bulk-import/bulk-import-starter/src/test/resources/example/emr/runjobflow-request.json @@ -226,6 +226,7 @@ "test-job", "sleeper-test-instance-table-id-test-job-EMR", "test-run", + "bulk_import/test-job-test-run.json", "EMR" ] } diff --git a/java/bulk-import/bulk-import-starter/src/test/resources/example/persistent-emr/addjobflow-request.json b/java/bulk-import/bulk-import-starter/src/test/resources/example/persistent-emr/addjobflow-request.json index 17e4cb7f8d4..6352c36ae78 100644 --- a/java/bulk-import/bulk-import-starter/src/test/resources/example/persistent-emr/addjobflow-request.json +++ b/java/bulk-import/bulk-import-starter/src/test/resources/example/persistent-emr/addjobflow-request.json @@ -17,6 +17,7 @@ "test-job", "test-cluster", "test-run", + "bulk_import/test-job-test-run.json", "EMR" ] } diff --git a/java/bulk-import/bulk-import-starter/src/test/resources/example/step-functions/startexecution-input.json b/java/bulk-import/bulk-import-starter/src/test/resources/example/step-functions/startexecution-input.json index 9b6a7045251..9eba93df122 100644 --- a/java/bulk-import/bulk-import-starter/src/test/resources/example/step-functions/startexecution-input.json +++ b/java/bulk-import/bulk-import-starter/src/test/resources/example/step-functions/startexecution-input.json @@ -94,6 +94,7 @@ "test-job", "state-machine-arn", "test-job-run", + "bulk_import/test-job-test-job-run.json", "EKS" ], "jobPodPrefix": "job-test-job", From c858cb21c24cb320025e485abc8f4ea128feb21f Mon Sep 17 00:00:00 2001 From: patchwork01 <110390516+patchwork01@users.noreply.github.com> Date: Fri, 17 Jul 2026 14:43:57 +0000 Subject: [PATCH 03/16] 7607 Inline job ID in BulkImportJobLoaderFromS3IT --- .../sleeper/bulkimport/runner/BulkImportJobLoaderFromS3IT.java | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/java/bulk-import/bulk-import-runner/src/test/java/sleeper/bulkimport/runner/BulkImportJobLoaderFromS3IT.java b/java/bulk-import/bulk-import-runner/src/test/java/sleeper/bulkimport/runner/BulkImportJobLoaderFromS3IT.java index 52925f862d0..fa83195c4f1 100644 --- a/java/bulk-import/bulk-import-runner/src/test/java/sleeper/bulkimport/runner/BulkImportJobLoaderFromS3IT.java +++ b/java/bulk-import/bulk-import-runner/src/test/java/sleeper/bulkimport/runner/BulkImportJobLoaderFromS3IT.java @@ -44,10 +44,9 @@ void setUp() { void shouldLoadBulkImportJobFromS3() { // Given String objectKey = "folder/test.json"; - String jobId = "load-job-id"; BulkImportJob bulkImportJob = BulkImportJob.builder() - .id(jobId) + .id("load-job-id") .tableId("test-table-id") .files(List.of("/load-job.parquet")) .build(); From e6311a61967dc53f0cb0aaae6844c7761ac5cd44 Mon Sep 17 00:00:00 2001 From: patchwork01 <110390516+patchwork01@users.noreply.github.com> Date: Fri, 17 Jul 2026 14:52:12 +0000 Subject: [PATCH 04/16] 7607 Split EKS bulk import system test --- .../suite/fixtures/SystemTestInstance.java | 16 +++- ...ImportST.java => EksAutoBulkImportST.java} | 20 +--- .../suite/EksFargateBulkImportST.java | 91 +++++++++++++++++++ 3 files changed, 107 insertions(+), 20 deletions(-) rename java/system-test/system-test-suite/src/test/java/sleeper/systemtest/suite/{EksBulkImportST.java => EksAutoBulkImportST.java} (80%) create mode 100644 java/system-test/system-test-suite/src/test/java/sleeper/systemtest/suite/EksFargateBulkImportST.java diff --git a/java/system-test/system-test-suite/src/main/java/sleeper/systemtest/suite/fixtures/SystemTestInstance.java b/java/system-test/system-test-suite/src/main/java/sleeper/systemtest/suite/fixtures/SystemTestInstance.java index 692975eec35..d6e1e9151b2 100644 --- a/java/system-test/system-test-suite/src/main/java/sleeper/systemtest/suite/fixtures/SystemTestInstance.java +++ b/java/system-test/system-test-suite/src/main/java/sleeper/systemtest/suite/fixtures/SystemTestInstance.java @@ -103,7 +103,8 @@ private SystemTestInstance() { .disableSchedules(Set.of(COMPACTION_TASK_CREATION)) .build(); public static final SystemTestInstanceConfiguration BULK_IMPORT_PERFORMANCE = usingSystemTestDefaults("emr", SystemTestInstance::createBulkImportPerformanceConfiguration); - public static final SystemTestInstanceConfiguration BULK_IMPORT_EKS = usingSystemTestDefaults("bi-eks", SystemTestInstance::createBulkImportOnEksConfiguration); + public static final SystemTestInstanceConfiguration BULK_IMPORT_EKS_FARGATE = usingSystemTestDefaults("bi-eks", SystemTestInstance::createBulkImportOnEksFargateConfiguration); + public static final SystemTestInstanceConfiguration BULK_IMPORT_EKS_AUTO = usingSystemTestDefaults("bi-eka", SystemTestInstance::createBulkImportOnEksAutoConfiguration); public static final SystemTestInstanceConfiguration BULK_IMPORT_PERFORMANCE_EKS = usingSystemTestDefaults("eksprf", SystemTestInstance::createBulkImportOnEksPerformanceConfiguration); public static final SystemTestInstanceConfiguration BULK_IMPORT_PERSISTENT_EMR = usingSystemTestDefaults("emrpst", SystemTestInstance::createBulkImportOnPersistentEmrConfiguration); public static final SystemTestInstanceConfiguration PARALLEL_COMPACTIONS = usingSystemTestDefaults("cptpll", SystemTestInstance::createCompactionInParallelConfiguration); @@ -230,10 +231,19 @@ private static SleeperInstanceConfiguration createBulkImportPerformanceConfigura return createInstanceConfiguration(properties); } - private static SleeperInstanceConfiguration createBulkImportOnEksConfiguration() { + private static SleeperInstanceConfiguration createBulkImportOnEksFargateConfiguration() { InstanceProperties properties = createInstanceProperties(); properties.setList(OPTIONAL_STACKS, List.of()); - setSystemTestTags(properties, "bulkImportOnEks", "Sleeper Maven system test bulk import on EKS"); + properties.set(BULK_IMPORT_EKS_CLUSTER_TYPE, EksClusterType.FARGATE.toString()); + setSystemTestTags(properties, "bulkImportOnEksFargate", "Sleeper Maven system test bulk import on EKS w/Fargate"); + return createInstanceConfiguration(properties); + } + + private static SleeperInstanceConfiguration createBulkImportOnEksAutoConfiguration() { + InstanceProperties properties = createInstanceProperties(); + properties.setList(OPTIONAL_STACKS, List.of()); + properties.set(BULK_IMPORT_EKS_CLUSTER_TYPE, EksClusterType.AUTOMODE.toString()); + setSystemTestTags(properties, "bulkImportOnEksAuto", "Sleeper Maven system test bulk import on EKS Auto Mode"); return createInstanceConfiguration(properties); } diff --git a/java/system-test/system-test-suite/src/test/java/sleeper/systemtest/suite/EksBulkImportST.java b/java/system-test/system-test-suite/src/test/java/sleeper/systemtest/suite/EksAutoBulkImportST.java similarity index 80% rename from java/system-test/system-test-suite/src/test/java/sleeper/systemtest/suite/EksBulkImportST.java rename to java/system-test/system-test-suite/src/test/java/sleeper/systemtest/suite/EksAutoBulkImportST.java index 5b96b00539a..db925d8ba9e 100644 --- a/java/system-test/system-test-suite/src/test/java/sleeper/systemtest/suite/EksBulkImportST.java +++ b/java/system-test/system-test-suite/src/test/java/sleeper/systemtest/suite/EksAutoBulkImportST.java @@ -36,45 +36,31 @@ import static org.assertj.core.api.Assertions.assertThat; import static sleeper.core.properties.instance.CdkDefinedInstanceProperty.BULK_IMPORT_EKS_JOB_QUEUE_URL; -import static sleeper.core.properties.instance.CommonProperty.LOG_RETENTION_IN_DAYS; import static sleeper.core.properties.table.TableProperty.BULK_IMPORT_MIN_LEAF_PARTITION_COUNT; import static sleeper.core.properties.table.TableProperty.PARTITION_SPLIT_MIN_ROWS; -import static sleeper.systemtest.suite.fixtures.SystemTestInstance.BULK_IMPORT_EKS; +import static sleeper.systemtest.suite.fixtures.SystemTestInstance.BULK_IMPORT_EKS_AUTO; @SystemTest @Slow2 // Slow because it needs to do two CDK deployments, one to add the EKS cluster and one to remove it. // Each CDK deployment takes around 20 minutes. // If we left the EKS cluster around, there would be extra costs as the control pane is persistent. -public class EksBulkImportST { +public class EksAutoBulkImportST { @BeforeEach void setUp(SleeperDsl sleeper, AfterTestReports reporting, SystemTestParameters parameters) { - if (parameters.isInstancePropertyOverridden(LOG_RETENTION_IN_DAYS)) { - return; - } - sleeper.connectToInstanceAddOnlineTable(BULK_IMPORT_EKS); + sleeper.connectToInstanceAddOnlineTable(BULK_IMPORT_EKS_AUTO); sleeper.enableOptionalStack(OptionalStack.EksBulkImportStack); reporting.reportIfTestFailed(SystemTestReports.SystemTestBuilder::ingestJobs); } @AfterEach void tearDown(SleeperDsl sleeper, SystemTestParameters parameters) { - if (parameters.isInstancePropertyOverridden(LOG_RETENTION_IN_DAYS)) { - return; - } sleeper.disableOptionalStack(OptionalStack.EksBulkImportStack); } @Test void shouldPreSplitPartitionTreeAndBulkImport(SleeperDsl sleeper, SystemTestParameters parameters) { - // This is intended to ignore this test when running in an environment where log retention must not be set. - // This is because we're currently unable to prevent an EKS cluster deployment from creating log groups with log - // retention set. See the following issue: - // https://github.com/gchq/sleeper/issues/3451 (Logs retention policy is not applied to all EKS cluster resources) - if (parameters.isInstancePropertyOverridden(LOG_RETENTION_IN_DAYS)) { - return; - } // Given sleeper.updateTableProperties(Map.of( BULK_IMPORT_MIN_LEAF_PARTITION_COUNT, "8", diff --git a/java/system-test/system-test-suite/src/test/java/sleeper/systemtest/suite/EksFargateBulkImportST.java b/java/system-test/system-test-suite/src/test/java/sleeper/systemtest/suite/EksFargateBulkImportST.java new file mode 100644 index 00000000000..0c2500e358f --- /dev/null +++ b/java/system-test/system-test-suite/src/test/java/sleeper/systemtest/suite/EksFargateBulkImportST.java @@ -0,0 +1,91 @@ +/* + * Copyright 2022-2026 Crown Copyright + * + * Licensed 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 sleeper.systemtest.suite; + +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; + +import sleeper.core.properties.model.OptionalStack; +import sleeper.core.row.Row; +import sleeper.core.util.PollWithRetries; +import sleeper.systemtest.dsl.SleeperDsl; +import sleeper.systemtest.dsl.extension.AfterTestReports; +import sleeper.systemtest.dsl.instance.SystemTestParameters; +import sleeper.systemtest.dsl.reporting.SystemTestReports; +import sleeper.systemtest.dsl.util.SystemTestSchema; +import sleeper.systemtest.suite.testutil.SystemTest; +import sleeper.systemtest.suite.testutil.parallel.Slow2; + +import java.time.Duration; +import java.util.Map; + +import static org.assertj.core.api.Assertions.assertThat; +import static sleeper.core.properties.instance.CdkDefinedInstanceProperty.BULK_IMPORT_EKS_JOB_QUEUE_URL; +import static sleeper.core.properties.table.TableProperty.BULK_IMPORT_MIN_LEAF_PARTITION_COUNT; +import static sleeper.core.properties.table.TableProperty.PARTITION_SPLIT_MIN_ROWS; +import static sleeper.systemtest.suite.fixtures.SystemTestInstance.BULK_IMPORT_EKS_FARGATE; + +@SystemTest +@Slow2 +// Slow because it needs to do two CDK deployments, one to add the EKS cluster and one to remove it. +// Each CDK deployment takes around 20 minutes. +// If we left the EKS cluster around, there would be extra costs as the control pane is persistent. +public class EksFargateBulkImportST { + + @BeforeEach + void setUp(SleeperDsl sleeper, AfterTestReports reporting, SystemTestParameters parameters) { + sleeper.connectToInstanceAddOnlineTable(BULK_IMPORT_EKS_FARGATE); + sleeper.enableOptionalStack(OptionalStack.EksBulkImportStack); + reporting.reportIfTestFailed(SystemTestReports.SystemTestBuilder::ingestJobs); + } + + @AfterEach + void tearDown(SleeperDsl sleeper, SystemTestParameters parameters) { + sleeper.disableOptionalStack(OptionalStack.EksBulkImportStack); + } + + @Test + void shouldPreSplitPartitionTreeAndBulkImport(SleeperDsl sleeper, SystemTestParameters parameters) { + // Given + sleeper.updateTableProperties(Map.of( + BULK_IMPORT_MIN_LEAF_PARTITION_COUNT, "8", + PARTITION_SPLIT_MIN_ROWS, "100")); + Iterable rows = sleeper.generateNumberedRows().iterableOverRange(0, 10_000); + sleeper.sourceFiles().create("test.parquet", rows); + + // When + sleeper.ingest().bulkImportByQueue() + .sendSourceFiles(BULK_IMPORT_EKS_JOB_QUEUE_URL, "test.parquet") + .waitForJobs(); + + // Then + assertThat(SystemTestSchema.sorted(sleeper.directQuery().allRowsInTable())) + .containsExactlyElementsOf(rows); + assertThat(sleeper.partitioning().tree().getLeafPartitions()) + .hasSize(8); + assertThat(sleeper.tableFiles().references()) + .hasSize(8); + assertThat(sleeper.eksBulkImportCheck().waitUntilExecutionsFinishedGetStatuses( + PollWithRetries.intervalAndPollingTimeout(Duration.ofSeconds(10), Duration.ofMinutes(5)))) + .containsOnly("SUCCEEDED"); + assertThat(sleeper.eksBulkImportCheck().getPods()) + .isEmpty(); + assertThat(sleeper.eksBulkImportCheck().getJobs()) + .isEmpty(); + } +} From bb800fef9116584803d8d16aa4229252e52210fa Mon Sep 17 00:00:00 2001 From: patchwork01 <110390516+patchwork01@users.noreply.github.com> Date: Fri, 17 Jul 2026 15:00:00 +0000 Subject: [PATCH 05/16] 7607 Update system-test-suites.md --- docs/development/system-test-suites.md | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/docs/development/system-test-suites.md b/docs/development/system-test-suites.md index 8dc31bcdba9..ac402bdff04 100644 --- a/docs/development/system-test-suites.md +++ b/docs/development/system-test-suites.md @@ -9,8 +9,8 @@ it takes to complete the nightly system tests. | AutoDeleteS3ObjectsST | CompactionOnEC2ST | MultipleTablesST | | AutoStopEcsTaskST | ECSStateStoreCommitterST | StateStoreCommitterThroughputST | | CompactionCreationST | ECSStateStoreCommitterThroughputST | -| EmrPersistentBulkImportST | EksBulkImportST | -| OptionalFeaturesDisabledST | +| EmrPersistentBulkImportST | EksAutoBulkImportST | +| OptionalFeaturesDisabledST | EksFargateBulkImportST | | RedeployOptionalStacksST | | Expensive1 | Expensive2 | Expensive3 | From ebba103e774b8a880ccb381b96d7ac65bfce8638 Mon Sep 17 00:00:00 2001 From: patchwork01 <110390516+patchwork01@users.noreply.github.com> Date: Mon, 20 Jul 2026 10:21:10 +0000 Subject: [PATCH 06/16] 7607 Add missing brackets in BulkImportArguments --- .../bulkimport/starter/executor/BulkImportArguments.java | 6 ++++-- 1 file changed, 4 insertions(+), 2 deletions(-) diff --git a/java/bulk-import/bulk-import-starter/src/main/java/sleeper/bulkimport/starter/executor/BulkImportArguments.java b/java/bulk-import/bulk-import-starter/src/main/java/sleeper/bulkimport/starter/executor/BulkImportArguments.java index d6cc3f44bb6..82537a9864d 100644 --- a/java/bulk-import/bulk-import-starter/src/main/java/sleeper/bulkimport/starter/executor/BulkImportArguments.java +++ b/java/bulk-import/bulk-import-starter/src/main/java/sleeper/bulkimport/starter/executor/BulkImportArguments.java @@ -137,10 +137,12 @@ public int hashCode() { @Override public boolean equals(Object obj) { - if (this == obj) + if (this == obj) { return true; - if (!(obj instanceof BulkImportArguments)) + } + 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); From 191bfeade5b8df55b6c8cb9a6e555a1feb687db2 Mon Sep 17 00:00:00 2001 From: patchwork01 <110390516+patchwork01@users.noreply.github.com> Date: Mon, 20 Jul 2026 10:23:25 +0000 Subject: [PATCH 07/16] 7607 Fix GenerateSystemTestSuiteDocumentationIT --- .../documentation/GenerateSystemTestSuiteDocumentationIT.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/java/system-test/system-test-suite/src/test/java/sleeper/systemtest/suite/documentation/GenerateSystemTestSuiteDocumentationIT.java b/java/system-test/system-test-suite/src/test/java/sleeper/systemtest/suite/documentation/GenerateSystemTestSuiteDocumentationIT.java index 091a399b994..a42e420a50e 100644 --- a/java/system-test/system-test-suite/src/test/java/sleeper/systemtest/suite/documentation/GenerateSystemTestSuiteDocumentationIT.java +++ b/java/system-test/system-test-suite/src/test/java/sleeper/systemtest/suite/documentation/GenerateSystemTestSuiteDocumentationIT.java @@ -120,7 +120,7 @@ private String readDocumentation() throws Exception { private static Stream parallelSystemTests() { return Stream.of( Arguments.of("Slow1", "AutoStopEcsTaskST"), - Arguments.of("Slow2", "EksBulkImportST"), + Arguments.of("Slow2", "EksAutoBulkImportST"), Arguments.of("Slow3", "MultipleTablesST"), Arguments.of("Expensive1", "CompactionDataFusionPerformanceST"), Arguments.of("Expensive2", "CompactionPerformanceST"), From 495eb2f313c63dfab7aada729c2a91461cc81ca9 Mon Sep 17 00:00:00 2001 From: patchwork01 <110390516+patchwork01@users.noreply.github.com> Date: Mon, 20 Jul 2026 13:40:58 +0000 Subject: [PATCH 08/16] 7607 Fix EMR Serverless arguments --- .../starter/executor/BulkImportArguments.java | 15 ++++++++++++--- .../executor/EmrServerlessPlatformExecutor.java | 15 +++------------ .../EmrServerlessPlatformExecutorWiremockIT.java | 2 ++ .../example/emr-serverless/jobrun-request.json | 1 + 4 files changed, 18 insertions(+), 15 deletions(-) diff --git a/java/bulk-import/bulk-import-starter/src/main/java/sleeper/bulkimport/starter/executor/BulkImportArguments.java b/java/bulk-import/bulk-import-starter/src/main/java/sleeper/bulkimport/starter/executor/BulkImportArguments.java index 82537a9864d..f524b83e061 100644 --- a/java/bulk-import/bulk-import-starter/src/main/java/sleeper/bulkimport/starter/executor/BulkImportArguments.java +++ b/java/bulk-import/bulk-import-starter/src/main/java/sleeper/bulkimport/starter/executor/BulkImportArguments.java @@ -30,6 +30,7 @@ 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 { @@ -63,16 +64,24 @@ public List sparkSubmitCommandForEKSCluster(String taskId, String jarLoc } private List sparkSubmitCommandForCluster(String taskId, String jarLocation, Map 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, jobFileObjectKey, bulkImportMode)) + Stream.of(jarLocation), + streamEntryPointArguments(taskId, bulkImportMode)) .flatMap(partialArgs -> partialArgs) .collect(Collectors.toUnmodifiableList()); } + public String[] entryPointArgumentsForServerless() { + return streamEntryPointArguments(instanceProperties.get(BULK_IMPORT_EMR_SERVERLESS_CLUSTER_NAME) + "-EMRS", "EMR") + .toArray(String[]::new); + } + + private Stream streamEntryPointArguments(String taskId, String bulkImportMode) { + return Stream.of(instanceProperties.get(CONFIG_BUCKET), bulkImportJob.getId(), taskId, jobRunId, jobFileObjectKey, bulkImportMode); + } + public String sparkSubmitParametersForServerless() { return sparkSubmitParameters( SparkConfigurationUtils.getSparkServerlessConfigurationFromInstanceProperties( diff --git a/java/bulk-import/bulk-import-starter/src/main/java/sleeper/bulkimport/starter/executor/EmrServerlessPlatformExecutor.java b/java/bulk-import/bulk-import-starter/src/main/java/sleeper/bulkimport/starter/executor/EmrServerlessPlatformExecutor.java index ec33b5aa093..1509d2ab896 100644 --- a/java/bulk-import/bulk-import-starter/src/main/java/sleeper/bulkimport/starter/executor/EmrServerlessPlatformExecutor.java +++ b/java/bulk-import/bulk-import-starter/src/main/java/sleeper/bulkimport/starter/executor/EmrServerlessPlatformExecutor.java @@ -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; @@ -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( diff --git a/java/bulk-import/bulk-import-starter/src/test/java/sleeper/bulkimport/starter/executor/EmrServerlessPlatformExecutorWiremockIT.java b/java/bulk-import/bulk-import-starter/src/test/java/sleeper/bulkimport/starter/executor/EmrServerlessPlatformExecutorWiremockIT.java index a22f5b4955c..6e8887fd79b 100644 --- a/java/bulk-import/bulk-import-starter/src/test/java/sleeper/bulkimport/starter/executor/EmrServerlessPlatformExecutorWiremockIT.java +++ b/java/bulk-import/bulk-import-starter/src/test/java/sleeper/bulkimport/starter/executor/EmrServerlessPlatformExecutorWiremockIT.java @@ -79,6 +79,7 @@ void shouldRunAServerlessJob(WireMockRuntimeInfo runtimeInfo) { BulkImportArguments arguments = BulkImportArguments.builder() .instanceProperties(instanceProperties) .bulkImportJob(job).jobRunId("run-id") + .jobFileObjectKey("folder/job.json") .build(); stubFor(post("/applications/application-id/jobruns").willReturn(aResponse().withStatus(200))); @@ -110,6 +111,7 @@ void shouldRunAServerlessJobWithJobSparkConfSet(WireMockRuntimeInfo runtimeInfo) BulkImportArguments arguments = BulkImportArguments.builder() .instanceProperties(instanceProperties) .bulkImportJob(job).jobRunId("run-id") + .jobFileObjectKey("folder/job.json") .build(); stubFor(post("/applications/application-id/jobruns").willReturn(aResponse().withStatus(200))); diff --git a/java/bulk-import/bulk-import-starter/src/test/resources/example/emr-serverless/jobrun-request.json b/java/bulk-import/bulk-import-starter/src/test/resources/example/emr-serverless/jobrun-request.json index 96bd9aca0dc..c79e0c1cdc1 100644 --- a/java/bulk-import/bulk-import-starter/src/test/resources/example/emr-serverless/jobrun-request.json +++ b/java/bulk-import/bulk-import-starter/src/test/resources/example/emr-serverless/jobrun-request.json @@ -8,6 +8,7 @@ "my-job", "my-application-EMRS", "run-id", + "folder/job.json", "EMR" ] } From 55f9e935c131688735f2a145902169ebb24f2f7f Mon Sep 17 00:00:00 2001 From: patchwork01 <110390516+patchwork01@users.noreply.github.com> Date: Mon, 20 Jul 2026 13:52:16 +0000 Subject: [PATCH 09/16] 7607 Refactor entry point arguments --- .../bulkimport/starter/executor/BulkImportArguments.java | 9 ++++----- 1 file changed, 4 insertions(+), 5 deletions(-) diff --git a/java/bulk-import/bulk-import-starter/src/main/java/sleeper/bulkimport/starter/executor/BulkImportArguments.java b/java/bulk-import/bulk-import-starter/src/main/java/sleeper/bulkimport/starter/executor/BulkImportArguments.java index f524b83e061..5d66869a8cd 100644 --- a/java/bulk-import/bulk-import-starter/src/main/java/sleeper/bulkimport/starter/executor/BulkImportArguments.java +++ b/java/bulk-import/bulk-import-starter/src/main/java/sleeper/bulkimport/starter/executor/BulkImportArguments.java @@ -68,18 +68,17 @@ private List sparkSubmitCommandForCluster(String taskId, String jarLocat Stream.of("spark-submit", "--deploy-mode", "cluster"), sparkSubmitParameters(baseSparkConfig), Stream.of(jarLocation), - streamEntryPointArguments(taskId, bulkImportMode)) + Stream.of(entryPointArguments(taskId, bulkImportMode))) .flatMap(partialArgs -> partialArgs) .collect(Collectors.toUnmodifiableList()); } public String[] entryPointArgumentsForServerless() { - return streamEntryPointArguments(instanceProperties.get(BULK_IMPORT_EMR_SERVERLESS_CLUSTER_NAME) + "-EMRS", "EMR") - .toArray(String[]::new); + return entryPointArguments(instanceProperties.get(BULK_IMPORT_EMR_SERVERLESS_CLUSTER_NAME) + "-EMRS", "EMR"); } - private Stream streamEntryPointArguments(String taskId, String bulkImportMode) { - return Stream.of(instanceProperties.get(CONFIG_BUCKET), bulkImportJob.getId(), taskId, jobRunId, jobFileObjectKey, bulkImportMode); + private String[] entryPointArguments(String taskId, String bulkImportMode) { + return new String[]{instanceProperties.get(CONFIG_BUCKET), bulkImportJob.getId(), taskId, jobRunId, jobFileObjectKey, bulkImportMode}; } public String sparkSubmitParametersForServerless() { From d9f9debe8d7c29bb3cef6d311ce60e00ee88a4ad Mon Sep 17 00:00:00 2001 From: patchwork01 <110390516+patchwork01@users.noreply.github.com> Date: Mon, 20 Jul 2026 14:14:05 +0000 Subject: [PATCH 10/16] 7607 Move EksAutoBulkImportST to Slow3 suite --- docs/development/system-test-suites.md | 10 +++++----- .../sleeper/systemtest/suite/EksAutoBulkImportST.java | 4 ++-- 2 files changed, 7 insertions(+), 7 deletions(-) diff --git a/docs/development/system-test-suites.md b/docs/development/system-test-suites.md index ac402bdff04..42f9d9bfc46 100644 --- a/docs/development/system-test-suites.md +++ b/docs/development/system-test-suites.md @@ -6,11 +6,11 @@ it takes to complete the nightly system tests. | Slow1 | Slow2 | Slow3 | |----------------------------|------------------------------------|---------------------------------| -| AutoDeleteS3ObjectsST | CompactionOnEC2ST | MultipleTablesST | -| AutoStopEcsTaskST | ECSStateStoreCommitterST | StateStoreCommitterThroughputST | -| CompactionCreationST | ECSStateStoreCommitterThroughputST | -| EmrPersistentBulkImportST | EksAutoBulkImportST | -| OptionalFeaturesDisabledST | EksFargateBulkImportST | +| AutoDeleteS3ObjectsST | CompactionOnEC2ST | EksAutoBulkImportST | +| AutoStopEcsTaskST | ECSStateStoreCommitterST | MultipleTablesST | +| CompactionCreationST | ECSStateStoreCommitterThroughputST | StateStoreCommitterThroughputST | +| EmrPersistentBulkImportST | EksFargateBulkImportST | +| OptionalFeaturesDisabledST | | RedeployOptionalStacksST | | Expensive1 | Expensive2 | Expensive3 | diff --git a/java/system-test/system-test-suite/src/test/java/sleeper/systemtest/suite/EksAutoBulkImportST.java b/java/system-test/system-test-suite/src/test/java/sleeper/systemtest/suite/EksAutoBulkImportST.java index db925d8ba9e..bb646c40730 100644 --- a/java/system-test/system-test-suite/src/test/java/sleeper/systemtest/suite/EksAutoBulkImportST.java +++ b/java/system-test/system-test-suite/src/test/java/sleeper/systemtest/suite/EksAutoBulkImportST.java @@ -29,7 +29,7 @@ import sleeper.systemtest.dsl.reporting.SystemTestReports; import sleeper.systemtest.dsl.util.SystemTestSchema; import sleeper.systemtest.suite.testutil.SystemTest; -import sleeper.systemtest.suite.testutil.parallel.Slow2; +import sleeper.systemtest.suite.testutil.parallel.Slow3; import java.time.Duration; import java.util.Map; @@ -41,7 +41,7 @@ import static sleeper.systemtest.suite.fixtures.SystemTestInstance.BULK_IMPORT_EKS_AUTO; @SystemTest -@Slow2 +@Slow3 // Slow because it needs to do two CDK deployments, one to add the EKS cluster and one to remove it. // Each CDK deployment takes around 20 minutes. // If we left the EKS cluster around, there would be extra costs as the control pane is persistent. From 01f49e7be1c0cd2d7a18d6b18f9f34eca4d77b19 Mon Sep 17 00:00:00 2001 From: patchwork01 <110390516+patchwork01@users.noreply.github.com> Date: Mon, 20 Jul 2026 14:17:22 +0000 Subject: [PATCH 11/16] 7607 Swap IngestPerformanceST and CompactionPerformanceST between suites --- docs/development/system-test-suites.md | 10 +++++----- .../systemtest/suite/CompactionPerformanceST.java | 4 ++-- .../sleeper/systemtest/suite/IngestPerformanceST.java | 4 ++-- 3 files changed, 9 insertions(+), 9 deletions(-) diff --git a/docs/development/system-test-suites.md b/docs/development/system-test-suites.md index 42f9d9bfc46..42eba999515 100644 --- a/docs/development/system-test-suites.md +++ b/docs/development/system-test-suites.md @@ -13,8 +13,8 @@ it takes to complete the nightly system tests. | OptionalFeaturesDisabledST | | RedeployOptionalStacksST | -| Expensive1 | Expensive2 | Expensive3 | -|-----------------------------------|----------------------------|-----------------------| -| CompactionDataFusionPerformanceST | CompactionPerformanceST | IngestPerformanceST | -| CompactionVeryLargeST | EksBulkImportPerformanceST | ParallelCompactionsST | -| | EmrBulkImportPerformanceST | +| Expensive1 | Expensive2 | Expensive3 | +|-----------------------------------|----------------------------|-------------------------| +| CompactionDataFusionPerformanceST | EksBulkImportPerformanceST | CompactionPerformanceST | +| CompactionVeryLargeST | EmrBulkImportPerformanceST | ParallelCompactionsST | +| | IngestPerformanceST | diff --git a/java/system-test/system-test-suite/src/test/java/sleeper/systemtest/suite/CompactionPerformanceST.java b/java/system-test/system-test-suite/src/test/java/sleeper/systemtest/suite/CompactionPerformanceST.java index afdb3e18ba1..3a4da95a5d4 100644 --- a/java/system-test/system-test-suite/src/test/java/sleeper/systemtest/suite/CompactionPerformanceST.java +++ b/java/system-test/system-test-suite/src/test/java/sleeper/systemtest/suite/CompactionPerformanceST.java @@ -27,7 +27,7 @@ import sleeper.systemtest.dsl.extension.AfterTestReports; import sleeper.systemtest.dsl.reporting.SystemTestReports; import sleeper.systemtest.suite.testutil.SystemTest; -import sleeper.systemtest.suite.testutil.parallel.Expensive2; +import sleeper.systemtest.suite.testutil.parallel.Expensive3; import java.time.Duration; import java.util.Map; @@ -41,7 +41,7 @@ @SystemTest // Expensive because it takes a long time to compact this many rows on fairly large ECS instances. -@Expensive2 +@Expensive3 public class CompactionPerformanceST { @BeforeEach diff --git a/java/system-test/system-test-suite/src/test/java/sleeper/systemtest/suite/IngestPerformanceST.java b/java/system-test/system-test-suite/src/test/java/sleeper/systemtest/suite/IngestPerformanceST.java index 075ce41233b..6f801fca777 100644 --- a/java/system-test/system-test-suite/src/test/java/sleeper/systemtest/suite/IngestPerformanceST.java +++ b/java/system-test/system-test-suite/src/test/java/sleeper/systemtest/suite/IngestPerformanceST.java @@ -24,7 +24,7 @@ import sleeper.systemtest.dsl.extension.AfterTestReports; import sleeper.systemtest.dsl.reporting.SystemTestReports; import sleeper.systemtest.suite.testutil.SystemTest; -import sleeper.systemtest.suite.testutil.parallel.Expensive3; +import sleeper.systemtest.suite.testutil.parallel.Expensive2; import java.time.Duration; @@ -37,7 +37,7 @@ @SystemTest // Expensive because it takes a long time to ingest this many rows on fairly large ECS instances. -@Expensive3 +@Expensive2 public class IngestPerformanceST { @BeforeEach From b36613c0ae1cfc44d6a75814da7fa50c3f6861a7 Mon Sep 17 00:00:00 2001 From: patchwork01 <110390516+patchwork01@users.noreply.github.com> Date: Mon, 20 Jul 2026 14:40:40 +0000 Subject: [PATCH 12/16] 7607 Fix GenerateSystemTestSuiteDocumentationIT --- .../GenerateSystemTestSuiteDocumentationIT.java | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/java/system-test/system-test-suite/src/test/java/sleeper/systemtest/suite/documentation/GenerateSystemTestSuiteDocumentationIT.java b/java/system-test/system-test-suite/src/test/java/sleeper/systemtest/suite/documentation/GenerateSystemTestSuiteDocumentationIT.java index a42e420a50e..8dbd1cff238 100644 --- a/java/system-test/system-test-suite/src/test/java/sleeper/systemtest/suite/documentation/GenerateSystemTestSuiteDocumentationIT.java +++ b/java/system-test/system-test-suite/src/test/java/sleeper/systemtest/suite/documentation/GenerateSystemTestSuiteDocumentationIT.java @@ -120,11 +120,11 @@ private String readDocumentation() throws Exception { private static Stream parallelSystemTests() { return Stream.of( Arguments.of("Slow1", "AutoStopEcsTaskST"), - Arguments.of("Slow2", "EksAutoBulkImportST"), + Arguments.of("Slow2", "EksFargateBulkImportST"), Arguments.of("Slow3", "MultipleTablesST"), Arguments.of("Expensive1", "CompactionDataFusionPerformanceST"), - Arguments.of("Expensive2", "CompactionPerformanceST"), - Arguments.of("Expensive3", "IngestPerformanceST")); + Arguments.of("Expensive2", "IngestPerformanceST"), + Arguments.of("Expensive3", "CompactionPerformanceST")); } private static String columnContaining(String documentation, String systemTestName) { From 7c2caa0f757ed07f7556ba57303dc82b21b2c530 Mon Sep 17 00:00:00 2001 From: patchwork01 <110390516+patchwork01@users.noreply.github.com> Date: Tue, 4 Aug 2026 08:02:09 +0000 Subject: [PATCH 13/16] 7607 Fix usage check of BulkImportJobDriver --- .../java/sleeper/bulkimport/runner/BulkImportJobDriver.java | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/java/bulk-import/bulk-import-runner/src/main/java/sleeper/bulkimport/runner/BulkImportJobDriver.java b/java/bulk-import/bulk-import-runner/src/main/java/sleeper/bulkimport/runner/BulkImportJobDriver.java index ab012789b4d..07d055b234c 100644 --- a/java/bulk-import/bulk-import-runner/src/main/java/sleeper/bulkimport/runner/BulkImportJobDriver.java +++ b/java/bulk-import/bulk-import-runner/src/main/java/sleeper/bulkimport/runner/BulkImportJobDriver.java @@ -228,9 +228,9 @@ 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:" + - " "); + if (args.length != 6) { + throw new IllegalArgumentException("Expected 6 arguments:" + + " "); } String configBucket = args[0]; String jobId = args[1]; From b56f60e8ff2072009f54efbb6df6d2512bdd45a1 Mon Sep 17 00:00:00 2001 From: patchwork01 <110390516+patchwork01@users.noreply.github.com> Date: Tue, 4 Aug 2026 08:46:31 +0000 Subject: [PATCH 14/16] 7607 Adjust comment on EksBulkImportPerformanceST --- .../sleeper/systemtest/suite/EksBulkImportPerformanceST.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/java/system-test/system-test-suite/src/test/java/sleeper/systemtest/suite/EksBulkImportPerformanceST.java b/java/system-test/system-test-suite/src/test/java/sleeper/systemtest/suite/EksBulkImportPerformanceST.java index 08e000e8d1d..25f5d5b1b8c 100644 --- a/java/system-test/system-test-suite/src/test/java/sleeper/systemtest/suite/EksBulkImportPerformanceST.java +++ b/java/system-test/system-test-suite/src/test/java/sleeper/systemtest/suite/EksBulkImportPerformanceST.java @@ -38,7 +38,7 @@ import static sleeper.systemtest.suite.testutil.FileReferenceSystemTestHelper.numberOfRowsIn; @SystemTest -// Expensive because it takes a lot of very costly EKS instances to import this many rows. +// Expensive because it takes a lot of very costly EC2 instances to import this many rows. @Expensive2 public class EksBulkImportPerformanceST { @BeforeEach From 09535d6afce41d61d9e156bcdb723268fb6c5041 Mon Sep 17 00:00:00 2001 From: patchwork01 <110390516+patchwork01@users.noreply.github.com> Date: Tue, 4 Aug 2026 08:52:56 +0000 Subject: [PATCH 15/16] 7607 Fix DirectEmrServerlessDriver --- .../starter/executor/BulkImportArguments.java | 8 ++++---- .../starter/executor/BulkImportExecutor.java | 6 +++++- .../ingest/DirectEmrServerlessDriver.java | 18 +++++++----------- 3 files changed, 16 insertions(+), 16 deletions(-) diff --git a/java/bulk-import/bulk-import-starter/src/main/java/sleeper/bulkimport/starter/executor/BulkImportArguments.java b/java/bulk-import/bulk-import-starter/src/main/java/sleeper/bulkimport/starter/executor/BulkImportArguments.java index 5d66869a8cd..b9d5bb23158 100644 --- a/java/bulk-import/bulk-import-starter/src/main/java/sleeper/bulkimport/starter/executor/BulkImportArguments.java +++ b/java/bulk-import/bulk-import-starter/src/main/java/sleeper/bulkimport/starter/executor/BulkImportArguments.java @@ -41,10 +41,10 @@ public class BulkImportArguments { private final String jobRunId; private BulkImportArguments(Builder builder) { - instanceProperties = builder.instanceProperties; - bulkImportJob = builder.bulkImportJob; - jobFileObjectKey = builder.jobFileObjectKey; - 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() { diff --git a/java/bulk-import/bulk-import-starter/src/main/java/sleeper/bulkimport/starter/executor/BulkImportExecutor.java b/java/bulk-import/bulk-import-starter/src/main/java/sleeper/bulkimport/starter/executor/BulkImportExecutor.java index 3426a2eb13f..9fd3240fbb4 100644 --- a/java/bulk-import/bulk-import-starter/src/main/java/sleeper/bulkimport/starter/executor/BulkImportExecutor.java +++ b/java/bulk-import/bulk-import-starter/src/main/java/sleeper/bulkimport/starter/executor/BulkImportExecutor.java @@ -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()); } @@ -75,7 +79,7 @@ public void runJob(BulkImportJob bulkImportJob, String jobRunId) { .jobRunId(jobRunId).build()); try { LOGGER.info("Writing job with id {} to JSON file", bulkImportJob.getId()); - String jobFileObjectKey = "bulk_import/" + bulkImportJob.getId() + "-" + jobRunId + ".json"; + String jobFileObjectKey = createJobFileObjectKey(bulkImportJob, jobRunId); writeJobToBucket.writeJobToBulkImportBucket(bulkImportJob, jobFileObjectKey); LOGGER.info("Submitting job with id {}", bulkImportJob.getId()); platformExecutor.runJobOnPlatform(BulkImportArguments.builder() diff --git a/java/system-test/system-test-drivers/src/main/java/sleeper/systemtest/drivers/ingest/DirectEmrServerlessDriver.java b/java/system-test/system-test-drivers/src/main/java/sleeper/systemtest/drivers/ingest/DirectEmrServerlessDriver.java index f3670762e8b..57e5c5a49a5 100644 --- a/java/system-test/system-test-drivers/src/main/java/sleeper/systemtest/drivers/ingest/DirectEmrServerlessDriver.java +++ b/java/system-test/system-test-drivers/src/main/java/sleeper/systemtest/drivers/ingest/DirectEmrServerlessDriver.java @@ -16,15 +16,14 @@ package sleeper.systemtest.drivers.ingest; -import software.amazon.awssdk.core.sync.RequestBody; import software.amazon.awssdk.services.dynamodb.DynamoDbClient; import software.amazon.awssdk.services.emrserverless.EmrServerlessClient; import software.amazon.awssdk.services.s3.S3Client; -import software.amazon.awssdk.services.s3.model.PutObjectRequest; import sleeper.bulkimport.core.job.BulkImportJob; -import sleeper.bulkimport.core.job.BulkImportJobSerDe; import sleeper.bulkimport.starter.executor.BulkImportArguments; +import sleeper.bulkimport.starter.executor.BulkImportExecutor; +import sleeper.bulkimport.starter.executor.BulkImportJobWriterToS3; import sleeper.bulkimport.starter.executor.EmrServerlessPlatformExecutor; import sleeper.core.tracker.ingest.job.IngestJobTracker; import sleeper.ingest.tracker.job.IngestJobTrackerFactory; @@ -35,8 +34,6 @@ import java.time.Instant; import java.util.UUID; -import static sleeper.core.properties.instance.CdkDefinedInstanceProperty.BULK_IMPORT_BUCKET; - public class DirectEmrServerlessDriver implements DirectBulkImportDriver { private final SystemTestInstanceContext instance; private final S3Client s3Client; @@ -54,15 +51,14 @@ public void sendJob(BulkImportJob job) { String jobRunId = UUID.randomUUID().toString(); jobTracker().jobValidated(job.toIngestJob().acceptedEventBuilder(Instant.now()) .jobRunId(jobRunId).build()); - s3Client.putObject(PutObjectRequest.builder() - .bucket(instance.getInstanceProperties().get(BULK_IMPORT_BUCKET)) - .key("bulk_import/" + job.getId() + "-" + jobRunId + ".json") - .build(), - RequestBody.fromString(new BulkImportJobSerDe().toJson(job))); + String jobFileObjectKey = BulkImportExecutor.createJobFileObjectKey(job, jobRunId); + new BulkImportJobWriterToS3(instance.getInstanceProperties(), s3Client) + .writeJobToBulkImportBucket(job, jobFileObjectKey); executor().runJobOnPlatform(BulkImportArguments.builder() + .instanceProperties(instance.getInstanceProperties()) .bulkImportJob(job) .jobRunId(jobRunId) - .instanceProperties(instance.getInstanceProperties()) + .jobFileObjectKey(jobFileObjectKey) .build()); } From 531e20588880c78546f34fc9a7bf24ba69a2222c Mon Sep 17 00:00:00 2001 From: patchwork01 <110390516+patchwork01@users.noreply.github.com> Date: Tue, 4 Aug 2026 09:04:46 +0000 Subject: [PATCH 16/16] 7607 Deploy EKS system test instances with EKS enabled --- .../systemtest/suite/fixtures/SystemTestInstance.java | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/java/system-test/system-test-suite/src/main/java/sleeper/systemtest/suite/fixtures/SystemTestInstance.java b/java/system-test/system-test-suite/src/main/java/sleeper/systemtest/suite/fixtures/SystemTestInstance.java index d6e1e9151b2..954db3f6b7e 100644 --- a/java/system-test/system-test-suite/src/main/java/sleeper/systemtest/suite/fixtures/SystemTestInstance.java +++ b/java/system-test/system-test-suite/src/main/java/sleeper/systemtest/suite/fixtures/SystemTestInstance.java @@ -233,7 +233,7 @@ private static SleeperInstanceConfiguration createBulkImportPerformanceConfigura private static SleeperInstanceConfiguration createBulkImportOnEksFargateConfiguration() { InstanceProperties properties = createInstanceProperties(); - properties.setList(OPTIONAL_STACKS, List.of()); + properties.setEnumList(OPTIONAL_STACKS, List.of(OptionalStack.EksBulkImportStack)); properties.set(BULK_IMPORT_EKS_CLUSTER_TYPE, EksClusterType.FARGATE.toString()); setSystemTestTags(properties, "bulkImportOnEksFargate", "Sleeper Maven system test bulk import on EKS w/Fargate"); return createInstanceConfiguration(properties); @@ -241,7 +241,7 @@ private static SleeperInstanceConfiguration createBulkImportOnEksFargateConfigur private static SleeperInstanceConfiguration createBulkImportOnEksAutoConfiguration() { InstanceProperties properties = createInstanceProperties(); - properties.setList(OPTIONAL_STACKS, List.of()); + properties.setEnumList(OPTIONAL_STACKS, List.of(OptionalStack.EksBulkImportStack)); properties.set(BULK_IMPORT_EKS_CLUSTER_TYPE, EksClusterType.AUTOMODE.toString()); setSystemTestTags(properties, "bulkImportOnEksAuto", "Sleeper Maven system test bulk import on EKS Auto Mode"); return createInstanceConfiguration(properties); @@ -249,7 +249,7 @@ private static SleeperInstanceConfiguration createBulkImportOnEksAutoConfigurati private static SleeperInstanceConfiguration createBulkImportOnEksPerformanceConfiguration() { InstanceProperties properties = createInstancePropertiesWithDefaults(); - properties.setList(OPTIONAL_STACKS, List.of()); + properties.setEnumList(OPTIONAL_STACKS, List.of(OptionalStack.EksBulkImportStack)); properties.set(BULK_IMPORT_EKS_CLUSTER_TYPE, EksClusterType.AUTOMODE.toString()); setSystemTestTags(properties, "bulkImportPerformanceOnEks", "Sleeper Maven system test bulk import performance on EKS"); return createInstanceConfiguration(properties);