FrankChen021 commented on code in PR #19534:
URL: https://github.com/apache/druid/pull/19534#discussion_r4082710267


##########
extensions-contrib/druid-iceberg-extensions/src/main/java/org/apache/druid/iceberg/input/IcebergInputSource.java:
##########
@@ -283,6 +309,37 @@ public InputSourceReader reader(
         File temporaryDirectory
     )
     {
+      if (!isLoaded) {
+        retrieveIcebergDatafiles();
+      }
+      if (hasDeleteFiles) {
+        // V2 path: create a combined reader across all file-scan tasks.
+        final List<IcebergNativeRecordReader> taskReaders = v2Tasks.stream()

Review Comment:
   [P2] V2 reads ignore the configured input format
   
   **Finding:** When any planned task has delete files, this branch routes 
every task through IcebergNativeRecordReader, but the reader is constructed 
only from InputRowSchema and the supplied InputFormat is discarded. 
ParquetInputFormat options such as flattenSpec and binaryAsString are therefore 
skipped for V2 tables, including data files without deletes, so a valid 
ingestion spec can produce different or missing indexed values from the 
existing warehouse-reader path.
   
   **Suggestion:** Preserve the configured Parquet input-format semantics on 
the delete-aware path, either by carrying the relevant options into the native 
reader or by applying delete filtering around the existing Druid input-format 
reader.



##########
extensions-contrib/druid-iceberg-extensions/src/main/java/org/apache/druid/iceberg/input/IcebergRecordConverter.java:
##########
@@ -0,0 +1,163 @@
+/*
+ * 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 org.apache.druid.data.input.InputRow;
+import org.apache.druid.data.input.InputRowSchema;
+import org.apache.druid.data.input.MapBasedInputRow;
+import org.apache.iceberg.Schema;
+import org.apache.iceberg.data.Record;
+import org.apache.iceberg.types.Type;
+import org.apache.iceberg.types.Types;
+import org.joda.time.DateTime;
+
+import java.nio.ByteBuffer;
+import java.time.LocalDate;
+import java.time.LocalDateTime;
+import java.time.LocalTime;
+import java.time.OffsetDateTime;
+import java.time.ZoneOffset;
+import java.util.ArrayList;
+import java.util.LinkedHashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.stream.Collectors;
+
+/**
+ * Static utility class for converting Iceberg {@link Record} objects to Druid 
{@link InputRow}.
+ */
+public class IcebergRecordConverter
+{
+  private IcebergRecordConverter()
+  {
+  }
+
+  /**
+   * Convert an Iceberg {@link Record} to a {@link Map} of field name → Java 
value.
+   * Iceberg primitive types are passed through; temporal and binary types are 
converted
+   * to representations suitable for Druid ingestion.
+   */
+  public static Map<String, Object> convertToMap(Record record, Schema schema)
+  {
+    Map<String, Object> result = new LinkedHashMap<>();
+    for (Types.NestedField field : schema.columns()) {
+      String name = field.name();
+      Object value = record.getField(name);
+      result.put(name, convertValue(value, field.type()));
+    }
+    return result;
+  }
+
+  /**
+   * Convert a raw-value map plus an {@link InputRowSchema} into a {@link 
MapBasedInputRow}.
+   * If the dimensions list is empty (auto mode), all columns except the 
timestamp column are used.
+   */
+  public static InputRow toInputRow(Map<String, Object> row, InputRowSchema 
inputRowSchema)
+  {
+    DateTime timestamp = 
inputRowSchema.getTimestampSpec().extractTimestamp(row);
+
+    List<String> dimensions = 
inputRowSchema.getDimensionsSpec().getDimensionNames();

Review Comment:
   [P2] Native conversion ignores dimension discovery rules
   
   **Finding:** The native V2 conversion uses only 
DimensionsSpec.getDimensionNames(), and when that list is empty it treats every 
non-timestamp column as a dimension. It never applies 
DimensionsSpec.useSchemaDiscovery(), includeAllDimensions, or 
dimensionExclusions, so V2 ingestion can omit dimensions that 
discovery/include-all should add or index fields that the ingestion spec 
explicitly excludes.
   
   **Suggestion:** Resolve dimensions with the same 
MapInputRowParser.findDimensions semantics as the normal input-format path, 
including discovery/include-all behavior and exclusions, before constructing 
MapBasedInputRow.



##########
extensions-contrib/druid-iceberg-extensions/src/main/java/org/apache/druid/iceberg/input/IcebergCatalog.java:
##########
@@ -114,42 +147,63 @@ public List<String> extractSnapshotDataFiles(
   )
   {
     Catalog catalog = retrieveCatalog();
-    Namespace namespace = Namespace.of(tableNamespace);
     String tableIdentifier = tableNamespace + "." + tableName;
-
-    List<String> dataFilePaths = new ArrayList<>();
+    List<FileScanTask> tasks = new ArrayList<>();
+    String tableSchemaJson;
 
     ClassLoader currCtxClassloader = 
Thread.currentThread().getContextClassLoader();
     try {
       
Thread.currentThread().setContextClassLoader(getClass().getClassLoader());
-      TableIdentifier icebergTableIdentifier = 
catalog.listTables(namespace).stream()
-                                                      .filter(tableId -> 
tableId.toString().equals(tableIdentifier))
-                                                      .findFirst()
-                                                      .orElseThrow(() -> new 
IAE(
-                                                          " Couldn't retrieve 
table identifier for '%s'. Please verify that the table exists in the given 
catalog",
-                                                          tableIdentifier
-                                                      ));
+
+      TableIdentifier icebergTableIdentifier = TableIdentifier.of(
+          Namespace.of(tableNamespace), tableName
+      );
 
       long start = System.currentTimeMillis();
-      TableScan tableScan = 
catalog.loadTable(icebergTableIdentifier).newScan();
+      Table table = catalog.loadTable(icebergTableIdentifier);
+      tableSchemaJson = SchemaParser.toJson(table.schema());
 
+      TableScan tableScan = table.newScan();
       if (icebergFilter != null) {
         tableScan = icebergFilter.filter(tableScan);
       }
       if (snapshotTime != null) {
         tableScan = tableScan.asOfTime(snapshotTime.getMillis());
       }
-
       tableScan = tableScan.caseSensitive(isCaseSensitive());
-      enforceResidualMode(tableScan, residualFilterMode);
-      try (CloseableIterable<FileScanTask> tasks = tableScan.planFiles()) {
-        for (FileScanTask task : tasks) {
-          dataFilePaths.add(task.file().location());
+
+      CloseableIterable<FileScanTask> taskIterable = tableScan.planFiles();

Review Comment:
   [P3] Close the planned scan iterable
   
   **Finding:** extractFileScanTasksWithSchema obtains the CloseableIterable 
returned by TableScan.planFiles() and iterates it without closing it. Iceberg 
wraps scan planning in a close-completion callback that finalizes planning 
metrics and closes nested manifest iterables; the old path used 
try-with-resources. Normal scans therefore skip completion, and exception or 
early-failure paths can retain planning resources across repeated ingestion 
attempts.
   
   **Suggestion:** Iterate the planned tasks inside try-with-resources and 
close the CloseableIterable on both normal completion and failure.



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