github-advanced-security[bot] commented on code in PR #19534:
URL: https://github.com/apache/druid/pull/19534#discussion_r4055849025


##########
extensions-contrib/druid-iceberg-extensions/src/main/java/org/apache/druid/iceberg/input/IcebergInputSource.java:
##########
@@ -196,21 +272,172 @@
 
   protected void retrieveIcebergDatafiles()
   {
-    List<String> snapshotDataFiles = icebergCatalog.extractSnapshotDataFiles(
+    IcebergCatalog.FileScanResult result = 
icebergCatalog.extractFileScanTasksWithSchema(
         getNamespace(),
         getTableName(),
         getIcebergFilter(),
         getSnapshotTime(),
         getResidualFilterMode()
     );
-    if (snapshotDataFiles.isEmpty()) {
-      delegateInputSource = new EmptyInputSource();
+
+    List<FileScanTask> tasks = result.getTasks();
+    boolean anyDeletes = tasks.stream().anyMatch(t -> !t.deletes().isEmpty());
+
+    if (anyDeletes) {
+      hasDeleteFiles = true;
+      v2Tasks = tasks;
+      tableSchemaJson = result.getTableSchemaJson();
+      // delegateInputSource remains null; the V2 path handles reading 
directly.
     } else {
-      delegateInputSource = warehouseSource.create(snapshotDataFiles);
+      // V1 path: extract file paths and delegate to the warehouse input 
source.
+      List<String> paths = tasks.stream()
+                                .map(t -> t.file().path().toString())
+                                .collect(Collectors.toList());
+      delegateInputSource = paths.isEmpty() ? new EmptyInputSource() : 
warehouseSource.create(paths);
     }
     isLoaded = true;
   }
 
+  // ---- V2 split encoding / decoding ----
+
+  /**
+   * Encodes a {@link FileScanTask} as a V2 split.
+   *
+   * <pre>
+   * parts[0] = "v2"
+   * parts[1] = data file path
+   * parts[2] = file format name ("PARQUET" | "ORC")
+   * parts[3] = data file size in bytes
+   * parts[4] = data file record count
+   * parts[5] = table schema as JSON (from SchemaParser.toJson)
+   * parts[6+] = delete file entries:
+   *   position delete  → "POS:<size>:<count>:<path>"
+   *   equality delete  → "EQ:<fieldIds>:<size>:<count>:<path>"
+   * </pre>
+   *
+   * The path is always the last component so that paths containing colons 
(e.g. {@code s3://…})
+   * are captured correctly by {@code split(":", N)} with a limit.
+   */
+  private InputSplit<List<String>> taskToV2Split(FileScanTask task)
+  {
+    List<String> parts = new ArrayList<>();
+    parts.add(V2_MARKER);
+    parts.add(task.file().path().toString());
+    parts.add(task.file().format().name());
+    parts.add(String.valueOf(task.file().fileSizeInBytes()));
+    parts.add(String.valueOf(task.file().recordCount()));
+    parts.add(tableSchemaJson);
+
+    for (DeleteFile deleteFile : task.deletes()) {
+      if (deleteFile.content() == FileContent.POSITION_DELETES) {
+        parts.add("POS:" + deleteFile.fileSizeInBytes()
+                  + ":" + deleteFile.recordCount()
+                  + ":" + deleteFile.path());

Review Comment:
   ## CodeQL / Deprecated method or constructor invocation
   
   Invoking [ContentFile.path](1) should be avoided because it has been 
deprecated.
   
   [Show more 
details](https://github.com/apache/druid/security/code-scanning/11992)



##########
extensions-contrib/druid-iceberg-extensions/src/main/java/org/apache/druid/iceberg/input/IcebergInputSource.java:
##########
@@ -196,21 +272,172 @@
 
   protected void retrieveIcebergDatafiles()
   {
-    List<String> snapshotDataFiles = icebergCatalog.extractSnapshotDataFiles(
+    IcebergCatalog.FileScanResult result = 
icebergCatalog.extractFileScanTasksWithSchema(
         getNamespace(),
         getTableName(),
         getIcebergFilter(),
         getSnapshotTime(),
         getResidualFilterMode()
     );
-    if (snapshotDataFiles.isEmpty()) {
-      delegateInputSource = new EmptyInputSource();
+
+    List<FileScanTask> tasks = result.getTasks();
+    boolean anyDeletes = tasks.stream().anyMatch(t -> !t.deletes().isEmpty());
+
+    if (anyDeletes) {
+      hasDeleteFiles = true;
+      v2Tasks = tasks;
+      tableSchemaJson = result.getTableSchemaJson();
+      // delegateInputSource remains null; the V2 path handles reading 
directly.
     } else {
-      delegateInputSource = warehouseSource.create(snapshotDataFiles);
+      // V1 path: extract file paths and delegate to the warehouse input 
source.
+      List<String> paths = tasks.stream()
+                                .map(t -> t.file().path().toString())
+                                .collect(Collectors.toList());
+      delegateInputSource = paths.isEmpty() ? new EmptyInputSource() : 
warehouseSource.create(paths);
     }
     isLoaded = true;
   }
 
+  // ---- V2 split encoding / decoding ----
+
+  /**
+   * Encodes a {@link FileScanTask} as a V2 split.
+   *
+   * <pre>
+   * parts[0] = "v2"
+   * parts[1] = data file path
+   * parts[2] = file format name ("PARQUET" | "ORC")
+   * parts[3] = data file size in bytes
+   * parts[4] = data file record count
+   * parts[5] = table schema as JSON (from SchemaParser.toJson)
+   * parts[6+] = delete file entries:
+   *   position delete  → "POS:<size>:<count>:<path>"
+   *   equality delete  → "EQ:<fieldIds>:<size>:<count>:<path>"
+   * </pre>
+   *
+   * The path is always the last component so that paths containing colons 
(e.g. {@code s3://…})
+   * are captured correctly by {@code split(":", N)} with a limit.
+   */
+  private InputSplit<List<String>> taskToV2Split(FileScanTask task)
+  {
+    List<String> parts = new ArrayList<>();
+    parts.add(V2_MARKER);
+    parts.add(task.file().path().toString());

Review Comment:
   ## CodeQL / Deprecated method or constructor invocation
   
   Invoking [ContentFile.path](1) should be avoided because it has been 
deprecated.
   
   [Show more 
details](https://github.com/apache/druid/security/code-scanning/11991)



##########
extensions-contrib/druid-iceberg-extensions/src/main/java/org/apache/druid/iceberg/input/IcebergCatalog.java:
##########
@@ -151,6 +179,61 @@
     finally {
       Thread.currentThread().setContextClassLoader(currCtxClassloader);
     }
-    return dataFilePaths;
+    return new FileScanResult(tasks, tableSchemaJson);
+  }
+
+  /**
+   * Extract the iceberg FileScanTasks (data + delete file info) upto the 
latest snapshot.
+   * This is a convenience wrapper around {@link 
#extractFileScanTasksWithSchema} that drops
+   * the schema JSON.
+   *
+   * @param tableNamespace     The catalog namespace
+   * @param tableName          The iceberg table name
+   * @param icebergFilter      Filter to apply before reading files
+   * @param snapshotTime       Datetime for snapshot-as-of
+   * @param residualFilterMode Controls how residual filters are handled
+   * @return list of FileScanTask objects (each carries data file + associated 
delete files)
+   */
+  public List<FileScanTask> extractFileScanTasks(
+      String tableNamespace,
+      String tableName,
+      IcebergFilter icebergFilter,
+      DateTime snapshotTime,
+      ResidualFilterMode residualFilterMode
+  )
+  {
+    return extractFileScanTasksWithSchema(
+        tableNamespace,
+        tableName,
+        icebergFilter,
+        snapshotTime,
+        residualFilterMode
+    ).getTasks();
+  }
+
+  /**
+   * Extract the iceberg data files upto the latest snapshot associated with 
the table.
+   * This is a backward-compatible wrapper around {@link 
#extractFileScanTasks}.
+   *
+   * @param tableNamespace     The catalog namespace under which the table is 
defined
+   * @param tableName          The iceberg table name
+   * @param icebergFilter      The iceberg filter that needs to be applied 
before reading the files
+   * @param snapshotTime       Datetime that will be used to fetch the most 
recent snapshot as of this time
+   * @param residualFilterMode Controls how residual filters are handled. When 
filtering on non-partition
+   *                           columns, residual rows may be returned that 
need row-level filtering.
+   * @return a list of data file paths
+   */
+  public List<String> extractSnapshotDataFiles(
+      String tableNamespace,
+      String tableName,
+      IcebergFilter icebergFilter,
+      DateTime snapshotTime,
+      ResidualFilterMode residualFilterMode
+  )
+  {
+    return extractFileScanTasks(tableNamespace, tableName, icebergFilter, 
snapshotTime, residualFilterMode)
+        .stream()
+        .map(t -> t.file().path().toString())

Review Comment:
   ## CodeQL / Deprecated method or constructor invocation
   
   Invoking [ContentFile.path](1) should be avoided because it has been 
deprecated.
   
   [Show more 
details](https://github.com/apache/druid/security/code-scanning/11989)



##########
extensions-contrib/druid-iceberg-extensions/src/main/java/org/apache/druid/iceberg/input/IcebergInputSource.java:
##########
@@ -196,21 +272,172 @@
 
   protected void retrieveIcebergDatafiles()
   {
-    List<String> snapshotDataFiles = icebergCatalog.extractSnapshotDataFiles(
+    IcebergCatalog.FileScanResult result = 
icebergCatalog.extractFileScanTasksWithSchema(
         getNamespace(),
         getTableName(),
         getIcebergFilter(),
         getSnapshotTime(),
         getResidualFilterMode()
     );
-    if (snapshotDataFiles.isEmpty()) {
-      delegateInputSource = new EmptyInputSource();
+
+    List<FileScanTask> tasks = result.getTasks();
+    boolean anyDeletes = tasks.stream().anyMatch(t -> !t.deletes().isEmpty());
+
+    if (anyDeletes) {
+      hasDeleteFiles = true;
+      v2Tasks = tasks;
+      tableSchemaJson = result.getTableSchemaJson();
+      // delegateInputSource remains null; the V2 path handles reading 
directly.
     } else {
-      delegateInputSource = warehouseSource.create(snapshotDataFiles);
+      // V1 path: extract file paths and delegate to the warehouse input 
source.
+      List<String> paths = tasks.stream()
+                                .map(t -> t.file().path().toString())
+                                .collect(Collectors.toList());
+      delegateInputSource = paths.isEmpty() ? new EmptyInputSource() : 
warehouseSource.create(paths);
     }
     isLoaded = true;
   }
 
+  // ---- V2 split encoding / decoding ----
+
+  /**
+   * Encodes a {@link FileScanTask} as a V2 split.
+   *
+   * <pre>
+   * parts[0] = "v2"
+   * parts[1] = data file path
+   * parts[2] = file format name ("PARQUET" | "ORC")
+   * parts[3] = data file size in bytes
+   * parts[4] = data file record count
+   * parts[5] = table schema as JSON (from SchemaParser.toJson)
+   * parts[6+] = delete file entries:
+   *   position delete  → "POS:<size>:<count>:<path>"
+   *   equality delete  → "EQ:<fieldIds>:<size>:<count>:<path>"
+   * </pre>
+   *
+   * The path is always the last component so that paths containing colons 
(e.g. {@code s3://…})
+   * are captured correctly by {@code split(":", N)} with a limit.
+   */
+  private InputSplit<List<String>> taskToV2Split(FileScanTask task)
+  {
+    List<String> parts = new ArrayList<>();
+    parts.add(V2_MARKER);
+    parts.add(task.file().path().toString());
+    parts.add(task.file().format().name());
+    parts.add(String.valueOf(task.file().fileSizeInBytes()));
+    parts.add(String.valueOf(task.file().recordCount()));
+    parts.add(tableSchemaJson);
+
+    for (DeleteFile deleteFile : task.deletes()) {
+      if (deleteFile.content() == FileContent.POSITION_DELETES) {
+        parts.add("POS:" + deleteFile.fileSizeInBytes()
+                  + ":" + deleteFile.recordCount()
+                  + ":" + deleteFile.path());
+      } else {
+        // EQUALITY_DELETES
+        String fieldIds = deleteFile.equalityFieldIds() == null
+                          ? ""
+                          : deleteFile.equalityFieldIds().stream()
+                                      .map(String::valueOf)
+                                      .collect(Collectors.joining(","));
+        parts.add("EQ:" + fieldIds
+                  + ":" + deleteFile.fileSizeInBytes()
+                  + ":" + deleteFile.recordCount()
+                  + ":" + deleteFile.path());
+      }
+    }
+    return new InputSplit<>(parts);
+  }
+
+  /**
+   * Parses a V2-encoded split back into an {@link IcebergFileTaskInputSource}.
+   */
+  private IcebergFileTaskInputSource decodeV2Split(List<String> parts)
+  {
+    // parts[0] = "v2", parts[1] = dataFilePath, parts[2] = format,
+    // parts[3] = dataFileSize, parts[4] = dataFileRecordCount,
+    // parts[5] = tableSchemaJson, parts[6+] = delete file tokens
+    String dataFilePath = parts.get(1);
+    String fileFormat = parts.get(2);
+    long dataFileSize = Long.parseLong(parts.get(3));
+    long dataFileCount = Long.parseLong(parts.get(4));
+    String schemaJson = parts.get(5);
+
+    List<IcebergFileTaskInputSource.DeleteFileInfo> deleteFiles = new 
ArrayList<>();
+    for (int i = 6; i < parts.size(); i++) {
+      String token = parts.get(i);
+      if (token.startsWith("POS:")) {
+        // "POS:<size>:<count>:<path>"  — split on at most 4 colons (path is 
last)
+        String[] toks = token.split(":", 4);
+        long size = Long.parseLong(toks[1]);
+        long count = Long.parseLong(toks[2]);
+        String path = toks[3];
+        deleteFiles.add(
+            new IcebergFileTaskInputSource.DeleteFileInfo(
+                path,
+                "POSITION_DELETES",
+                null,
+                size,
+                count));
+      } else if (token.startsWith("EQ:")) {
+        // "EQ:<fieldIds>:<size>:<count>:<path>"
+        String[] toks = token.split(":", 5);
+        String fieldIdsStr = toks[1];
+        long size = Long.parseLong(toks[2]);

Review Comment:
   ## CodeQL / Missing catch of NumberFormatException
   
   Potential uncaught 'java.lang.NumberFormatException'.
   
   [Show more 
details](https://github.com/apache/druid/security/code-scanning/11985)



##########
extensions-contrib/druid-iceberg-extensions/src/main/java/org/apache/druid/iceberg/input/IcebergFileTaskInputSource.java:
##########
@@ -0,0 +1,287 @@
+/*
+ * 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;
+  }
+
+  @Override
+  public boolean isSplittable()
+  {
+    return false;
+  }
+
+  @Override
+  public InputSourceReader reader(
+      InputRowSchema inputRowSchema,
+      @Nullable InputFormat inputFormat,
+      File temporaryDirectory
+  )
+  {
+    return new IcebergNativeRecordReader(
+        dataFilePath,
+        fileFormat,
+        dataFileSizeInBytes,
+        dataFileRecordCount,
+        deleteFiles,
+        tableSchemaJson,
+        tableNamespace,
+        tableName,
+        icebergCatalog,
+        warehouseSource,
+        inputRowSchema
+    );
+  }
+
+  public String getDataFilePath()
+  {
+    return dataFilePath;
+  }
+
+  public String getFileFormat()
+  {
+    return fileFormat;
+  }
+
+  public long getDataFileSizeInBytes()
+  {
+    return dataFileSizeInBytes;
+  }
+
+  public long getDataFileRecordCount()
+  {
+    return dataFileRecordCount;
+  }
+
+  public List<DeleteFileInfo> getDeleteFiles()
+  {
+    return deleteFiles;
+  }
+
+  public String getTableSchemaJson()
+  {
+    return tableSchemaJson;
+  }
+
+  public String getTableNamespace()
+  {
+    return tableNamespace;
+  }
+
+  public String getTableName()
+  {
+    return tableName;
+  }
+
+  /**
+   * Metadata describing one delete file associated with a data file.
+   */
+  public static class DeleteFileInfo
+  {
+    /** Fully-qualified path to the delete file (same format as stored in 
Iceberg metadata). */
+    @JsonProperty
+    private final String path;
+
+    /**
+     * Delete file content type as a string: {@code "POSITION_DELETES"} or
+     * {@code "EQUALITY_DELETES"} (mirrors {@link 
org.apache.iceberg.FileContent}).
+     */
+    @JsonProperty
+    private final String content;
+
+    /** Field IDs used for equality matching (non-null only for 
EQUALITY_DELETES). */
+    @JsonProperty
+    private final List<Integer> equalityFieldIds;
+
+    @JsonProperty
+    private final long fileSizeInBytes;
+
+    @JsonProperty
+    private final long recordCount;
+
+    @JsonCreator
+    public DeleteFileInfo(
+        @JsonProperty("path") String path,
+        @JsonProperty("content") String content,
+        @JsonProperty("equalityFieldIds") @Nullable List<Integer> 
equalityFieldIds,
+        @JsonProperty("fileSizeInBytes") long fileSizeInBytes,
+        @JsonProperty("recordCount") long recordCount
+    )
+    {
+      this.path = path;
+      this.content = content;
+      this.equalityFieldIds = equalityFieldIds != null ? equalityFieldIds : 
Collections.emptyList();
+      this.fileSizeInBytes = fileSizeInBytes;
+      this.recordCount = recordCount;
+    }
+
+    /** Create a {@link DeleteFileInfo} from an Iceberg {@link DeleteFile} 
object. */
+    public static DeleteFileInfo fromDeleteFile(DeleteFile deleteFile)
+    {
+      List<Integer> fieldIds = Collections.emptyList();
+      if (deleteFile.content() == 
org.apache.iceberg.FileContent.EQUALITY_DELETES
+          && deleteFile.equalityFieldIds() != null) {
+        fieldIds = deleteFile.equalityFieldIds();
+      }
+      return new DeleteFileInfo(
+          deleteFile.path().toString(),

Review Comment:
   ## CodeQL / Deprecated method or constructor invocation
   
   Invoking [ContentFile.path](1) should be avoided because it has been 
deprecated.
   
   [Show more 
details](https://github.com/apache/druid/security/code-scanning/11988)



##########
extensions-contrib/druid-iceberg-extensions/src/main/java/org/apache/druid/iceberg/input/IcebergInputSource.java:
##########
@@ -196,21 +272,172 @@
 
   protected void retrieveIcebergDatafiles()
   {
-    List<String> snapshotDataFiles = icebergCatalog.extractSnapshotDataFiles(
+    IcebergCatalog.FileScanResult result = 
icebergCatalog.extractFileScanTasksWithSchema(
         getNamespace(),
         getTableName(),
         getIcebergFilter(),
         getSnapshotTime(),
         getResidualFilterMode()
     );
-    if (snapshotDataFiles.isEmpty()) {
-      delegateInputSource = new EmptyInputSource();
+
+    List<FileScanTask> tasks = result.getTasks();
+    boolean anyDeletes = tasks.stream().anyMatch(t -> !t.deletes().isEmpty());
+
+    if (anyDeletes) {
+      hasDeleteFiles = true;
+      v2Tasks = tasks;
+      tableSchemaJson = result.getTableSchemaJson();
+      // delegateInputSource remains null; the V2 path handles reading 
directly.
     } else {
-      delegateInputSource = warehouseSource.create(snapshotDataFiles);
+      // V1 path: extract file paths and delegate to the warehouse input 
source.
+      List<String> paths = tasks.stream()
+                                .map(t -> t.file().path().toString())
+                                .collect(Collectors.toList());
+      delegateInputSource = paths.isEmpty() ? new EmptyInputSource() : 
warehouseSource.create(paths);
     }
     isLoaded = true;
   }
 
+  // ---- V2 split encoding / decoding ----
+
+  /**
+   * Encodes a {@link FileScanTask} as a V2 split.
+   *
+   * <pre>
+   * parts[0] = "v2"
+   * parts[1] = data file path
+   * parts[2] = file format name ("PARQUET" | "ORC")
+   * parts[3] = data file size in bytes
+   * parts[4] = data file record count
+   * parts[5] = table schema as JSON (from SchemaParser.toJson)
+   * parts[6+] = delete file entries:
+   *   position delete  → "POS:<size>:<count>:<path>"
+   *   equality delete  → "EQ:<fieldIds>:<size>:<count>:<path>"
+   * </pre>
+   *
+   * The path is always the last component so that paths containing colons 
(e.g. {@code s3://…})
+   * are captured correctly by {@code split(":", N)} with a limit.
+   */
+  private InputSplit<List<String>> taskToV2Split(FileScanTask task)
+  {
+    List<String> parts = new ArrayList<>();
+    parts.add(V2_MARKER);
+    parts.add(task.file().path().toString());
+    parts.add(task.file().format().name());
+    parts.add(String.valueOf(task.file().fileSizeInBytes()));
+    parts.add(String.valueOf(task.file().recordCount()));
+    parts.add(tableSchemaJson);
+
+    for (DeleteFile deleteFile : task.deletes()) {
+      if (deleteFile.content() == FileContent.POSITION_DELETES) {
+        parts.add("POS:" + deleteFile.fileSizeInBytes()
+                  + ":" + deleteFile.recordCount()
+                  + ":" + deleteFile.path());
+      } else {
+        // EQUALITY_DELETES
+        String fieldIds = deleteFile.equalityFieldIds() == null
+                          ? ""
+                          : deleteFile.equalityFieldIds().stream()
+                                      .map(String::valueOf)
+                                      .collect(Collectors.joining(","));
+        parts.add("EQ:" + fieldIds
+                  + ":" + deleteFile.fileSizeInBytes()
+                  + ":" + deleteFile.recordCount()
+                  + ":" + deleteFile.path());
+      }
+    }
+    return new InputSplit<>(parts);
+  }
+
+  /**
+   * Parses a V2-encoded split back into an {@link IcebergFileTaskInputSource}.
+   */
+  private IcebergFileTaskInputSource decodeV2Split(List<String> parts)
+  {
+    // parts[0] = "v2", parts[1] = dataFilePath, parts[2] = format,
+    // parts[3] = dataFileSize, parts[4] = dataFileRecordCount,
+    // parts[5] = tableSchemaJson, parts[6+] = delete file tokens
+    String dataFilePath = parts.get(1);
+    String fileFormat = parts.get(2);
+    long dataFileSize = Long.parseLong(parts.get(3));

Review Comment:
   ## CodeQL / Missing catch of NumberFormatException
   
   Potential uncaught 'java.lang.NumberFormatException'.
   
   [Show more 
details](https://github.com/apache/druid/security/code-scanning/11981)



##########
extensions-contrib/druid-iceberg-extensions/src/main/java/org/apache/druid/iceberg/input/IcebergInputSource.java:
##########
@@ -196,21 +272,172 @@
 
   protected void retrieveIcebergDatafiles()
   {
-    List<String> snapshotDataFiles = icebergCatalog.extractSnapshotDataFiles(
+    IcebergCatalog.FileScanResult result = 
icebergCatalog.extractFileScanTasksWithSchema(
         getNamespace(),
         getTableName(),
         getIcebergFilter(),
         getSnapshotTime(),
         getResidualFilterMode()
     );
-    if (snapshotDataFiles.isEmpty()) {
-      delegateInputSource = new EmptyInputSource();
+
+    List<FileScanTask> tasks = result.getTasks();
+    boolean anyDeletes = tasks.stream().anyMatch(t -> !t.deletes().isEmpty());
+
+    if (anyDeletes) {
+      hasDeleteFiles = true;
+      v2Tasks = tasks;
+      tableSchemaJson = result.getTableSchemaJson();
+      // delegateInputSource remains null; the V2 path handles reading 
directly.
     } else {
-      delegateInputSource = warehouseSource.create(snapshotDataFiles);
+      // V1 path: extract file paths and delegate to the warehouse input 
source.
+      List<String> paths = tasks.stream()
+                                .map(t -> t.file().path().toString())
+                                .collect(Collectors.toList());
+      delegateInputSource = paths.isEmpty() ? new EmptyInputSource() : 
warehouseSource.create(paths);
     }
     isLoaded = true;
   }
 
+  // ---- V2 split encoding / decoding ----
+
+  /**
+   * Encodes a {@link FileScanTask} as a V2 split.
+   *
+   * <pre>
+   * parts[0] = "v2"
+   * parts[1] = data file path
+   * parts[2] = file format name ("PARQUET" | "ORC")
+   * parts[3] = data file size in bytes
+   * parts[4] = data file record count
+   * parts[5] = table schema as JSON (from SchemaParser.toJson)
+   * parts[6+] = delete file entries:
+   *   position delete  → "POS:<size>:<count>:<path>"
+   *   equality delete  → "EQ:<fieldIds>:<size>:<count>:<path>"
+   * </pre>
+   *
+   * The path is always the last component so that paths containing colons 
(e.g. {@code s3://…})
+   * are captured correctly by {@code split(":", N)} with a limit.
+   */
+  private InputSplit<List<String>> taskToV2Split(FileScanTask task)
+  {
+    List<String> parts = new ArrayList<>();
+    parts.add(V2_MARKER);
+    parts.add(task.file().path().toString());
+    parts.add(task.file().format().name());
+    parts.add(String.valueOf(task.file().fileSizeInBytes()));
+    parts.add(String.valueOf(task.file().recordCount()));
+    parts.add(tableSchemaJson);
+
+    for (DeleteFile deleteFile : task.deletes()) {
+      if (deleteFile.content() == FileContent.POSITION_DELETES) {
+        parts.add("POS:" + deleteFile.fileSizeInBytes()
+                  + ":" + deleteFile.recordCount()
+                  + ":" + deleteFile.path());
+      } else {
+        // EQUALITY_DELETES
+        String fieldIds = deleteFile.equalityFieldIds() == null
+                          ? ""
+                          : deleteFile.equalityFieldIds().stream()
+                                      .map(String::valueOf)
+                                      .collect(Collectors.joining(","));
+        parts.add("EQ:" + fieldIds
+                  + ":" + deleteFile.fileSizeInBytes()
+                  + ":" + deleteFile.recordCount()
+                  + ":" + deleteFile.path());
+      }
+    }
+    return new InputSplit<>(parts);
+  }
+
+  /**
+   * Parses a V2-encoded split back into an {@link IcebergFileTaskInputSource}.
+   */
+  private IcebergFileTaskInputSource decodeV2Split(List<String> parts)
+  {
+    // parts[0] = "v2", parts[1] = dataFilePath, parts[2] = format,
+    // parts[3] = dataFileSize, parts[4] = dataFileRecordCount,
+    // parts[5] = tableSchemaJson, parts[6+] = delete file tokens
+    String dataFilePath = parts.get(1);
+    String fileFormat = parts.get(2);
+    long dataFileSize = Long.parseLong(parts.get(3));
+    long dataFileCount = Long.parseLong(parts.get(4));
+    String schemaJson = parts.get(5);
+
+    List<IcebergFileTaskInputSource.DeleteFileInfo> deleteFiles = new 
ArrayList<>();
+    for (int i = 6; i < parts.size(); i++) {
+      String token = parts.get(i);
+      if (token.startsWith("POS:")) {
+        // "POS:<size>:<count>:<path>"  — split on at most 4 colons (path is 
last)
+        String[] toks = token.split(":", 4);
+        long size = Long.parseLong(toks[1]);
+        long count = Long.parseLong(toks[2]);
+        String path = toks[3];
+        deleteFiles.add(
+            new IcebergFileTaskInputSource.DeleteFileInfo(
+                path,
+                "POSITION_DELETES",
+                null,
+                size,
+                count));
+      } else if (token.startsWith("EQ:")) {
+        // "EQ:<fieldIds>:<size>:<count>:<path>"
+        String[] toks = token.split(":", 5);
+        String fieldIdsStr = toks[1];
+        long size = Long.parseLong(toks[2]);
+        long count = Long.parseLong(toks[3]);
+        String path = toks[4];
+        List<Integer> fieldIds = fieldIdsStr.isEmpty()
+                                 ? Collections.emptyList()
+                                 : Arrays.stream(fieldIdsStr.split(","))
+                                         .map(Integer::parseInt)

Review Comment:
   ## CodeQL / Missing catch of NumberFormatException
   
   Potential uncaught 'java.lang.NumberFormatException'.
   
   [Show more 
details](https://github.com/apache/druid/security/code-scanning/11987)



##########
extensions-contrib/druid-iceberg-extensions/src/main/java/org/apache/druid/iceberg/input/IcebergInputSource.java:
##########
@@ -196,21 +272,172 @@
 
   protected void retrieveIcebergDatafiles()
   {
-    List<String> snapshotDataFiles = icebergCatalog.extractSnapshotDataFiles(
+    IcebergCatalog.FileScanResult result = 
icebergCatalog.extractFileScanTasksWithSchema(
         getNamespace(),
         getTableName(),
         getIcebergFilter(),
         getSnapshotTime(),
         getResidualFilterMode()
     );
-    if (snapshotDataFiles.isEmpty()) {
-      delegateInputSource = new EmptyInputSource();
+
+    List<FileScanTask> tasks = result.getTasks();
+    boolean anyDeletes = tasks.stream().anyMatch(t -> !t.deletes().isEmpty());
+
+    if (anyDeletes) {
+      hasDeleteFiles = true;
+      v2Tasks = tasks;
+      tableSchemaJson = result.getTableSchemaJson();
+      // delegateInputSource remains null; the V2 path handles reading 
directly.
     } else {
-      delegateInputSource = warehouseSource.create(snapshotDataFiles);
+      // V1 path: extract file paths and delegate to the warehouse input 
source.
+      List<String> paths = tasks.stream()
+                                .map(t -> t.file().path().toString())
+                                .collect(Collectors.toList());
+      delegateInputSource = paths.isEmpty() ? new EmptyInputSource() : 
warehouseSource.create(paths);
     }
     isLoaded = true;
   }
 
+  // ---- V2 split encoding / decoding ----
+
+  /**
+   * Encodes a {@link FileScanTask} as a V2 split.
+   *
+   * <pre>
+   * parts[0] = "v2"
+   * parts[1] = data file path
+   * parts[2] = file format name ("PARQUET" | "ORC")
+   * parts[3] = data file size in bytes
+   * parts[4] = data file record count
+   * parts[5] = table schema as JSON (from SchemaParser.toJson)
+   * parts[6+] = delete file entries:
+   *   position delete  → "POS:<size>:<count>:<path>"
+   *   equality delete  → "EQ:<fieldIds>:<size>:<count>:<path>"
+   * </pre>
+   *
+   * The path is always the last component so that paths containing colons 
(e.g. {@code s3://…})
+   * are captured correctly by {@code split(":", N)} with a limit.
+   */
+  private InputSplit<List<String>> taskToV2Split(FileScanTask task)
+  {
+    List<String> parts = new ArrayList<>();
+    parts.add(V2_MARKER);
+    parts.add(task.file().path().toString());
+    parts.add(task.file().format().name());
+    parts.add(String.valueOf(task.file().fileSizeInBytes()));
+    parts.add(String.valueOf(task.file().recordCount()));
+    parts.add(tableSchemaJson);
+
+    for (DeleteFile deleteFile : task.deletes()) {
+      if (deleteFile.content() == FileContent.POSITION_DELETES) {
+        parts.add("POS:" + deleteFile.fileSizeInBytes()
+                  + ":" + deleteFile.recordCount()
+                  + ":" + deleteFile.path());
+      } else {
+        // EQUALITY_DELETES
+        String fieldIds = deleteFile.equalityFieldIds() == null
+                          ? ""
+                          : deleteFile.equalityFieldIds().stream()
+                                      .map(String::valueOf)
+                                      .collect(Collectors.joining(","));
+        parts.add("EQ:" + fieldIds
+                  + ":" + deleteFile.fileSizeInBytes()
+                  + ":" + deleteFile.recordCount()
+                  + ":" + deleteFile.path());
+      }
+    }
+    return new InputSplit<>(parts);
+  }
+
+  /**
+   * Parses a V2-encoded split back into an {@link IcebergFileTaskInputSource}.
+   */
+  private IcebergFileTaskInputSource decodeV2Split(List<String> parts)
+  {
+    // parts[0] = "v2", parts[1] = dataFilePath, parts[2] = format,
+    // parts[3] = dataFileSize, parts[4] = dataFileRecordCount,
+    // parts[5] = tableSchemaJson, parts[6+] = delete file tokens
+    String dataFilePath = parts.get(1);
+    String fileFormat = parts.get(2);
+    long dataFileSize = Long.parseLong(parts.get(3));
+    long dataFileCount = Long.parseLong(parts.get(4));
+    String schemaJson = parts.get(5);
+
+    List<IcebergFileTaskInputSource.DeleteFileInfo> deleteFiles = new 
ArrayList<>();
+    for (int i = 6; i < parts.size(); i++) {
+      String token = parts.get(i);
+      if (token.startsWith("POS:")) {
+        // "POS:<size>:<count>:<path>"  — split on at most 4 colons (path is 
last)
+        String[] toks = token.split(":", 4);
+        long size = Long.parseLong(toks[1]);

Review Comment:
   ## CodeQL / Missing catch of NumberFormatException
   
   Potential uncaught 'java.lang.NumberFormatException'.
   
   [Show more 
details](https://github.com/apache/druid/security/code-scanning/11983)



##########
extensions-contrib/druid-iceberg-extensions/src/main/java/org/apache/druid/iceberg/input/IcebergInputSource.java:
##########
@@ -196,21 +272,172 @@
 
   protected void retrieveIcebergDatafiles()
   {
-    List<String> snapshotDataFiles = icebergCatalog.extractSnapshotDataFiles(
+    IcebergCatalog.FileScanResult result = 
icebergCatalog.extractFileScanTasksWithSchema(
         getNamespace(),
         getTableName(),
         getIcebergFilter(),
         getSnapshotTime(),
         getResidualFilterMode()
     );
-    if (snapshotDataFiles.isEmpty()) {
-      delegateInputSource = new EmptyInputSource();
+
+    List<FileScanTask> tasks = result.getTasks();
+    boolean anyDeletes = tasks.stream().anyMatch(t -> !t.deletes().isEmpty());
+
+    if (anyDeletes) {
+      hasDeleteFiles = true;
+      v2Tasks = tasks;
+      tableSchemaJson = result.getTableSchemaJson();
+      // delegateInputSource remains null; the V2 path handles reading 
directly.
     } else {
-      delegateInputSource = warehouseSource.create(snapshotDataFiles);
+      // V1 path: extract file paths and delegate to the warehouse input 
source.
+      List<String> paths = tasks.stream()
+                                .map(t -> t.file().path().toString())
+                                .collect(Collectors.toList());
+      delegateInputSource = paths.isEmpty() ? new EmptyInputSource() : 
warehouseSource.create(paths);
     }
     isLoaded = true;
   }
 
+  // ---- V2 split encoding / decoding ----
+
+  /**
+   * Encodes a {@link FileScanTask} as a V2 split.
+   *
+   * <pre>
+   * parts[0] = "v2"
+   * parts[1] = data file path
+   * parts[2] = file format name ("PARQUET" | "ORC")
+   * parts[3] = data file size in bytes
+   * parts[4] = data file record count
+   * parts[5] = table schema as JSON (from SchemaParser.toJson)
+   * parts[6+] = delete file entries:
+   *   position delete  → "POS:<size>:<count>:<path>"
+   *   equality delete  → "EQ:<fieldIds>:<size>:<count>:<path>"
+   * </pre>
+   *
+   * The path is always the last component so that paths containing colons 
(e.g. {@code s3://…})
+   * are captured correctly by {@code split(":", N)} with a limit.
+   */
+  private InputSplit<List<String>> taskToV2Split(FileScanTask task)
+  {
+    List<String> parts = new ArrayList<>();
+    parts.add(V2_MARKER);
+    parts.add(task.file().path().toString());
+    parts.add(task.file().format().name());
+    parts.add(String.valueOf(task.file().fileSizeInBytes()));
+    parts.add(String.valueOf(task.file().recordCount()));
+    parts.add(tableSchemaJson);
+
+    for (DeleteFile deleteFile : task.deletes()) {
+      if (deleteFile.content() == FileContent.POSITION_DELETES) {
+        parts.add("POS:" + deleteFile.fileSizeInBytes()
+                  + ":" + deleteFile.recordCount()
+                  + ":" + deleteFile.path());
+      } else {
+        // EQUALITY_DELETES
+        String fieldIds = deleteFile.equalityFieldIds() == null
+                          ? ""
+                          : deleteFile.equalityFieldIds().stream()
+                                      .map(String::valueOf)
+                                      .collect(Collectors.joining(","));
+        parts.add("EQ:" + fieldIds
+                  + ":" + deleteFile.fileSizeInBytes()
+                  + ":" + deleteFile.recordCount()
+                  + ":" + deleteFile.path());
+      }
+    }
+    return new InputSplit<>(parts);
+  }
+
+  /**
+   * Parses a V2-encoded split back into an {@link IcebergFileTaskInputSource}.
+   */
+  private IcebergFileTaskInputSource decodeV2Split(List<String> parts)
+  {
+    // parts[0] = "v2", parts[1] = dataFilePath, parts[2] = format,
+    // parts[3] = dataFileSize, parts[4] = dataFileRecordCount,
+    // parts[5] = tableSchemaJson, parts[6+] = delete file tokens
+    String dataFilePath = parts.get(1);
+    String fileFormat = parts.get(2);
+    long dataFileSize = Long.parseLong(parts.get(3));
+    long dataFileCount = Long.parseLong(parts.get(4));
+    String schemaJson = parts.get(5);
+
+    List<IcebergFileTaskInputSource.DeleteFileInfo> deleteFiles = new 
ArrayList<>();
+    for (int i = 6; i < parts.size(); i++) {
+      String token = parts.get(i);
+      if (token.startsWith("POS:")) {
+        // "POS:<size>:<count>:<path>"  — split on at most 4 colons (path is 
last)
+        String[] toks = token.split(":", 4);
+        long size = Long.parseLong(toks[1]);
+        long count = Long.parseLong(toks[2]);
+        String path = toks[3];
+        deleteFiles.add(
+            new IcebergFileTaskInputSource.DeleteFileInfo(
+                path,
+                "POSITION_DELETES",
+                null,
+                size,
+                count));
+      } else if (token.startsWith("EQ:")) {
+        // "EQ:<fieldIds>:<size>:<count>:<path>"
+        String[] toks = token.split(":", 5);
+        String fieldIdsStr = toks[1];
+        long size = Long.parseLong(toks[2]);
+        long count = Long.parseLong(toks[3]);

Review Comment:
   ## CodeQL / Missing catch of NumberFormatException
   
   Potential uncaught 'java.lang.NumberFormatException'.
   
   [Show more 
details](https://github.com/apache/druid/security/code-scanning/11986)



##########
extensions-contrib/druid-iceberg-extensions/src/test/java/org/apache/druid/iceberg/input/IcebergInputSourceTest.java:
##########
@@ -311,13 +320,584 @@
         "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
+    );
+
+    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
+    );
+
+    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
+    );
+
+    // 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
+    );
+
+    // 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
+    InputSplit<List<String>> v1Split = new 
InputSplit<>(ImmutableList.of(dataFilePath));

Review Comment:
   ## CodeQL / Unread local variable
   
   Variable 'InputSplit<List<String>> v1Split' is never read.
   
   [Show more 
details](https://github.com/apache/druid/security/code-scanning/11979)



##########
extensions-contrib/druid-iceberg-extensions/src/main/java/org/apache/druid/iceberg/input/IcebergInputSource.java:
##########
@@ -196,21 +272,172 @@
 
   protected void retrieveIcebergDatafiles()
   {
-    List<String> snapshotDataFiles = icebergCatalog.extractSnapshotDataFiles(
+    IcebergCatalog.FileScanResult result = 
icebergCatalog.extractFileScanTasksWithSchema(
         getNamespace(),
         getTableName(),
         getIcebergFilter(),
         getSnapshotTime(),
         getResidualFilterMode()
     );
-    if (snapshotDataFiles.isEmpty()) {
-      delegateInputSource = new EmptyInputSource();
+
+    List<FileScanTask> tasks = result.getTasks();
+    boolean anyDeletes = tasks.stream().anyMatch(t -> !t.deletes().isEmpty());
+
+    if (anyDeletes) {
+      hasDeleteFiles = true;
+      v2Tasks = tasks;
+      tableSchemaJson = result.getTableSchemaJson();
+      // delegateInputSource remains null; the V2 path handles reading 
directly.
     } else {
-      delegateInputSource = warehouseSource.create(snapshotDataFiles);
+      // V1 path: extract file paths and delegate to the warehouse input 
source.
+      List<String> paths = tasks.stream()
+                                .map(t -> t.file().path().toString())

Review Comment:
   ## CodeQL / Deprecated method or constructor invocation
   
   Invoking [ContentFile.path](1) should be avoided because it has been 
deprecated.
   
   [Show more 
details](https://github.com/apache/druid/security/code-scanning/11990)



##########
extensions-contrib/druid-iceberg-extensions/src/main/java/org/apache/druid/iceberg/input/IcebergInputSource.java:
##########
@@ -196,21 +272,172 @@
 
   protected void retrieveIcebergDatafiles()
   {
-    List<String> snapshotDataFiles = icebergCatalog.extractSnapshotDataFiles(
+    IcebergCatalog.FileScanResult result = 
icebergCatalog.extractFileScanTasksWithSchema(
         getNamespace(),
         getTableName(),
         getIcebergFilter(),
         getSnapshotTime(),
         getResidualFilterMode()
     );
-    if (snapshotDataFiles.isEmpty()) {
-      delegateInputSource = new EmptyInputSource();
+
+    List<FileScanTask> tasks = result.getTasks();
+    boolean anyDeletes = tasks.stream().anyMatch(t -> !t.deletes().isEmpty());
+
+    if (anyDeletes) {
+      hasDeleteFiles = true;
+      v2Tasks = tasks;
+      tableSchemaJson = result.getTableSchemaJson();
+      // delegateInputSource remains null; the V2 path handles reading 
directly.
     } else {
-      delegateInputSource = warehouseSource.create(snapshotDataFiles);
+      // V1 path: extract file paths and delegate to the warehouse input 
source.
+      List<String> paths = tasks.stream()
+                                .map(t -> t.file().path().toString())
+                                .collect(Collectors.toList());
+      delegateInputSource = paths.isEmpty() ? new EmptyInputSource() : 
warehouseSource.create(paths);
     }
     isLoaded = true;
   }
 
+  // ---- V2 split encoding / decoding ----
+
+  /**
+   * Encodes a {@link FileScanTask} as a V2 split.
+   *
+   * <pre>
+   * parts[0] = "v2"
+   * parts[1] = data file path
+   * parts[2] = file format name ("PARQUET" | "ORC")
+   * parts[3] = data file size in bytes
+   * parts[4] = data file record count
+   * parts[5] = table schema as JSON (from SchemaParser.toJson)
+   * parts[6+] = delete file entries:
+   *   position delete  → "POS:<size>:<count>:<path>"
+   *   equality delete  → "EQ:<fieldIds>:<size>:<count>:<path>"
+   * </pre>
+   *
+   * The path is always the last component so that paths containing colons 
(e.g. {@code s3://…})
+   * are captured correctly by {@code split(":", N)} with a limit.
+   */
+  private InputSplit<List<String>> taskToV2Split(FileScanTask task)
+  {
+    List<String> parts = new ArrayList<>();
+    parts.add(V2_MARKER);
+    parts.add(task.file().path().toString());
+    parts.add(task.file().format().name());
+    parts.add(String.valueOf(task.file().fileSizeInBytes()));
+    parts.add(String.valueOf(task.file().recordCount()));
+    parts.add(tableSchemaJson);
+
+    for (DeleteFile deleteFile : task.deletes()) {
+      if (deleteFile.content() == FileContent.POSITION_DELETES) {
+        parts.add("POS:" + deleteFile.fileSizeInBytes()
+                  + ":" + deleteFile.recordCount()
+                  + ":" + deleteFile.path());
+      } else {
+        // EQUALITY_DELETES
+        String fieldIds = deleteFile.equalityFieldIds() == null
+                          ? ""
+                          : deleteFile.equalityFieldIds().stream()
+                                      .map(String::valueOf)
+                                      .collect(Collectors.joining(","));
+        parts.add("EQ:" + fieldIds
+                  + ":" + deleteFile.fileSizeInBytes()
+                  + ":" + deleteFile.recordCount()
+                  + ":" + deleteFile.path());

Review Comment:
   ## CodeQL / Deprecated method or constructor invocation
   
   Invoking [ContentFile.path](1) should be avoided because it has been 
deprecated.
   
   [Show more 
details](https://github.com/apache/druid/security/code-scanning/11993)



##########
extensions-contrib/druid-iceberg-extensions/src/main/java/org/apache/druid/iceberg/input/IcebergInputSource.java:
##########
@@ -196,21 +272,172 @@
 
   protected void retrieveIcebergDatafiles()
   {
-    List<String> snapshotDataFiles = icebergCatalog.extractSnapshotDataFiles(
+    IcebergCatalog.FileScanResult result = 
icebergCatalog.extractFileScanTasksWithSchema(
         getNamespace(),
         getTableName(),
         getIcebergFilter(),
         getSnapshotTime(),
         getResidualFilterMode()
     );
-    if (snapshotDataFiles.isEmpty()) {
-      delegateInputSource = new EmptyInputSource();
+
+    List<FileScanTask> tasks = result.getTasks();
+    boolean anyDeletes = tasks.stream().anyMatch(t -> !t.deletes().isEmpty());
+
+    if (anyDeletes) {
+      hasDeleteFiles = true;
+      v2Tasks = tasks;
+      tableSchemaJson = result.getTableSchemaJson();
+      // delegateInputSource remains null; the V2 path handles reading 
directly.
     } else {
-      delegateInputSource = warehouseSource.create(snapshotDataFiles);
+      // V1 path: extract file paths and delegate to the warehouse input 
source.
+      List<String> paths = tasks.stream()
+                                .map(t -> t.file().path().toString())
+                                .collect(Collectors.toList());
+      delegateInputSource = paths.isEmpty() ? new EmptyInputSource() : 
warehouseSource.create(paths);
     }
     isLoaded = true;
   }
 
+  // ---- V2 split encoding / decoding ----
+
+  /**
+   * Encodes a {@link FileScanTask} as a V2 split.
+   *
+   * <pre>
+   * parts[0] = "v2"
+   * parts[1] = data file path
+   * parts[2] = file format name ("PARQUET" | "ORC")
+   * parts[3] = data file size in bytes
+   * parts[4] = data file record count
+   * parts[5] = table schema as JSON (from SchemaParser.toJson)
+   * parts[6+] = delete file entries:
+   *   position delete  → "POS:<size>:<count>:<path>"
+   *   equality delete  → "EQ:<fieldIds>:<size>:<count>:<path>"
+   * </pre>
+   *
+   * The path is always the last component so that paths containing colons 
(e.g. {@code s3://…})
+   * are captured correctly by {@code split(":", N)} with a limit.
+   */
+  private InputSplit<List<String>> taskToV2Split(FileScanTask task)
+  {
+    List<String> parts = new ArrayList<>();
+    parts.add(V2_MARKER);
+    parts.add(task.file().path().toString());
+    parts.add(task.file().format().name());
+    parts.add(String.valueOf(task.file().fileSizeInBytes()));
+    parts.add(String.valueOf(task.file().recordCount()));
+    parts.add(tableSchemaJson);
+
+    for (DeleteFile deleteFile : task.deletes()) {
+      if (deleteFile.content() == FileContent.POSITION_DELETES) {
+        parts.add("POS:" + deleteFile.fileSizeInBytes()
+                  + ":" + deleteFile.recordCount()
+                  + ":" + deleteFile.path());
+      } else {
+        // EQUALITY_DELETES
+        String fieldIds = deleteFile.equalityFieldIds() == null
+                          ? ""
+                          : deleteFile.equalityFieldIds().stream()
+                                      .map(String::valueOf)
+                                      .collect(Collectors.joining(","));
+        parts.add("EQ:" + fieldIds
+                  + ":" + deleteFile.fileSizeInBytes()
+                  + ":" + deleteFile.recordCount()
+                  + ":" + deleteFile.path());
+      }
+    }
+    return new InputSplit<>(parts);
+  }
+
+  /**
+   * Parses a V2-encoded split back into an {@link IcebergFileTaskInputSource}.
+   */
+  private IcebergFileTaskInputSource decodeV2Split(List<String> parts)
+  {
+    // parts[0] = "v2", parts[1] = dataFilePath, parts[2] = format,
+    // parts[3] = dataFileSize, parts[4] = dataFileRecordCount,
+    // parts[5] = tableSchemaJson, parts[6+] = delete file tokens
+    String dataFilePath = parts.get(1);
+    String fileFormat = parts.get(2);
+    long dataFileSize = Long.parseLong(parts.get(3));
+    long dataFileCount = Long.parseLong(parts.get(4));

Review Comment:
   ## CodeQL / Missing catch of NumberFormatException
   
   Potential uncaught 'java.lang.NumberFormatException'.
   
   [Show more 
details](https://github.com/apache/druid/security/code-scanning/11982)



##########
extensions-contrib/druid-iceberg-extensions/src/main/java/org/apache/druid/iceberg/input/IcebergInputSource.java:
##########
@@ -196,21 +272,172 @@
 
   protected void retrieveIcebergDatafiles()
   {
-    List<String> snapshotDataFiles = icebergCatalog.extractSnapshotDataFiles(
+    IcebergCatalog.FileScanResult result = 
icebergCatalog.extractFileScanTasksWithSchema(
         getNamespace(),
         getTableName(),
         getIcebergFilter(),
         getSnapshotTime(),
         getResidualFilterMode()
     );
-    if (snapshotDataFiles.isEmpty()) {
-      delegateInputSource = new EmptyInputSource();
+
+    List<FileScanTask> tasks = result.getTasks();
+    boolean anyDeletes = tasks.stream().anyMatch(t -> !t.deletes().isEmpty());
+
+    if (anyDeletes) {
+      hasDeleteFiles = true;
+      v2Tasks = tasks;
+      tableSchemaJson = result.getTableSchemaJson();
+      // delegateInputSource remains null; the V2 path handles reading 
directly.
     } else {
-      delegateInputSource = warehouseSource.create(snapshotDataFiles);
+      // V1 path: extract file paths and delegate to the warehouse input 
source.
+      List<String> paths = tasks.stream()
+                                .map(t -> t.file().path().toString())
+                                .collect(Collectors.toList());
+      delegateInputSource = paths.isEmpty() ? new EmptyInputSource() : 
warehouseSource.create(paths);
     }
     isLoaded = true;
   }
 
+  // ---- V2 split encoding / decoding ----
+
+  /**
+   * Encodes a {@link FileScanTask} as a V2 split.
+   *
+   * <pre>
+   * parts[0] = "v2"
+   * parts[1] = data file path
+   * parts[2] = file format name ("PARQUET" | "ORC")
+   * parts[3] = data file size in bytes
+   * parts[4] = data file record count
+   * parts[5] = table schema as JSON (from SchemaParser.toJson)
+   * parts[6+] = delete file entries:
+   *   position delete  → "POS:<size>:<count>:<path>"
+   *   equality delete  → "EQ:<fieldIds>:<size>:<count>:<path>"
+   * </pre>
+   *
+   * The path is always the last component so that paths containing colons 
(e.g. {@code s3://…})
+   * are captured correctly by {@code split(":", N)} with a limit.
+   */
+  private InputSplit<List<String>> taskToV2Split(FileScanTask task)
+  {
+    List<String> parts = new ArrayList<>();
+    parts.add(V2_MARKER);
+    parts.add(task.file().path().toString());
+    parts.add(task.file().format().name());
+    parts.add(String.valueOf(task.file().fileSizeInBytes()));
+    parts.add(String.valueOf(task.file().recordCount()));
+    parts.add(tableSchemaJson);
+
+    for (DeleteFile deleteFile : task.deletes()) {
+      if (deleteFile.content() == FileContent.POSITION_DELETES) {
+        parts.add("POS:" + deleteFile.fileSizeInBytes()
+                  + ":" + deleteFile.recordCount()
+                  + ":" + deleteFile.path());
+      } else {
+        // EQUALITY_DELETES
+        String fieldIds = deleteFile.equalityFieldIds() == null
+                          ? ""
+                          : deleteFile.equalityFieldIds().stream()
+                                      .map(String::valueOf)
+                                      .collect(Collectors.joining(","));
+        parts.add("EQ:" + fieldIds
+                  + ":" + deleteFile.fileSizeInBytes()
+                  + ":" + deleteFile.recordCount()
+                  + ":" + deleteFile.path());
+      }
+    }
+    return new InputSplit<>(parts);
+  }
+
+  /**
+   * Parses a V2-encoded split back into an {@link IcebergFileTaskInputSource}.
+   */
+  private IcebergFileTaskInputSource decodeV2Split(List<String> parts)
+  {
+    // parts[0] = "v2", parts[1] = dataFilePath, parts[2] = format,
+    // parts[3] = dataFileSize, parts[4] = dataFileRecordCount,
+    // parts[5] = tableSchemaJson, parts[6+] = delete file tokens
+    String dataFilePath = parts.get(1);
+    String fileFormat = parts.get(2);
+    long dataFileSize = Long.parseLong(parts.get(3));
+    long dataFileCount = Long.parseLong(parts.get(4));
+    String schemaJson = parts.get(5);
+
+    List<IcebergFileTaskInputSource.DeleteFileInfo> deleteFiles = new 
ArrayList<>();
+    for (int i = 6; i < parts.size(); i++) {
+      String token = parts.get(i);
+      if (token.startsWith("POS:")) {
+        // "POS:<size>:<count>:<path>"  — split on at most 4 colons (path is 
last)
+        String[] toks = token.split(":", 4);
+        long size = Long.parseLong(toks[1]);
+        long count = Long.parseLong(toks[2]);
+        String path = toks[3];
+        deleteFiles.add(
+            new IcebergFileTaskInputSource.DeleteFileInfo(
+                path,
+                "POSITION_DELETES",
+                null,
+                size,
+                count));
+      } else if (token.startsWith("EQ:")) {
+        // "EQ:<fieldIds>:<size>:<count>:<path>"
+        String[] toks = token.split(":", 5);
+        String fieldIdsStr = toks[1];
+        long size = Long.parseLong(toks[2]);
+        long count = Long.parseLong(toks[3]);
+        String path = toks[4];
+        List<Integer> fieldIds = fieldIdsStr.isEmpty()
+                                 ? Collections.emptyList()
+                                 : Arrays.stream(fieldIdsStr.split(","))
+                                         .map(Integer::parseInt)
+                                         .collect(Collectors.toList());
+        deleteFiles.add(
+            new IcebergFileTaskInputSource.DeleteFileInfo(
+            path,
+            "EQUALITY_DELETES",
+            fieldIds,
+            size,
+            count));
+      }
+    }
+
+    return new IcebergFileTaskInputSource(
+        dataFilePath,
+        fileFormat,
+        dataFileSize,
+        dataFileCount,
+        deleteFiles,
+        schemaJson,
+        namespace,
+        tableName,
+        icebergCatalog,
+        warehouseSource
+    );
+  }
+
+  /** Build an {@link IcebergNativeRecordReader} from a live {@link 
FileScanTask}. */
+  private IcebergNativeRecordReader taskToNativeReader(FileScanTask task, 
InputRowSchema inputRowSchema)
+  {
+    List<IcebergFileTaskInputSource.DeleteFileInfo> deleteInfos = 
task.deletes().stream()
+                                                                      
.map(IcebergFileTaskInputSource.DeleteFileInfo::fromDeleteFile)
+                                                                      
.collect(Collectors.toList());
+    return new IcebergNativeRecordReader(
+        task.file().path().toString(),

Review Comment:
   ## CodeQL / Deprecated method or constructor invocation
   
   Invoking [ContentFile.path](1) should be avoided because it has been 
deprecated.
   
   [Show more 
details](https://github.com/apache/druid/security/code-scanning/11994)



##########
extensions-contrib/druid-iceberg-extensions/src/main/java/org/apache/druid/iceberg/input/IcebergInputSource.java:
##########
@@ -196,21 +272,172 @@
 
   protected void retrieveIcebergDatafiles()
   {
-    List<String> snapshotDataFiles = icebergCatalog.extractSnapshotDataFiles(
+    IcebergCatalog.FileScanResult result = 
icebergCatalog.extractFileScanTasksWithSchema(
         getNamespace(),
         getTableName(),
         getIcebergFilter(),
         getSnapshotTime(),
         getResidualFilterMode()
     );
-    if (snapshotDataFiles.isEmpty()) {
-      delegateInputSource = new EmptyInputSource();
+
+    List<FileScanTask> tasks = result.getTasks();
+    boolean anyDeletes = tasks.stream().anyMatch(t -> !t.deletes().isEmpty());
+
+    if (anyDeletes) {
+      hasDeleteFiles = true;
+      v2Tasks = tasks;
+      tableSchemaJson = result.getTableSchemaJson();
+      // delegateInputSource remains null; the V2 path handles reading 
directly.
     } else {
-      delegateInputSource = warehouseSource.create(snapshotDataFiles);
+      // V1 path: extract file paths and delegate to the warehouse input 
source.
+      List<String> paths = tasks.stream()
+                                .map(t -> t.file().path().toString())
+                                .collect(Collectors.toList());
+      delegateInputSource = paths.isEmpty() ? new EmptyInputSource() : 
warehouseSource.create(paths);
     }
     isLoaded = true;
   }
 
+  // ---- V2 split encoding / decoding ----
+
+  /**
+   * Encodes a {@link FileScanTask} as a V2 split.
+   *
+   * <pre>
+   * parts[0] = "v2"
+   * parts[1] = data file path
+   * parts[2] = file format name ("PARQUET" | "ORC")
+   * parts[3] = data file size in bytes
+   * parts[4] = data file record count
+   * parts[5] = table schema as JSON (from SchemaParser.toJson)
+   * parts[6+] = delete file entries:
+   *   position delete  → "POS:<size>:<count>:<path>"
+   *   equality delete  → "EQ:<fieldIds>:<size>:<count>:<path>"
+   * </pre>
+   *
+   * The path is always the last component so that paths containing colons 
(e.g. {@code s3://…})
+   * are captured correctly by {@code split(":", N)} with a limit.
+   */
+  private InputSplit<List<String>> taskToV2Split(FileScanTask task)
+  {
+    List<String> parts = new ArrayList<>();
+    parts.add(V2_MARKER);
+    parts.add(task.file().path().toString());
+    parts.add(task.file().format().name());
+    parts.add(String.valueOf(task.file().fileSizeInBytes()));
+    parts.add(String.valueOf(task.file().recordCount()));
+    parts.add(tableSchemaJson);
+
+    for (DeleteFile deleteFile : task.deletes()) {
+      if (deleteFile.content() == FileContent.POSITION_DELETES) {
+        parts.add("POS:" + deleteFile.fileSizeInBytes()
+                  + ":" + deleteFile.recordCount()
+                  + ":" + deleteFile.path());
+      } else {
+        // EQUALITY_DELETES
+        String fieldIds = deleteFile.equalityFieldIds() == null
+                          ? ""
+                          : deleteFile.equalityFieldIds().stream()
+                                      .map(String::valueOf)
+                                      .collect(Collectors.joining(","));
+        parts.add("EQ:" + fieldIds
+                  + ":" + deleteFile.fileSizeInBytes()
+                  + ":" + deleteFile.recordCount()
+                  + ":" + deleteFile.path());
+      }
+    }
+    return new InputSplit<>(parts);
+  }
+
+  /**
+   * Parses a V2-encoded split back into an {@link IcebergFileTaskInputSource}.
+   */
+  private IcebergFileTaskInputSource decodeV2Split(List<String> parts)
+  {
+    // parts[0] = "v2", parts[1] = dataFilePath, parts[2] = format,
+    // parts[3] = dataFileSize, parts[4] = dataFileRecordCount,
+    // parts[5] = tableSchemaJson, parts[6+] = delete file tokens
+    String dataFilePath = parts.get(1);
+    String fileFormat = parts.get(2);
+    long dataFileSize = Long.parseLong(parts.get(3));
+    long dataFileCount = Long.parseLong(parts.get(4));
+    String schemaJson = parts.get(5);
+
+    List<IcebergFileTaskInputSource.DeleteFileInfo> deleteFiles = new 
ArrayList<>();
+    for (int i = 6; i < parts.size(); i++) {
+      String token = parts.get(i);
+      if (token.startsWith("POS:")) {
+        // "POS:<size>:<count>:<path>"  — split on at most 4 colons (path is 
last)
+        String[] toks = token.split(":", 4);
+        long size = Long.parseLong(toks[1]);
+        long count = Long.parseLong(toks[2]);

Review Comment:
   ## CodeQL / Missing catch of NumberFormatException
   
   Potential uncaught 'java.lang.NumberFormatException'.
   
   [Show more 
details](https://github.com/apache/druid/security/code-scanning/11984)



##########
extensions-contrib/druid-iceberg-extensions/src/test/java/org/apache/druid/iceberg/input/IcebergInputSourceTest.java:
##########
@@ -395,6 +975,146 @@
     icebergTable.newAppend().appendFile(dataFile).commit();
   }
 
+  /**
+   * Writes a Parquet data file with the given rows and returns the resulting 
{@link DataFile}.
+   */
+  private DataFile writeParquetDataFile(
+      Table table,
+      Schema schema,
+      List<Map<String, Object>> rows,
+      String filePath
+  ) throws IOException
+  {
+    OutputFile outputFile = table.io().newOutputFile(filePath);
+    DataWriter<GenericRecord> dataWriter =
+        Parquet.writeData(outputFile)
+               .schema(schema)
+               .createWriterFunc(GenericParquetWriter::create)
+               .overwrite()
+               .withSpec(PartitionSpec.unpartitioned())
+               .build();
+    try {
+      for (Map<String, Object> row : rows) {
+        GenericRecord record = GenericRecord.create(schema);
+        for (Types.NestedField field : schema.columns()) {
+          record.setField(field.name(), row.get(field.name()));
+        }
+        dataWriter.write(record);
+      }
+    }
+    finally {
+      dataWriter.close();
+    }
+    return dataWriter.toDataFile();
+  }
+
+  /**
+   * Writes a Parquet position-delete file that marks {@code position} in 
{@code dataFilePath} as
+   * deleted.  Returns the resulting {@link DeleteFile} metadata object.
+   */
+  private DeleteFile writePositionDeleteFile(
+      Table table,
+      String dataFilePath,
+      long position,
+      String deleteFilePath
+  ) throws IOException
+  {
+    return writePositionDeleteFile(table, dataFilePath, 
ImmutableList.of(position), deleteFilePath);
+  }
+
+  private DeleteFile writePositionDeleteFile(
+      Table table,
+      String dataFilePath,
+      List<Long> positions,
+      String deleteFilePath
+  ) throws IOException
+  {
+    OutputFile outputFile = table.io().newOutputFile(deleteFilePath);
+    PositionDeleteWriter<GenericRecord> writer =
+        Parquet.writeDeletes(outputFile)
+               .createWriterFunc(GenericParquetWriter::create)
+               .overwrite()
+               .withSpec(PartitionSpec.unpartitioned())
+               .buildPositionWriter();
+    try {
+      PositionDelete<GenericRecord> posDelete = PositionDelete.create();
+      for (long pos : positions) {
+        posDelete.set(dataFilePath, pos, null);

Review Comment:
   ## CodeQL / Deprecated method or constructor invocation
   
   Invoking [PositionDelete.set](1) should be avoided because it has been 
deprecated.
   
   [Show more 
details](https://github.com/apache/druid/security/code-scanning/11995)



##########
extensions-contrib/druid-iceberg-extensions/src/test/java/org/apache/druid/iceberg/input/IcebergInputSourceTest.java:
##########
@@ -311,13 +320,584 @@
         "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
+    );
+
+    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
+    );
+
+    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
+    );
+
+    // 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
+    );
+
+    // 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
+    InputSplit<List<String>> v1Split = new 
InputSplit<>(ImmutableList.of(dataFilePath));
+    // Reinitialise so delegateInputSource is available
+    IcebergInputSource v1Source = new IcebergInputSource(
+        v2TableName,
+        NAMESPACE,
+        null,
+        testCatalog,
+        new LocalInputSourceFactory(),
+        null,
+        null
+    );

Review Comment:
   ## CodeQL / Unread local variable
   
   Variable 'IcebergInputSource v1Source' is never read.
   
   [Show more 
details](https://github.com/apache/druid/security/code-scanning/11980)



-- 
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