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


##########
extensions-contrib/druid-iceberg-extensions/src/main/java/org/apache/druid/iceberg/input/IcebergInputSource.java:
##########
@@ -292,51 +349,242 @@
         @Nullable SplitHintSpec splitHintSpec
     ) throws IOException
     {
+      if (!isLoaded) {
+        retrieveIcebergDatafiles();
+      }
+      if (hasDeleteFiles) {
+        return v2Tasks.stream().map(this::taskToV2Split);
+      }
       return warehouseInputSource().createSplits(inputFormat, splitHintSpec);
     }
 
     @Override
-    public int estimateNumSplits(InputFormat inputFormat, @Nullable 
SplitHintSpec splitHintSpec) throws IOException
+    public int estimateNumSplits(InputFormat inputFormat, @Nullable 
SplitHintSpec splitHintSpec)
+        throws IOException
     {
+      if (!isLoaded) {
+        retrieveIcebergDatafiles();
+      }
+      if (hasDeleteFiles) {
+        return v2Tasks.size();
+      }
       return warehouseInputSource().estimateNumSplits(inputFormat, 
splitHintSpec);
     }
 
+    /**
+     * Returns an {@link InputSource} for exactly the data described by {@code 
inputSplit}.
+     *
+     * <ul>
+     *   <li>V2 split (first element is {@value #V2_MARKER}): decoded into an
+     *       {@link IcebergFileTaskInputSource} that uses the native Iceberg 
reader.</li>
+     *   <li>V1 split (no marker): delegated to the warehouse {@link 
SplittableInputSource} as
+     *       before—fully backward-compatible.</li>
+     * </ul>
+     */
     @Override
     public InputSource withSplit(InputSplit<List<String>> inputSplit)
     {
+      List<String> parts = inputSplit.get();
+      if (!parts.isEmpty() && V2_MARKER.equals(parts.get(0))) {
+        return decodeV2Split(parts);
+      }
       return warehouseInputSource().withSplit(inputSplit);
     }
 
     @Override
     public SplitHintSpec getSplitHintSpecOrDefault(@Nullable SplitHintSpec 
splitHintSpec)
     {
-      return warehouseInputSource().getSplitHintSpecOrDefault(splitHintSpec);
+      if (delegateInputSource != null) {
+        return warehouseInputSource().getSplitHintSpecOrDefault(splitHintSpec);
+      }
+      // V2 path or not yet loaded — use interface default
+      return splitHintSpec == null ? 
SplittableInputSource.DEFAULT_SPLIT_HINT_SPEC : splitHintSpec;
     }
 
     private SplittableInputSource warehouseInputSource()
     {
-      if (!isLoaded) {
-        final List<String> snapshotDataFiles = 
icebergCatalog.extractSnapshotDataFiles(
-            getNamespace(),
-            getTableName(),
-            getIcebergFilter(),
-            getSnapshotTime(),
-            getResidualFilterMode()
-        );
-        if (snapshotDataFiles.isEmpty()) {
-          delegateInputSource = new EmptyInputSource();
+      return delegateInputSource;
+    }
+
+    private void retrieveIcebergDatafiles()
+    {
+      IcebergCatalog.FileScanResult result = 
icebergCatalog.extractFileScanTasksWithSchema(
+          getNamespace(),
+          getTableName(),
+          getIcebergFilter(),
+          getSnapshotTime(),
+          getResidualFilterMode()
+      );
+
+      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 {
+        // V1 path: extract file paths and delegate to the warehouse input 
source.
+        List<String> paths = tasks.stream()
+                                  .map(t -> t.file().location())
+                                  .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().location());
+      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.location());
         } else {
-          delegateInputSource = warehouseSource.create(snapshotDataFiles);
+          // 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.location());
         }
-        isLoaded = true;
       }
-      return delegateInputSource;
+      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 dataFileCount;
+      String schemaJson = parts.get(5);
+
+      List<IcebergFileTaskInputSource.DeleteFileInfo> deleteFiles = new 
ArrayList<>();
+      try {
+        dataFileSize = Long.parseLong(parts.get(3));
+        dataFileCount = Long.parseLong(parts.get(4));
+
+        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/12000)



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