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


##########
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 getDimensionNames(), and 
when that list is empty it treats every non-timestamp column as a dimension. It 
never honors DimensionsSpec.useSchemaDiscovery(), isIncludeAllDimensions(), or 
getDimensionExclusions(). A V2 ingestion can therefore omit columns that schema 
discovery/include-all should add, or index excluded columns as dimensions, 
producing different rows from the existing V1 Parquet reader.
   
   **Suggestion:** Build the dimension list with the same 
MapInputRowParser.findDimensions semantics, including schema 
discovery/include-all behavior and exclusions, before constructing 
MapBasedInputRow.



##########
extensions-contrib/druid-iceberg-extensions/src/main/java/org/apache/druid/iceberg/input/IcebergInputSource.java:
##########
@@ -116,6 +140,34 @@ public InputSourceReader reader(
     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 path drops configured input-format behavior
   
   **Finding:** When any planned task has delete files, this branch routes 
every task through IcebergNativeRecordReader, but the reader is constructed 
with only InputRowSchema and the InputFormat argument is discarded. That 
changes the documented Parquet ingestion behavior for V2 tables: 
ParquetInputFormat options such as flattenSpec and binaryAsString are no longer 
applied, and even tasks without deletes are affected because the whole source 
switches to this branch.
   
   **Suggestion:** Carry the configured Parquet input-format options into the 
native reader, or preserve the existing warehouse Parquet reader and layer the 
Iceberg delete filtering around its rows so the same InputFormat semantics are 
used on both paths.



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