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]

Reply via email to