Skip to content
Draft
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
Original file line number Diff line number Diff line change
Expand Up @@ -51,6 +51,7 @@
import sleeper.core.schema.type.ListType;
import sleeper.core.schema.type.LongType;
import sleeper.core.schema.type.MapType;
import sleeper.core.schema.type.PrimitiveType;
import sleeper.core.schema.type.StringType;
import sleeper.core.statestore.FileReference;
import sleeper.core.statestore.StateStore;
Expand Down Expand Up @@ -107,14 +108,28 @@

class BulkImportJobDriverIT {

private static Stream<Arguments> getStreamOfBulkImportJobRunners() {
private static Stream<Named<BulkImportJobRunner>> getBulkImportJobRunners() {
return Stream.of(
Arguments.of(Named.of("BulkImportJobDataframeDriver",
(BulkImportJobRunner) BulkImportJobDataframeDriver::createFileReferences)),
Arguments.of(Named.of("BulkImportJobRDDDriver",
(BulkImportJobRunner) BulkImportJobRDDDriver::createFileReferences)),
Arguments.of(Named.of("BulkImportDataframeLocalSortDriver",
(BulkImportJobRunner) BulkImportDataframeLocalSortDriver::createFileReferences)));
Named.of("BulkImportJobDataframeDriver",
(BulkImportJobRunner) BulkImportJobDataframeDriver::createFileReferences),
Named.of("BulkImportJobRDDDriver",
(BulkImportJobRunner) BulkImportJobRDDDriver::createFileReferences),
Named.of("BulkImportDataframeLocalSortDriver",
(BulkImportJobRunner) BulkImportDataframeLocalSortDriver::createFileReferences));
}

private static Stream<Arguments> getStreamOfBulkImportJobRunners() {
return getBulkImportJobRunners().map(Arguments::of);
}

private static Stream<Arguments> getStreamOfBulkImportJobRunnersAndKeyTypes() {
return getBulkImportJobRunners().flatMap(runner -> Stream.of(
Arguments.of(runner, Named.of("LongType", new KeyTypeTestData(
new LongType(), 1L, 2L))),
Arguments.of(runner, Named.of("StringType", new KeyTypeTestData(
new StringType(), "A", "B"))),
Arguments.of(runner, Named.of("ByteArrayType", new KeyTypeTestData(
new ByteArrayType(), new byte[]{1}, new byte[]{2})))));
}

@TempDir
Expand Down Expand Up @@ -193,6 +208,25 @@ void shouldImportDataSinglePartition(BulkImportJobRunner runner) throws Exceptio
ingestFinishedStatus(summary(startTime, endTime, 200, 200), 1))));
}

@ParameterizedTest
@MethodSource("getStreamOfBulkImportJobRunnersAndKeyTypes")
void shouldImportDataWithSupportedRowAndSortKeyTypes(
BulkImportJobRunner runner, KeyTypeTestData keyType) throws Exception {
// Given
tableProperties.setSchema(getSchemaWithKeyType(keyType.type()));
update(stateStore()).initialise(tableProperties);
List<Row> rows = getRows(keyType);
writeRowsToFile(rows, dataDir + "/import/a.parquet");

// When
BulkImportJob job = jobForTable().id("my-job")
.files(List.of(dataDir + "/import/a.parquet")).build();
runJob(runner, job);

// Then
assertThat(readRowsInPartitionTreeOrder()).isEqualTo(sorted(rows));
}

@ParameterizedTest
@MethodSource("getStreamOfBulkImportJobRunners")
void shouldImportDataSinglePartitionIdenticalRowKeyDifferentSortKeys(BulkImportJobRunner runner) throws Exception {
Expand Down Expand Up @@ -458,6 +492,29 @@ private static Schema getSchema() {
.build();
}

private static Schema getSchemaWithKeyType(PrimitiveType keyType) {
return Schema.builder()
.rowKeyFields(new Field("key", keyType))
.sortKeyFields(new Field("sort", keyType))
.valueFields(new Field("value", new StringType()))
.build();
}

private static List<Row> getRows(KeyTypeTestData keyType) {
return List.of(
row(keyType.higherValue(), keyType.higherValue(), "higher key"),
row(keyType.lowerValue(), keyType.higherValue(), "higher sort"),
row(keyType.lowerValue(), keyType.lowerValue(), "lower sort"));
}

private static Row row(Object key, Object sort, String value) {
Row row = new Row();
row.put("key", key);
row.put("sort", sort);
row.put("value", value);
return row;
}

private static List<Row> getRows() {
List<Row> rows = new ArrayList<>(200);
for (int i = 0; i < 100; i++) {
Expand Down Expand Up @@ -571,4 +628,7 @@ private String createDir(String name) {
}
return path.toString();
}

private record KeyTypeTestData(PrimitiveType type, Object lowerValue, Object higherValue) {
}
}
Loading