FrankChen021 commented on code in PR #19534: URL: https://github.com/apache/druid/pull/19534#discussion_r4111068858
########## extensions-contrib/druid-iceberg-extensions/src/main/java/org/apache/druid/iceberg/input/IcebergFileTaskInputSource.java: ########## @@ -0,0 +1,288 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you 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 org.apache.druid.iceberg.input; + +import com.fasterxml.jackson.annotation.JsonCreator; +import com.fasterxml.jackson.annotation.JsonProperty; +import com.fasterxml.jackson.annotation.JsonTypeName; +import org.apache.druid.data.input.InputFormat; +import org.apache.druid.data.input.InputRowSchema; +import org.apache.druid.data.input.InputSource; +import org.apache.druid.data.input.InputSourceFactory; +import org.apache.druid.data.input.InputSourceReader; +import org.apache.druid.data.input.InputSplit; +import org.apache.iceberg.DeleteFile; + +import javax.annotation.Nullable; +import java.io.File; +import java.util.Collections; +import java.util.List; + +/** + * A single-split {@link InputSource} representing one Iceberg V2 file scan task (one data file + * plus its associated delete files). Instances of this class are created by + * {@link IcebergInputSource#withSplit(InputSplit)} when the split carries a {@code "v2"} marker, + * and they are embedded in per-task specs that are distributed to Druid workers. + * + * <p>The table schema is carried as a JSON string (serialized on the coordinator at planning + * time via {@code SchemaParser.toJson}). Workers deserialize it without contacting the catalog. + * + * <p>The {@code icebergCatalog} field is retained as a fallback source of {@link org.apache.iceberg.io.FileIO} + * for non-local storage paths (S3, HDFS) in Phase 1. For local filesystem paths + * {@link WarehouseFileIO} is used instead. + * + */ +@JsonTypeName(IcebergFileTaskInputSource.TYPE_KEY) +public class IcebergFileTaskInputSource implements InputSource +{ + public static final String TYPE_KEY = "icebergFileTask"; + + @JsonProperty + private final String dataFilePath; + + @JsonProperty + private final String fileFormat; + + @JsonProperty + private final long dataFileSizeInBytes; + + @JsonProperty + private final long dataFileRecordCount; + + @JsonProperty + private final List<DeleteFileInfo> deleteFiles; + + /** JSON-serialized Iceberg Schema captured at coordinator planning time. */ + @JsonProperty + private final String tableSchemaJson; + + @JsonProperty + private final String tableNamespace; + + @JsonProperty + private final String tableName; + + /** + * Catalog retained as a fallback FileIO source for non-local storage paths (S3, HDFS). + * For local filesystem paths {@link WarehouseFileIO} is used instead. + */ + @JsonProperty + private final IcebergCatalog icebergCatalog; + + @JsonProperty + private final InputSourceFactory warehouseSource; + + @JsonCreator + public IcebergFileTaskInputSource( + @JsonProperty("dataFilePath") String dataFilePath, + @JsonProperty("fileFormat") String fileFormat, + @JsonProperty("dataFileSizeInBytes") long dataFileSizeInBytes, + @JsonProperty("dataFileRecordCount") long dataFileRecordCount, + @JsonProperty("deleteFiles") List<DeleteFileInfo> deleteFiles, + @JsonProperty("tableSchemaJson") String tableSchemaJson, + @JsonProperty("tableNamespace") String tableNamespace, + @JsonProperty("tableName") String tableName, + @JsonProperty("icebergCatalog") IcebergCatalog icebergCatalog, + @JsonProperty("warehouseSource") InputSourceFactory warehouseSource + ) + { + this.dataFilePath = dataFilePath; + this.fileFormat = fileFormat; + this.dataFileSizeInBytes = dataFileSizeInBytes; + this.dataFileRecordCount = dataFileRecordCount; + this.deleteFiles = deleteFiles != null ? deleteFiles : Collections.emptyList(); + this.tableSchemaJson = tableSchemaJson; + this.tableNamespace = tableNamespace; + this.tableName = tableName; + this.icebergCatalog = icebergCatalog; + this.warehouseSource = warehouseSource; + } + + @Override + public boolean needsFormat() + { + // Handles its own reading through Iceberg's native Parquet/ORC reader + return false; Review Comment: [P1] Preserve the configured input format on V2 split workers **Finding:** A parallel-index V2 split is deserialized as this source, and Druid passes an InputFormat to a task only when needsFormat() is true. Returning false therefore makes IcebergNativeRecordReader receive null on every worker, so binaryAsString is silently disabled and flattenSpec fields are silently ignored even though the source-level reader and documentation claim to honor or reject those settings. This makes parallel ingestion produce different dimensions and binary values from single-task ingestion. **Suggestion:** Carry the relevant input-format settings in the split source or make the split source participate in Druid's input-format propagation so the native reader receives the configured format on worker tasks. ########## extensions-contrib/druid-iceberg-extensions/src/main/java/org/apache/druid/iceberg/input/IcebergInputSource.java: ########## @@ -283,6 +309,37 @@ public InputSourceReader reader( File temporaryDirectory ) { + if (!isLoaded) { + retrieveIcebergDatafiles(); + } + if (hasDeleteFiles) { + // V2 path: create a combined reader across all file-scan tasks. + final List<IcebergNativeRecordReader> taskReaders = v2Tasks.stream() + .map(task -> taskToNativeReader(task, inputRowSchema, inputFormat)) + .collect(Collectors.toList()); + return new InputSourceReader() + { + @Override + public CloseableIterator<InputRow> read(InputStats inputStats) throws IOException + { + List<CloseableIterator<InputRow>> iters = new ArrayList<>(); + for (IcebergNativeRecordReader r : taskReaders) { Review Comment: [P2] Open V2 readers lazily instead of all at once **Finding:** The unsplit V2 reader calls read() for every task before returning the concatenated iterator. IcebergNativeRecordReader.read() opens its Parquet iterator during streamRows(), so a table with many data files holds one file handle per data file (and can leak already-open iterators if a later reader fails) until the outer iterator is closed. This can exhaust worker file descriptors for ordinary single-task ingestion. **Suggestion:** Build the concatenated iterator lazily and close each task reader before opening the next one, with cleanup for readers opened before an exception. ########## extensions-contrib/druid-iceberg-extensions/src/test/java/org/apache/druid/iceberg/input/IcebergInputSourceTest.java: ########## @@ -329,13 +346,765 @@ public void testResidualFilterModeFailWithPartitionedTableNonPartitionColumn() t "Expect residual error to be thrown" ); } + /** + * Creates a V2 format table with 3 records, writes a position-delete file that deletes row at + * position 1, and verifies that only 2 rows survive. + */ + @Test + public void testInputSourceV2WithPositionDeletes() throws IOException + { + tearDown(); + String v2TableName = "v2PosDeleteTable"; + tableIdentifier = TableIdentifier.of(Namespace.of(NAMESPACE), v2TableName); + + // Schema includes a timestamp column for InputRowSchema + Schema v2Schema = new Schema( + Types.NestedField.required(1, "id", Types.StringType.get()), + Types.NestedField.required(2, "name", Types.StringType.get()), + Types.NestedField.optional(3, "__time", Types.LongType.get()) + ); + + List<Map<String, Object>> rows = ImmutableList.of( + ImmutableMap.of("id", "1", "name", "Alice", "__time", 0L), + ImmutableMap.of("id", "123988", "name", "Foo", "__time", 1000L), // will be deleted + ImmutableMap.of("id", "3", "name", "Charlie", "__time", 2000L) + ); + + // Create V2 table and write data + Table table = testCatalog.retrieveCatalog().createTable( + tableIdentifier, + v2Schema, + PartitionSpec.unpartitioned(), + ImmutableMap.of(TableProperties.FORMAT_VERSION, "2") + ); + + String dataFilePath = table.location() + "/data/" + UUID.randomUUID() + ".parquet"; + DataFile dataFile = writeParquetDataFile(table, v2Schema, rows, dataFilePath); + table.newAppend().appendFile(dataFile).commit(); + + // Write a position-delete file: delete the row at position 1 (0-indexed) + String posDeletePath = table.location() + "/delete-files/" + UUID.randomUUID() + ".parquet"; + DeleteFile posDeleteFile = writePositionDeleteFile(table, dataFilePath, 1L, posDeletePath); + table.newRowDelta().addDeletes(posDeleteFile).commit(); + + // Read via IcebergInputSource + InputRowSchema inputRowSchema = new InputRowSchema( + new TimestampSpec("__time", "millis", null), + new DimensionsSpec(DimensionsSpec.getDefaultSchemas(ImmutableList.of("id", "name"))), + ColumnsFilter.all() + ); + + IcebergInputSource inputSource = new IcebergInputSource( + v2TableName, + NAMESPACE, + null, + testCatalog, + new LocalInputSourceFactory(), + null, + null, + null, + null + ); + + List<InputRow> result = new ArrayList<>(); + try (CloseableIterator<InputRow> it = inputSource.reader(inputRowSchema, null, FileUtils.createTempDir()).read(null)) { + it.forEachRemaining(result::add); + } + + Assertions.assertEquals(2, result.size(), "Position delete should remove exactly one row"); + List<String> ids = result.stream() + .map(r -> r.getDimension("id").get(0)) + .collect(Collectors.toList()); + Assertions.assertTrue(ids.contains("1"), "Row 'Alice' should survive"); + Assertions.assertFalse(ids.contains("123988"), "Row 'Foo' (id=123988) should be deleted"); + Assertions.assertTrue(ids.contains("3"), "Row 'Charlie' should survive"); + } + + /** + * Creates a V2 format table, writes an equality-delete file that deletes any row where + * {@code id = "123988"}, and verifies the row is absent from the reader output. + */ + @Test + public void testInputSourceV2WithEqualityDeletes() throws IOException + { + tearDown(); + String v2TableName = "v2EqDeleteTable"; + tableIdentifier = TableIdentifier.of(Namespace.of(NAMESPACE), v2TableName); + + Schema v2Schema = new Schema( + Types.NestedField.required(1, "id", Types.StringType.get()), + Types.NestedField.required(2, "name", Types.StringType.get()), + Types.NestedField.optional(3, "__time", Types.LongType.get()) + ); + + List<Map<String, Object>> rows = ImmutableList.of( + ImmutableMap.of("id", "1", "name", "Alice", "__time", 0L), + ImmutableMap.of("id", "123988", "name", "Foo", "__time", 1000L), // will be deleted + ImmutableMap.of("id", "3", "name", "Charlie", "__time", 2000L) + ); + + Table table = testCatalog.retrieveCatalog().createTable( + tableIdentifier, + v2Schema, + PartitionSpec.unpartitioned(), + ImmutableMap.of(TableProperties.FORMAT_VERSION, "2") + ); + + String dataFilePath = table.location() + "/data/" + UUID.randomUUID() + ".parquet"; + DataFile dataFile = writeParquetDataFile(table, v2Schema, rows, dataFilePath); + table.newAppend().appendFile(dataFile).commit(); + + // Equality-delete schema: just the "id" field (field ID 1) + Schema eqDeleteSchema = v2Schema.select(ImmutableList.of("id")); + String eqDeletePath = table.location() + "/delete-files/" + UUID.randomUUID() + ".parquet"; + DeleteFile eqDeleteFile = writeEqualityDeleteFile( + table, + eqDeleteSchema, + ImmutableList.of(1), // field ID for "id" + ImmutableList.of(ImmutableMap.of("id", "123988")), + eqDeletePath + ); + table.newRowDelta().addDeletes(eqDeleteFile).commit(); + + InputRowSchema inputRowSchema = new InputRowSchema( + new TimestampSpec("__time", "millis", null), + new DimensionsSpec(DimensionsSpec.getDefaultSchemas(ImmutableList.of("id", "name"))), + ColumnsFilter.all() + ); + + IcebergInputSource inputSource = new IcebergInputSource( + v2TableName, + NAMESPACE, + null, + testCatalog, + new LocalInputSourceFactory(), + null, + null, + null, + null + ); + + List<InputRow> result = new ArrayList<>(); + try (CloseableIterator<InputRow> it = inputSource.reader(inputRowSchema, null, FileUtils.createTempDir()).read(null)) { + it.forEachRemaining(result::add); + } + + Assertions.assertEquals(2, result.size(), "Equality delete should remove exactly one row"); + List<String> ids = result.stream() + .map(r -> r.getDimension("id").get(0)) + .collect(Collectors.toList()); + Assertions.assertFalse(ids.contains("123988"), "Row with id='123988' should be equality-deleted"); + Assertions.assertTrue(ids.contains("1"), "Other rows should survive"); + Assertions.assertTrue(ids.contains("3"), "Other rows should survive"); + } + + /** + * V2 format table with NO delete files should fall through to the V1 path-based approach. + */ + @Test + public void testInputSourceV2WithNoDeleteFiles() throws IOException + { + tearDown(); + String v2TableName = "v2NoDeleteTable"; + tableIdentifier = TableIdentifier.of(Namespace.of(NAMESPACE), v2TableName); + + Table table = testCatalog.retrieveCatalog().createTable( + tableIdentifier, + tableSchema, + PartitionSpec.unpartitioned(), + ImmutableMap.of(TableProperties.FORMAT_VERSION, "2") + ); + + GenericRecord record = GenericRecord.create(tableSchema); + record.setField("id", "123988"); + record.setField("name", "Foo"); + writeAndCommit(table, tableSchema, ImmutableList.of(record)); + + IcebergInputSource inputSource = new IcebergInputSource( + v2TableName, + NAMESPACE, + null, + testCatalog, + new LocalInputSourceFactory(), + null, + null, + null, + null + ); + + // No delete files → falls through to V1 path → splits are LocalInputSource backed + Stream<InputSplit<List<String>>> splits = inputSource.createSplits(null, new MaxSizeSplitHintSpec(null, null)); + List<InputSource> splitSources = splits.map(inputSource::withSplit).collect(Collectors.toList()); + + Assertions.assertEquals(1, splitSources.size()); + // V1 path returns a LocalInputSource (not IcebergFileTaskInputSource) + Assertions.assertFalse(splitSources.get(0) instanceof IcebergFileTaskInputSource, "V2 table without delete files should use V1 (path-based) path"); + } + + /** + * Unit-level test for the V2 split encoding/decoding contract. + */ + @Test + public void testInputSourceV2SplitEncoding() throws IOException + { + tearDown(); + String v2TableName = "v2SplitEncTable"; + tableIdentifier = TableIdentifier.of(Namespace.of(NAMESPACE), v2TableName); + + Schema v2Schema = new Schema( + Types.NestedField.required(1, "id", Types.StringType.get()), + Types.NestedField.required(2, "name", Types.StringType.get()), + Types.NestedField.optional(3, "__time", Types.LongType.get()) + ); + + List<Map<String, Object>> rows = ImmutableList.of( + ImmutableMap.of("id", "a", "name", "A", "__time", 0L), + ImmutableMap.of("id", "b", "name", "B", "__time", 1000L) + ); + + Table table = testCatalog.retrieveCatalog().createTable( + tableIdentifier, + v2Schema, + PartitionSpec.unpartitioned(), + ImmutableMap.of(TableProperties.FORMAT_VERSION, "2") + ); + + String dataFilePath = table.location() + "/data/" + UUID.randomUUID() + ".parquet"; + DataFile dataFile = writeParquetDataFile(table, v2Schema, rows, dataFilePath); + table.newAppend().appendFile(dataFile).commit(); + + // Add a position delete file so the V2 path is triggered + String posDeletePath = table.location() + "/delete-files/" + UUID.randomUUID() + ".parquet"; + DeleteFile posDeleteFile = writePositionDeleteFile(table, dataFilePath, 0L, posDeletePath); + table.newRowDelta().addDeletes(posDeleteFile).commit(); + + IcebergInputSource inputSource = new IcebergInputSource( + v2TableName, + NAMESPACE, + null, + testCatalog, + new LocalInputSourceFactory(), + null, + null, + null, + null + ); + + // Trigger planning + List<InputSplit<List<String>>> splits = inputSource.createSplits(null, new MaxSizeSplitHintSpec(null, null)) + .collect(Collectors.toList()); + Assertions.assertEquals(1, splits.size()); + + // Verify the split carries the v2 marker + List<String> splitParts = splits.get(0).get(); + Assertions.assertEquals("v2", splitParts.get(0), "First element must be v2 marker"); + // parts[5] must be the table schema JSON + Assertions.assertNotNull(splitParts.get(5), "Split must carry table schema JSON at index 5"); + Assertions.assertTrue(splitParts.get(5).contains("\"type\"") || splitParts.get(5).contains("fields"), "Schema JSON must look like an Iceberg schema object"); + // The POS: entry must appear at index 6 or later + boolean hasPosEntry = splitParts.stream().anyMatch(p -> p.startsWith("POS:")); + Assertions.assertTrue(hasPosEntry, "V2 split with position delete must contain a POS: entry"); + + // withSplit must return IcebergFileTaskInputSource + InputSource splitSource = inputSource.withSplit(splits.get(0)); + Assertions.assertTrue(splitSource instanceof IcebergFileTaskInputSource, "withSplit on a v2 split must return IcebergFileTaskInputSource"); + + // A V1 split (no marker) should return a non-IcebergFileTaskInputSource. + // Not possible to force with the current table, so just verify the encoding roundtrip: + IcebergFileTaskInputSource decodedSource = (IcebergFileTaskInputSource) splitSource; + Assertions.assertEquals(dataFilePath, decodedSource.getDataFilePath()); + Assertions.assertEquals("PARQUET", decodedSource.getFileFormat()); + Assertions.assertNotNull(decodedSource.getTableSchemaJson(), "Decoded source must carry table schema JSON"); + Assertions.assertEquals(1, decodedSource.getDeleteFiles().size()); + Assertions.assertTrue(decodedSource.getDeleteFiles().get(0).isPositionDelete()); + Assertions.assertEquals(NAMESPACE, decodedSource.getTableNamespace()); + Assertions.assertEquals(v2TableName, decodedSource.getTableName()); + } + + /** + * Two sequential position-delete files against the same data file — both must be applied. + */ + @Test + public void testInputSourceV2MultiplePositionDeleteFiles() throws IOException + { + tearDown(); + String v2TableName = "v2MultiPosDeleteTable"; + tableIdentifier = TableIdentifier.of(Namespace.of(NAMESPACE), v2TableName); + + Schema v2Schema = new Schema( + Types.NestedField.required(1, "id", Types.StringType.get()), + Types.NestedField.required(2, "name", Types.StringType.get()), + Types.NestedField.optional(3, "__time", Types.LongType.get()) + ); + + List<Map<String, Object>> rows = ImmutableList.of( + ImmutableMap.of("id", "1", "name", "Alice", "__time", 0L), + ImmutableMap.of("id", "2", "name", "Bob", "__time", 1000L), // deleted by file 1 + ImmutableMap.of("id", "3", "name", "Charlie", "__time", 2000L), + ImmutableMap.of("id", "4", "name", "Dave", "__time", 3000L), // deleted by file 2 + ImmutableMap.of("id", "5", "name", "Eve", "__time", 4000L) + ); + + Table table = testCatalog.retrieveCatalog().createTable( + tableIdentifier, + v2Schema, + PartitionSpec.unpartitioned(), + ImmutableMap.of(TableProperties.FORMAT_VERSION, "2") + ); + + String dataFilePath = table.location() + "/data/" + UUID.randomUUID() + ".parquet"; + DataFile dataFile = writeParquetDataFile(table, v2Schema, rows, dataFilePath); + table.newAppend().appendFile(dataFile).commit(); + + // First delete file: remove row at position 1 (Bob) + String posDeletePath1 = table.location() + "/delete-files/" + UUID.randomUUID() + ".parquet"; + DeleteFile posDeleteFile1 = writePositionDeleteFile(table, dataFilePath, 1L, posDeletePath1); + table.newRowDelta().addDeletes(posDeleteFile1).commit(); + + // Second delete file: remove row at position 3 (Dave) + String posDeletePath2 = table.location() + "/delete-files/" + UUID.randomUUID() + ".parquet"; + DeleteFile posDeleteFile2 = writePositionDeleteFile(table, dataFilePath, 3L, posDeletePath2); + table.newRowDelta().addDeletes(posDeleteFile2).commit(); + + InputRowSchema inputRowSchema = new InputRowSchema( + new TimestampSpec("__time", "millis", null), + new DimensionsSpec(DimensionsSpec.getDefaultSchemas(ImmutableList.of("id", "name"))), + ColumnsFilter.all() + ); + + IcebergInputSource inputSource = new IcebergInputSource( + v2TableName, + NAMESPACE, + null, + testCatalog, + new LocalInputSourceFactory(), + null, + null, + null, + null + ); + + List<InputRow> result = new ArrayList<>(); + try (CloseableIterator<InputRow> it = inputSource.reader(inputRowSchema, null, FileUtils.createTempDir()).read(null)) { + it.forEachRemaining(result::add); + } + + Assertions.assertEquals(3, result.size(), "Both delete files must be applied; 3 rows should survive"); + List<String> ids = result.stream().map(r -> r.getDimension("id").get(0)).collect(Collectors.toList()); + Assertions.assertTrue(ids.contains("1")); + Assertions.assertFalse(ids.contains("2"), "Bob (pos 1) should be deleted by first delete file"); + Assertions.assertTrue(ids.contains("3")); + Assertions.assertFalse(ids.contains("4"), "Dave (pos 3) should be deleted by second delete file"); + Assertions.assertTrue(ids.contains("5")); + } + + /** + * All rows in a data file are position-deleted — reader must return zero rows without error. + */ + @Test + public void testInputSourceV2AllRowsDeleted() throws IOException + { + tearDown(); + String v2TableName = "v2AllDeletedTable"; + tableIdentifier = TableIdentifier.of(Namespace.of(NAMESPACE), v2TableName); + + Schema v2Schema = new Schema( + Types.NestedField.required(1, "id", Types.StringType.get()), + Types.NestedField.required(2, "name", Types.StringType.get()), + Types.NestedField.optional(3, "__time", Types.LongType.get()) + ); + + List<Map<String, Object>> rows = ImmutableList.of( + ImmutableMap.of("id", "1", "name", "Alice", "__time", 0L), + ImmutableMap.of("id", "2", "name", "Bob", "__time", 1000L) + ); + + Table table = testCatalog.retrieveCatalog().createTable( + tableIdentifier, + v2Schema, + PartitionSpec.unpartitioned(), + ImmutableMap.of(TableProperties.FORMAT_VERSION, "2") + ); + + String dataFilePath = table.location() + "/data/" + UUID.randomUUID() + ".parquet"; + DataFile dataFile = writeParquetDataFile(table, v2Schema, rows, dataFilePath); + table.newAppend().appendFile(dataFile).commit(); + + // Delete every row + String posDeletePath = table.location() + "/delete-files/" + UUID.randomUUID() + ".parquet"; + DeleteFile posDeleteFile = writePositionDeleteFile(table, dataFilePath, ImmutableList.of(0L, 1L), posDeletePath); + table.newRowDelta().addDeletes(posDeleteFile).commit(); + + InputRowSchema inputRowSchema = new InputRowSchema( + new TimestampSpec("__time", "millis", null), + new DimensionsSpec(DimensionsSpec.getDefaultSchemas(ImmutableList.of("id", "name"))), + ColumnsFilter.all() + ); + + IcebergInputSource inputSource = new IcebergInputSource( + v2TableName, + NAMESPACE, + null, + testCatalog, + new LocalInputSourceFactory(), + null, + null, + null, + null + ); + + List<InputRow> result = new ArrayList<>(); + try (CloseableIterator<InputRow> it = inputSource.reader(inputRowSchema, null, FileUtils.createTempDir()).read(null)) { + it.forEachRemaining(result::add); + } + + Assertions.assertEquals(0, result.size(), "All rows deleted — reader must return zero rows"); + } + + /** + * Both a position-delete and an equality-delete apply to the same data file. + */ + @Test + public void testInputSourceV2MixedDeleteTypes() throws IOException + { + tearDown(); + String v2TableName = "v2MixedDeleteTable"; + tableIdentifier = TableIdentifier.of(Namespace.of(NAMESPACE), v2TableName); + + Schema v2Schema = new Schema( + Types.NestedField.required(1, "id", Types.StringType.get()), + Types.NestedField.required(2, "name", Types.StringType.get()), + Types.NestedField.optional(3, "__time", Types.LongType.get()) + ); + + List<Map<String, Object>> rows = ImmutableList.of( + ImmutableMap.of("id", "1", "name", "Alice", "__time", 0L), + ImmutableMap.of("id", "2", "name", "Bob", "__time", 1000L), // removed by position delete + ImmutableMap.of("id", "3", "name", "Charlie", "__time", 2000L) // removed by equality delete + ); + + Table table = testCatalog.retrieveCatalog().createTable( + tableIdentifier, + v2Schema, + PartitionSpec.unpartitioned(), + ImmutableMap.of(TableProperties.FORMAT_VERSION, "2") + ); + + String dataFilePath = table.location() + "/data/" + UUID.randomUUID() + ".parquet"; + DataFile dataFile = writeParquetDataFile(table, v2Schema, rows, dataFilePath); + table.newAppend().appendFile(dataFile).commit(); + + // Position delete: remove Bob (position 1) + String posDeletePath = table.location() + "/delete-files/" + UUID.randomUUID() + ".parquet"; + DeleteFile posDeleteFile = writePositionDeleteFile(table, dataFilePath, 1L, posDeletePath); + table.newRowDelta().addDeletes(posDeleteFile).commit(); + + // Equality delete: remove Charlie (id = "3") + Schema eqDeleteSchema = v2Schema.select(ImmutableList.of("id")); + String eqDeletePath = table.location() + "/delete-files/" + UUID.randomUUID() + ".parquet"; + DeleteFile eqDeleteFile = writeEqualityDeleteFile( + table, + eqDeleteSchema, + ImmutableList.of(1), + ImmutableList.of(ImmutableMap.of("id", "3")), + eqDeletePath + ); + table.newRowDelta().addDeletes(eqDeleteFile).commit(); + + InputRowSchema inputRowSchema = new InputRowSchema( + new TimestampSpec("__time", "millis", null), + new DimensionsSpec(DimensionsSpec.getDefaultSchemas(ImmutableList.of("id", "name"))), + ColumnsFilter.all() + ); + + IcebergInputSource inputSource = new IcebergInputSource( + v2TableName, + NAMESPACE, + null, + testCatalog, + new LocalInputSourceFactory(), + null, + null, + null, + null + ); + + List<InputRow> result = new ArrayList<>(); + try (CloseableIterator<InputRow> it = inputSource.reader(inputRowSchema, null, FileUtils.createTempDir()).read(null)) { + it.forEachRemaining(result::add); + } + + Assertions.assertEquals(1, result.size(), "Position and equality delete must both be applied; only Alice survives"); + Assertions.assertEquals("1", result.get(0).getDimension("id").get(0)); + } + + /** + * Two data files each with their own position-delete file — confirms per-split correctness. + */ + @Test + public void testInputSourceV2MultipleDataFilesWithDeletes() throws IOException + { + tearDown(); + String v2TableName = "v2MultiDataFileTable"; + tableIdentifier = TableIdentifier.of(Namespace.of(NAMESPACE), v2TableName); + + Schema v2Schema = new Schema( + Types.NestedField.required(1, "id", Types.StringType.get()), + Types.NestedField.required(2, "name", Types.StringType.get()), + Types.NestedField.optional(3, "__time", Types.LongType.get()) + ); + + Table table = testCatalog.retrieveCatalog().createTable( + tableIdentifier, + v2Schema, + PartitionSpec.unpartitioned(), + ImmutableMap.of(TableProperties.FORMAT_VERSION, "2") + ); + + // Data file 1: [Alice, Bob] — delete Alice (pos 0), keep Bob + List<Map<String, Object>> rows1 = ImmutableList.of( + ImmutableMap.of("id", "1", "name", "Alice", "__time", 0L), + ImmutableMap.of("id", "2", "name", "Bob", "__time", 1000L) + ); + String dataFilePath1 = table.location() + "/data/" + UUID.randomUUID() + ".parquet"; + DataFile dataFile1 = writeParquetDataFile(table, v2Schema, rows1, dataFilePath1); + + // Data file 2: [Charlie, Dave] — delete Dave (pos 1), keep Charlie + List<Map<String, Object>> rows2 = ImmutableList.of( + ImmutableMap.of("id", "3", "name", "Charlie", "__time", 2000L), + ImmutableMap.of("id", "4", "name", "Dave", "__time", 3000L) + ); + String dataFilePath2 = table.location() + "/data/" + UUID.randomUUID() + ".parquet"; + DataFile dataFile2 = writeParquetDataFile(table, v2Schema, rows2, dataFilePath2); + + table.newAppend().appendFile(dataFile1).appendFile(dataFile2).commit(); + + String posDeletePath1 = table.location() + "/delete-files/" + UUID.randomUUID() + ".parquet"; + DeleteFile posDeleteFile1 = writePositionDeleteFile(table, dataFilePath1, 0L, posDeletePath1); + + String posDeletePath2 = table.location() + "/delete-files/" + UUID.randomUUID() + ".parquet"; + DeleteFile posDeleteFile2 = writePositionDeleteFile(table, dataFilePath2, 1L, posDeletePath2); + + table.newRowDelta().addDeletes(posDeleteFile1).addDeletes(posDeleteFile2).commit(); + + InputRowSchema inputRowSchema = new InputRowSchema( + new TimestampSpec("__time", "millis", null), + new DimensionsSpec(DimensionsSpec.getDefaultSchemas(ImmutableList.of("id", "name"))), + ColumnsFilter.all() + ); + + IcebergInputSource inputSource = new IcebergInputSource( + v2TableName, + NAMESPACE, + null, + testCatalog, + new LocalInputSourceFactory(), + null, + null, + null, + null + ); + + List<InputRow> result = new ArrayList<>(); + try (CloseableIterator<InputRow> it = inputSource.reader(inputRowSchema, null, FileUtils.createTempDir()).read(null)) { + it.forEachRemaining(result::add); + } + + Assertions.assertEquals(2, result.size(), "One row deleted per data file; 2 should survive"); + List<String> ids = result.stream().map(r -> r.getDimension("id").get(0)).collect(Collectors.toList()); + Assertions.assertFalse(ids.contains("1"), "Alice (file1, pos 0) should be deleted"); + Assertions.assertTrue(ids.contains("2"), "Bob should survive"); + Assertions.assertTrue(ids.contains("3"), "Charlie should survive"); + Assertions.assertFalse(ids.contains("4"), "Dave (file2, pos 1) should be deleted"); + } + + /** + * On the V2 (delete-aware) path, a configured {@code flattenSpec} with fields can't be applied + * by the native reader, so ingestion must fail fast rather than silently ignore it. + */ + @Test + public void testInputSourceV2RejectsFlattenSpecWithFields() throws IOException + { + tearDown(); + String v2TableName = "v2FlattenRejectTable"; + tableIdentifier = TableIdentifier.of(Namespace.of(NAMESPACE), v2TableName); + + Schema v2Schema = new Schema( + Types.NestedField.required(1, "id", Types.StringType.get()), + Types.NestedField.required(2, "name", Types.StringType.get()), + Types.NestedField.optional(3, "__time", Types.LongType.get()) + ); + + List<Map<String, Object>> rows = ImmutableList.of( + ImmutableMap.of("id", "1", "name", "Alice", "__time", 0L) + ); + + Table table = testCatalog.retrieveCatalog().createTable( + tableIdentifier, + v2Schema, + PartitionSpec.unpartitioned(), + ImmutableMap.of(TableProperties.FORMAT_VERSION, "2") + ); + + String dataFilePath = table.location() + "/data/" + UUID.randomUUID() + ".parquet"; + DataFile dataFile = writeParquetDataFile(table, v2Schema, rows, dataFilePath); + table.newAppend().appendFile(dataFile).commit(); + + // Any delete file, even one that deletes nothing, is enough to route through the V2 path. + String posDeletePath = table.location() + "/delete-files/" + UUID.randomUUID() + ".parquet"; + DeleteFile posDeleteFile = writePositionDeleteFile(table, dataFilePath, 99L, posDeletePath); + table.newRowDelta().addDeletes(posDeleteFile).commit(); + + InputRowSchema inputRowSchema = new InputRowSchema( + new TimestampSpec("__time", "millis", null), + new DimensionsSpec(DimensionsSpec.getDefaultSchemas(ImmutableList.of("id", "name"))), + ColumnsFilter.all() + ); + + IcebergInputSource inputSource = new IcebergInputSource( + v2TableName, + NAMESPACE, + null, + testCatalog, + new LocalInputSourceFactory(), + null, + null, + null, + null + ); + + JSONPathSpec flattenSpec = new JSONPathSpec( + true, + ImmutableList.of(JSONPathFieldSpec.createRootField("name")) + ); + FakeNestedInputFormat inputFormat = new FakeNestedInputFormat(flattenSpec, false); + + DruidException exception = Assertions.assertThrows( + DruidException.class, + () -> inputSource.reader(inputRowSchema, inputFormat, FileUtils.createTempDir()) + ); + Assertions.assertTrue( + exception.getMessage().contains("flattenSpec"), + "Expect flattenSpec error to be thrown" + ); + } + + /** + * On the V2 (delete-aware) path, {@code binaryAsString} is cheap to replicate directly and + * must be honored so binary columns match what the V1 (input-format-based) path would produce. + */ + @Test + public void testInputSourceV2HonorsBinaryAsString() throws IOException + { + tearDown(); + String v2TableName = "v2BinaryAsStringTable"; + tableIdentifier = TableIdentifier.of(Namespace.of(NAMESPACE), v2TableName); + + Schema v2Schema = new Schema( + Types.NestedField.required(1, "id", Types.StringType.get()), + Types.NestedField.optional(2, "payload", Types.BinaryType.get()), + Types.NestedField.optional(3, "__time", Types.LongType.get()) + ); + + byte[] payloadBytes = "hello".getBytes(StandardCharsets.UTF_8); + List<Map<String, Object>> rows = ImmutableList.of( + ImmutableMap.of("id", "1", "payload", ByteBuffer.wrap(payloadBytes), "__time", 0L) + ); + + Table table = testCatalog.retrieveCatalog().createTable( + tableIdentifier, + v2Schema, + PartitionSpec.unpartitioned(), + ImmutableMap.of(TableProperties.FORMAT_VERSION, "2") + ); + + String dataFilePath = table.location() + "/data/" + UUID.randomUUID() + ".parquet"; + DataFile dataFile = writeParquetDataFile(table, v2Schema, rows, dataFilePath); + table.newAppend().appendFile(dataFile).commit(); + + // Any delete file, even one that deletes nothing, is enough to route through the V2 path. + String posDeletePath = table.location() + "/delete-files/" + UUID.randomUUID() + ".parquet"; + DeleteFile posDeleteFile = writePositionDeleteFile(table, dataFilePath, 99L, posDeletePath); + table.newRowDelta().addDeletes(posDeleteFile).commit(); + + InputRowSchema inputRowSchema = new InputRowSchema( + new TimestampSpec("__time", "millis", null), + new DimensionsSpec(DimensionsSpec.getDefaultSchemas(ImmutableList.of("id", "payload"))), + ColumnsFilter.all() + ); + + IcebergInputSource inputSource = new IcebergInputSource( + v2TableName, + NAMESPACE, + null, + testCatalog, + new LocalInputSourceFactory(), + null, + null, + null, + null + ); + + FakeNestedInputFormat inputFormat = new FakeNestedInputFormat(null, true); + + List<InputRow> result = new ArrayList<>(); + try (CloseableIterator<InputRow> it = + inputSource.reader(inputRowSchema, inputFormat, FileUtils.createTempDir()).read(null)) { + it.forEachRemaining(result::add); + } + + Assertions.assertEquals(1, result.size()); + Object payloadValue = ((MapBasedInputRow) result.get(0)).getEvent().get("payload"); + Assertions.assertEquals("hello", payloadValue, "binaryAsString=true should decode binary column as UTF-8 string"); + } + + /** + * Minimal {@link NestedInputFormat} test double that also exposes a Parquet-style + * {@code getBinaryAsString()} accessor, mirroring {@code ParquetInputFormat}'s shape without + * depending on the parquet-extensions module. + */ + private static class FakeNestedInputFormat extends NestedInputFormat Review Comment: [P2] Make the reflective test input format accessible **Finding:** The binaryAsString regression test passes a private FakeNestedInputFormat class to IcebergNativeRecordReader, which invokes getBinaryAsString reflectively from a different top-level class. Method.invoke performs access checks on the declaring class, so this private test double causes IllegalAccessException (wrapped as RE) instead of returning the configured value and makes testInputSourceV2HonorsBinaryAsString fail. **Suggestion:** Use a public test double or avoid reflective access by exposing a shared accessor contract that the reader can call directly. -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
