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]