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]
