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]