github-advanced-security[bot] commented on code in PR #19534:
URL: https://github.com/apache/druid/pull/19534#discussion_r4058896350
##########
extensions-contrib/druid-iceberg-extensions/src/main/java/org/apache/druid/iceberg/input/IcebergInputSource.java:
##########
@@ -196,21 +273,180 @@
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().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 {
+ // 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());
+ }
+ }
+ 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/11996)
--
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]