This is an automated email from the ASF dual-hosted git repository.
abhishekrb19 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/druid.git
The following commit(s) were added to refs/heads/master by this push:
new f81c3c19c9d feat: Adds an opt-in vectorized read path to the Iceberg
input source using iceberg-arrow. (#19510)
f81c3c19c9d is described below
commit f81c3c19c9da7d19a7e1305c0e4ff4aaa87ca534
Author: Shekhar Rajak <[email protected]>
AuthorDate: Mon Sep 21 19:42:55 2026 +0530
feat: Adds an opt-in vectorized read path to the Iceberg input source using
iceberg-arrow. (#19510)
Adds an opt-in vectorized read path to the Iceberg input source, backed by
iceberg-arrow.
IcebergInputSource ("type": "iceberg") is now a facade over two internal
reader-mode delegates, selected by the useArrowReader property:
StandardDelegate — the existing behavior: resolve the snapshot to a list of
live data file paths and read them through warehouseSource with the configured
inputFormat.
ArrowDelegate — builds a filtered/projected TableScan and reads it with
Iceberg's ArrowReader, converting ColumnarBatch vectors into Druid InputRows
and tracking processed bytes from the Arrow buffers.
Default values: useArrowReader : false and arrowBatchSize: 1024
Apache Iceberg ingestion supports an opt-in vectorized reader backed by
iceberg-arrow. Enable it with "useArrowReader": true; tune batching with
arrowBatchSize (default 1024). The Arrow path supports Parquet data files,
dictionary-encoded values, schema evolution, time-travel snapshots, column
projection, and predicate pushdown. Snapshots containing positional or equality
delete files are currently unsupported and are rejected with a clear error; use
the standard delete-aware reader f [...]
---
docs/development/extensions-contrib/iceberg.md | 28 +
.../iceberg/IcebergRestCatalogIngestionTest.java | 38 +-
.../druid-iceberg-extensions/pom.xml | 23 +
.../input/IcebergArrowInputSourceReader.java | 476 ++++++++++++++++
.../apache/druid/iceberg/input/IcebergCatalog.java | 108 +++-
.../druid/iceberg/input/IcebergInputSource.java | 402 +++++++++++---
.../input/IcebergArrowInputSourceReaderTest.java | 614 +++++++++++++++++++++
.../input/IcebergInputSourceArrowModeTest.java | 427 ++++++++++++++
.../iceberg/input/IcebergInputSourceTest.java | 26 +-
licenses.yaml | 69 +++
pom.xml | 17 +
11 files changed, 2098 insertions(+), 130 deletions(-)
diff --git a/docs/development/extensions-contrib/iceberg.md
b/docs/development/extensions-contrib/iceberg.md
index 7373d0fadc8..d4f1a939ed3 100644
--- a/docs/development/extensions-contrib/iceberg.md
+++ b/docs/development/extensions-contrib/iceberg.md
@@ -194,6 +194,34 @@ Example:
When `residualFilterMode` is set to `fail` and a residual filter is detected,
the job will fail with an error message indicating which filter expression
produced the residual. This helps ensure data quality by preventing unintended
rows from being ingested.
+## Arrow vectorized reader
+
+By default the Iceberg input source resolves the snapshot to a list of data
file paths and reads them through the `warehouseSource`. Setting
`useArrowReader` to `true` reads the table scan directly with Iceberg's
vectorized Arrow reader instead, which avoids the per-file input format layer.
+
+| Property | Description | Default |
+|----------|-------------|---------|
+| `useArrowReader` | Read the table with Iceberg's vectorized Arrow reader. |
`false` |
+| `arrowBatchSize` | Rows per Arrow batch. Only applies when `useArrowReader`
is `true`. | `1024` |
+
+Example:
+```json
+{
+ "type": "iceberg",
+ "tableName": "events",
+ "namespace": "analytics",
+ "icebergCatalog": { ... },
+ "useArrowReader": true,
+ "arrowBatchSize": 2048
+}
+```
+
+Note the following when `useArrowReader` is `true`:
+
+- The table must store its data files as Parquet. Ingestion fails with an
error if it finds a data file in any other format.
+- `warehouseSource` is not required and is unused. It is still required when
`useArrowReader` is `false`.
+- No `inputFormat` is needed, since the reader works from Iceberg metadata.
+- The input source is not splittable, so the table is read by a single task
even in a parallel ingestion. Keep `useArrowReader` set to `false` if you rely
on parallel batch ingestion for throughput.
+
## Known limitations
This section lists the known limitations that apply to the Iceberg extension.
diff --git
a/embedded-tests/src/test/java/org/apache/druid/testing/embedded/iceberg/IcebergRestCatalogIngestionTest.java
b/embedded-tests/src/test/java/org/apache/druid/testing/embedded/iceberg/IcebergRestCatalogIngestionTest.java
index 502a516b9a2..e95a8185078 100644
---
a/embedded-tests/src/test/java/org/apache/druid/testing/embedded/iceberg/IcebergRestCatalogIngestionTest.java
+++
b/embedded-tests/src/test/java/org/apache/druid/testing/embedded/iceberg/IcebergRestCatalogIngestionTest.java
@@ -23,6 +23,7 @@ import com.google.common.collect.ImmutableList;
import com.google.common.collect.ImmutableMap;
import org.apache.druid.data.input.parquet.ParquetExtensionsModule;
import org.apache.druid.java.util.common.StringUtils;
+import org.apache.druid.java.util.common.logger.Logger;
import org.apache.druid.query.http.SqlTaskStatus;
import org.apache.druid.testing.embedded.EmbeddedBroker;
import org.apache.druid.testing.embedded.EmbeddedCoordinator;
@@ -46,6 +47,7 @@ import org.junit.jupiter.api.Test;
import java.io.IOException;
import java.util.UUID;
+import java.util.concurrent.TimeUnit;
/**
* Ingestion test for Iceberg tables via a REST catalog.
@@ -56,6 +58,7 @@ public class IcebergRestCatalogIngestionTest extends
EmbeddedClusterTestBase
{
private static final String ICEBERG_NAMESPACE = "default";
private static final String ICEBERG_TABLE_NAME = "test_events";
+ private static final Logger log = new
Logger(IcebergRestCatalogIngestionTest.class);
private static final Schema ICEBERG_SCHEMA = new Schema(
Types.NestedField.required(1, "event_time", Types.StringType.get()),
@@ -136,7 +139,23 @@ public class IcebergRestCatalogIngestionTest extends
EmbeddedClusterTestBase
@Test
public void testIngestFromIcebergRestCatalog()
{
+ ingestFromIcebergRestCatalog(dataSource, false);
+ }
+
+ @Test
+ public void testIngestFromIcebergRestCatalogWithArrowReader()
+ {
+ ingestFromIcebergRestCatalog(dataSource + "_arrow", true);
+ }
+
+ private void ingestFromIcebergRestCatalog(final String targetDataSource,
final boolean useArrowReader)
+ {
+ final long startNanos = System.nanoTime();
final String catalogUri = icebergCatalog.getCatalogUri();
+ final String warehouseSource = useArrowReader ? "" :
"\"warehouseSource\":{\"type\":\"local\"}";
+ final String arrowReaderOptions = useArrowReader
+ ?
"\"useArrowReader\":true,\"arrowBatchSize\":2"
+ : "";
final String sql = StringUtils.format(
"INSERT INTO %s\n"
@@ -151,7 +170,7 @@ public class IcebergRestCatalogIngestionTest extends
EmbeddedClusterTestBase
+ "\"namespace\":\"%s\","
+ "\"icebergCatalog\":{\"type\":\"rest\",\"catalogUri\":\"%s\","
+
"\"catalogProperties\":{\"io-impl\":\"org.apache.iceberg.hadoop.HadoopFileIO\"}},"
- + "\"warehouseSource\":{\"type\":\"local\"}}',\n"
+ + "%s%s}',\n"
+ " '{\"type\":\"parquet\"}',\n"
+ " '[{\"type\":\"string\",\"name\":\"event_time\"},"
+ "{\"type\":\"string\",\"name\":\"name\"},"
@@ -159,22 +178,31 @@ public class IcebergRestCatalogIngestionTest extends
EmbeddedClusterTestBase
+ " )\n"
+ ")\n"
+ "PARTITIONED BY ALL TIME",
- dataSource,
+ targetDataSource,
ICEBERG_TABLE_NAME,
ICEBERG_NAMESPACE,
- catalogUri
+ catalogUri,
+ warehouseSource,
+ arrowReaderOptions
);
final SqlTaskStatus taskStatus = msqApis.submitTaskSql(sql);
cluster.callApi().waitForTaskToSucceed(taskStatus.getTaskId(), overlord);
- cluster.callApi().waitForAllSegmentsToBeAvailable(dataSource, coordinator,
broker);
+ cluster.callApi().waitForAllSegmentsToBeAvailable(targetDataSource,
coordinator, broker);
cluster.callApi().verifySqlQuery(
"SELECT __time, \"name\", \"value\" FROM %s ORDER BY __time",
- dataSource,
+ targetDataSource,
"2024-01-01T00:00:00.000Z,alice,100\n"
+ "2024-01-01T01:00:00.000Z,bob,200\n"
+ "2024-01-01T02:00:00.000Z,charlie,300"
);
+
+ log.info(
+ "Iceberg REST catalog ingestion timing: useArrowReader[%s],
dataSource[%s], elapsedMs[%d]",
+ useArrowReader,
+ targetDataSource,
+ TimeUnit.NANOSECONDS.toMillis(System.nanoTime() - startNanos)
+ );
}
}
diff --git a/extensions-contrib/druid-iceberg-extensions/pom.xml
b/extensions-contrib/druid-iceberg-extensions/pom.xml
index 9e15b53631b..6b5286d9fd6 100644
--- a/extensions-contrib/druid-iceberg-extensions/pom.xml
+++ b/extensions-contrib/druid-iceberg-extensions/pom.xml
@@ -811,6 +811,29 @@
<scope>provided</scope>
</dependency>
+ <dependency>
+ <groupId>org.apache.iceberg</groupId>
+ <artifactId>iceberg-arrow</artifactId>
+ <version>${iceberg.core.version}</version>
+ <exclusions>
+ <exclusion>
+ <groupId>org.apache.arrow</groupId>
+ <artifactId>arrow-memory-netty</artifactId>
+ </exclusion>
+ </exclusions>
+ </dependency>
+ <dependency>
+ <groupId>org.apache.arrow</groupId>
+ <artifactId>arrow-vector</artifactId>
+ <version>${arrow.version}</version>
+ </dependency>
+ <dependency>
+ <groupId>org.apache.arrow</groupId>
+ <artifactId>arrow-memory-unsafe</artifactId>
+ <version>${arrow.version}</version>
+ <scope>runtime</scope>
+ </dependency>
+
<dependency>
<groupId>org.apache.iceberg</groupId>
<artifactId>iceberg-parquet</artifactId>
diff --git
a/extensions-contrib/druid-iceberg-extensions/src/main/java/org/apache/druid/iceberg/input/IcebergArrowInputSourceReader.java
b/extensions-contrib/druid-iceberg-extensions/src/main/java/org/apache/druid/iceberg/input/IcebergArrowInputSourceReader.java
new file mode 100644
index 00000000000..6266be5b29a
--- /dev/null
+++
b/extensions-contrib/druid-iceberg-extensions/src/main/java/org/apache/druid/iceberg/input/IcebergArrowInputSourceReader.java
@@ -0,0 +1,476 @@
+/*
+ * 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 com.google.common.collect.Maps;
+import org.apache.arrow.vector.FieldVector;
+import org.apache.druid.data.input.ColumnsFilter;
+import org.apache.druid.data.input.InputRow;
+import org.apache.druid.data.input.InputRowListPlusRawValues;
+import org.apache.druid.data.input.InputRowSchema;
+import org.apache.druid.data.input.InputSourceReader;
+import org.apache.druid.data.input.InputStats;
+import org.apache.druid.data.input.MapBasedInputRow;
+import org.apache.druid.data.input.impl.MapInputRowParser;
+import org.apache.druid.error.DruidException;
+import org.apache.druid.iceberg.filter.IcebergFilter;
+import org.apache.druid.java.util.common.parsers.CloseableIterator;
+import org.apache.iceberg.CombinedScanTask;
+import org.apache.iceberg.FileScanTask;
+import org.apache.iceberg.Schema;
+import org.apache.iceberg.Table;
+import org.apache.iceberg.TableScan;
+import org.apache.iceberg.arrow.vectorized.ArrowReader;
+import org.apache.iceberg.arrow.vectorized.ColumnVector;
+import org.apache.iceberg.arrow.vectorized.ColumnarBatch;
+import org.apache.iceberg.io.CloseableIterable;
+import org.apache.iceberg.types.Type;
+import org.apache.iceberg.types.Types;
+import org.apache.iceberg.util.TableScanUtil;
+import org.joda.time.DateTime;
+
+import javax.annotation.Nullable;
+import java.io.IOException;
+import java.util.List;
+import java.util.Map;
+import java.util.NoSuchElementException;
+import java.util.concurrent.TimeUnit;
+import java.util.stream.Collectors;
+
+/**
+ * Reads an Iceberg table via iceberg-arrow's {@link ArrowReader}, yielding
{@link InputRow} objects.
+ *
+ * Type coercion and compatible schema evolution are handled by the Iceberg
library. Druid only consumes
+ * the resulting {@link ColumnarBatch} batches and maps them to {@link
MapBasedInputRow}.
+ *
+ * Column projection and predicate push-down are applied at scan planning time
so only requested
+ * columns and matching files are read from storage.
+ *
+ * Note: iceberg-arrow currently supports Parquet data files only. ORC and
Avro files will throw
+ * {@link UnsupportedOperationException} at read time. Delete-file snapshots
are rejected because
+ * iceberg-arrow does not apply equality or positional deletes.
+ */
+public class IcebergArrowInputSourceReader implements InputSourceReader
+{
+ static final int DEFAULT_BATCH_SIZE = 1024;
+
+ private final Table table;
+ @Nullable
+ private final IcebergFilter icebergFilter;
+ @Nullable
+ private final DateTime snapshotTime;
+ private final boolean caseSensitive;
+ private final InputRowSchema schema;
+ private final int batchSize;
+ private final ResidualFilterMode residualFilterMode;
+
+ public IcebergArrowInputSourceReader(
+ final Table table,
+ @Nullable final IcebergFilter icebergFilter,
+ @Nullable final DateTime snapshotTime,
+ final boolean caseSensitive,
+ final InputRowSchema schema,
+ final int batchSize
+ )
+ {
+ this(table, icebergFilter, snapshotTime, caseSensitive, schema, batchSize,
ResidualFilterMode.IGNORE);
+ }
+
+ public IcebergArrowInputSourceReader(
+ final Table table,
+ @Nullable final IcebergFilter icebergFilter,
+ @Nullable final DateTime snapshotTime,
+ final boolean caseSensitive,
+ final InputRowSchema schema,
+ final int batchSize,
+ final ResidualFilterMode residualFilterMode
+ )
+ {
+ this.table = table;
+ this.icebergFilter = icebergFilter;
+ this.snapshotTime = snapshotTime;
+ this.caseSensitive = caseSensitive;
+ this.schema = schema;
+ this.batchSize = batchSize;
+ this.residualFilterMode = residualFilterMode;
+ }
+
+ @Override
+ public CloseableIterator<InputRow> read(@Nullable final InputStats
inputStats) throws IOException
+ {
+ final TableScan scan = buildScan();
+ validateNoDeleteFiles(scan);
+ validateDecimalPrecision(scan);
+ final CloseableIterable<CombinedScanTask> tasks = TableScanUtil.planTasks(
+ scan.planFiles(),
+ scan.targetSplitSize(),
+ scan.splitLookback(),
+ scan.splitOpenFileCost()
+ );
+ final ClassLoader extensionClassLoader =
IcebergArrowInputSourceReader.class.getClassLoader();
+ final ClassLoader originalClassLoader =
Thread.currentThread().getContextClassLoader();
+ ArrowReader arrowReader = null;
+ boolean ownershipTransferred = false;
+ try {
+ Thread.currentThread().setContextClassLoader(extensionClassLoader);
+ arrowReader = new ArrowReader(scan, batchSize, true);
+ final org.apache.iceberg.io.CloseableIterator<ColumnarBatch> batchIter =
arrowReader.open(tasks);
+ final CloseableIterator<InputRow> iterator = new ArrowInputRowIterator(
+ batchIter,
+ arrowReader,
+ tasks,
+ inputStats != null ? inputStats : new NoopInputStats(),
+ scan.schema(),
+ extensionClassLoader
+ );
+ ownershipTransferred = true;
+ return iterator;
+ }
+ finally {
+ Thread.currentThread().setContextClassLoader(originalClassLoader);
+ if (!ownershipTransferred) {
+ try {
+ if (arrowReader != null) {
+ arrowReader.close();
+ }
+ }
+ finally {
+ tasks.close();
+ }
+ }
+ }
+ }
+
+ private void validateNoDeleteFiles(final TableScan scan) throws IOException
+ {
+ try (CloseableIterable<FileScanTask> fileTasks = scan.planFiles()) {
+ for (FileScanTask fileTask : fileTasks) {
+ if (!fileTask.deletes().isEmpty()) {
+ throw DruidException.forPersona(DruidException.Persona.USER)
+ .ofCategory(DruidException.Category.UNSUPPORTED)
+ .build(
+ "Arrow reader does not support Iceberg
snapshots with delete files. "
+ + "Use a delete-aware input path."
+ );
+ }
+ }
+ }
+ }
+
+ private void validateDecimalPrecision(final TableScan scan)
+ {
+ for (final Types.NestedField field : scan.schema().columns()) {
+ if (field.type().typeId() == Type.TypeID.DECIMAL
+ && ((Types.DecimalType) field.type()).precision() > 18) {
+ throw DruidException.forPersona(DruidException.Persona.USER)
+ .ofCategory(DruidException.Category.UNSUPPORTED)
+ .build(
+ "Arrow reader does not support decimal fields
with precision greater than 18. "
+ + "Use the standard Iceberg reader."
+ );
+ }
+ }
+ }
+
+ @Override
+ public CloseableIterator<InputRowListPlusRawValues> sample() throws
IOException
+ {
+ final CloseableIterator<InputRow> rows = read(new NoopInputStats());
+ return new CloseableIterator<InputRowListPlusRawValues>()
+ {
+ @Override
+ public boolean hasNext()
+ {
+ return rows.hasNext();
+ }
+
+ @Override
+ public InputRowListPlusRawValues next()
+ {
+ final InputRow row = rows.next();
+ return InputRowListPlusRawValues.of(row, ((MapBasedInputRow)
row).getEvent());
+ }
+
+ @Override
+ public void close() throws IOException
+ {
+ rows.close();
+ }
+ };
+ }
+
+ private TableScan buildScan()
+ {
+ TableScan scan = table.newScan().caseSensitive(caseSensitive);
+
+ if (snapshotTime != null) {
+ scan = scan.asOfTime(snapshotTime.getMillis());
+ }
+
+ final List<String> projection = projectedColumns(scan.schema());
+ if (projection != null) {
+ scan = scan.select(projection);
+ }
+ if (icebergFilter != null) {
+ scan = icebergFilter.filter(scan);
+ if (residualFilterMode == ResidualFilterMode.IGNORE) {
+ scan = scan.ignoreResiduals();
+ }
+ }
+ return scan;
+ }
+
+ /** Projection authority is ColumnsFilter, not DimensionsSpec. Mirrors
DeltaInputSource#pruneSchema. */
+ @Nullable
+ private List<String> projectedColumns(final Schema scanSchema)
+ {
+ final ColumnsFilter filter = schema.getColumnsFilter();
+ final List<String> allColumns = scanSchema.columns().stream()
+ .map(Types.NestedField::name)
+ .collect(Collectors.toList());
+ final List<String> filtered = allColumns.stream()
+ .filter(filter::apply)
+ .collect(Collectors.toList());
+ if (filtered.equals(allColumns)) {
+ return null;
+ }
+ final String tsCol = schema.getTimestampSpec().getTimestampColumn();
+ if (tsCol != null && allColumns.contains(tsCol) &&
!filtered.contains(tsCol)) {
+ filtered.add(tsCol);
+ }
+ return filtered;
+ }
+
+ private InputRow batchRowToInputRow(
+ final ColumnarBatch batch,
+ final int rowIdx,
+ final Schema readSchema
+ )
+ {
+ final int numCols = batch.numCols();
+ final Map<String, Object> event = Maps.newHashMapWithExpectedSize(numCols);
+ for (int col = 0; col < numCols; col++) {
+ final ColumnVector column = batch.column(col);
+ final Types.NestedField field = readSchema.columns().get(col);
+ if (!column.isNullAt(rowIdx)) {
+ event.put(field.name(), extractValue(column, field.type(), rowIdx));
+ }
+ }
+ final long timestamp =
schema.getTimestampSpec().extractTimestamp(event).getMillis();
+ final List<String> dimensions = resolveDimensions(readSchema);
+ return new MapBasedInputRow(timestamp, dimensions, event);
+ }
+
+ private List<String> resolveDimensions(final Schema readSchema)
+ {
+ return MapInputRowParser.findDimensions(
+ schema.getTimestampSpec(),
+ schema.getDimensionsSpec(),
+
readSchema.columns().stream().map(Types.NestedField::name).collect(Collectors.toSet())
+ );
+ }
+
+ /**
+ * Type-safe extraction from Iceberg column accessors so physical dictionary
encoding is not exposed.
+ * Covers all scalar types supported by iceberg-arrow 1.10.0.
+ */
+ static Object extractValue(final ColumnVector column, final Type type, final
int idx)
+ {
+ switch (type.typeId()) {
+ case BOOLEAN:
+ return column.getBoolean(idx);
+ case INTEGER:
+ return column.getInt(idx);
+ case LONG:
+ return column.getLong(idx);
+ case FLOAT:
+ return (double) column.getFloat(idx);
+ case DOUBLE:
+ return column.getDouble(idx);
+ case STRING:
+ return column.getString(idx);
+ case BINARY:
+ case FIXED:
+ case UUID:
+ return column.getBinary(idx);
+ case DATE:
+ return TimeUnit.DAYS.toMillis(column.getInt(idx));
+ case TIME:
+ return TimeUnit.MICROSECONDS.toMillis(column.getLong(idx));
+ case TIMESTAMP:
+ return TimeUnit.MICROSECONDS.toMillis(column.getLong(idx));
+ case TIMESTAMP_NANO:
+ return TimeUnit.NANOSECONDS.toMillis(column.getLong(idx));
+ case DECIMAL:
+ final Types.DecimalType decimalType = (Types.DecimalType) type;
+ return column.getDecimal(idx, decimalType.precision(),
decimalType.scale());
+ default:
+ throw new IllegalArgumentException("Unsupported Iceberg type: " +
type);
+ }
+ }
+
+ private static final class NoopInputStats implements InputStats
+ {
+ @Override
+ public void incrementProcessedBytes(final long incrementByValue)
+ {
+ }
+
+ @Override
+ public long getProcessedBytes()
+ {
+ return 0;
+ }
+ }
+
+ private class ArrowInputRowIterator implements CloseableIterator<InputRow>
+ {
+ private final org.apache.iceberg.io.CloseableIterator<ColumnarBatch>
batchIter;
+ private final ArrowReader arrowReader;
+ private final CloseableIterable<CombinedScanTask> tasks;
+ private final InputStats inputStats;
+ private final Schema readSchema;
+ private final ClassLoader extensionClassLoader;
+
+ private ColumnarBatch currentBatch = null;
+ private int rowIndexInBatch = 0;
+ private boolean exhausted = false;
+
+ ArrowInputRowIterator(
+ final org.apache.iceberg.io.CloseableIterator<ColumnarBatch> batchIter,
+ final ArrowReader arrowReader,
+ final CloseableIterable<CombinedScanTask> tasks,
+ final InputStats inputStats,
+ final Schema readSchema,
+ final ClassLoader extensionClassLoader
+ )
+ {
+ this.batchIter = batchIter;
+ this.arrowReader = arrowReader;
+ this.tasks = tasks;
+ this.inputStats = inputStats;
+ this.readSchema = readSchema;
+ this.extensionClassLoader = extensionClassLoader;
+ }
+
+ @Override
+ public boolean hasNext()
+ {
+ final ClassLoader originalClassLoader =
Thread.currentThread().getContextClassLoader();
+ try {
+ Thread.currentThread().setContextClassLoader(extensionClassLoader);
+ return hasNextInternal();
+ }
+ finally {
+ Thread.currentThread().setContextClassLoader(originalClassLoader);
+ }
+ }
+
+ private boolean hasNextInternal()
+ {
+ if (exhausted) {
+ return false;
+ }
+ if (currentBatch != null && rowIndexInBatch < currentBatch.numRows()) {
+ return true;
+ }
+ return loadNextBatch();
+ }
+
+ @Override
+ public InputRow next()
+ {
+ final ClassLoader originalClassLoader =
Thread.currentThread().getContextClassLoader();
+ try {
+ Thread.currentThread().setContextClassLoader(extensionClassLoader);
+ if (!hasNextInternal()) {
+ throw new NoSuchElementException();
+ }
+ return batchRowToInputRow(currentBatch, rowIndexInBatch++, readSchema);
+ }
+ finally {
+ Thread.currentThread().setContextClassLoader(originalClassLoader);
+ }
+ }
+
+ private boolean loadNextBatch()
+ {
+ try {
+ while (batchIter.hasNext()) {
+ currentBatch = batchIter.next();
+ rowIndexInBatch = 0;
+ if (currentBatch.numRows() > 0) {
+
inputStats.incrementProcessedBytes(estimateBatchBytes(currentBatch));
+ return true;
+ }
+ }
+ }
+ catch (RuntimeException e) {
+ if (e.getMessage() != null && e.getMessage().contains("vector")) {
+ throw DruidException.forPersona(DruidException.Persona.USER)
+ .ofCategory(DruidException.Category.UNSUPPORTED)
+ .build(
+ e,
+ "Arrow reader does not support snapshots
with data files written using "
+ + "different schemas. Use the standard
Iceberg reader."
+ );
+ }
+ throw e;
+ }
+ exhausted = true;
+ return false;
+ }
+
+ private long estimateBatchBytes(final ColumnarBatch batch)
+ {
+ long bytes = 0;
+ for (int col = 0; col < batch.numCols(); col++) {
+ final FieldVector vector = batch.column(col).getFieldVector();
+ if (vector != null) {
+ bytes += vector.getBufferSize();
+ }
+ }
+ return bytes;
+ }
+
+ @Override
+ public void close() throws IOException
+ {
+ final ClassLoader originalClassLoader =
Thread.currentThread().getContextClassLoader();
+ try {
+ Thread.currentThread().setContextClassLoader(extensionClassLoader);
+ try {
+ batchIter.close();
+ }
+ finally {
+ try {
+ arrowReader.close();
+ }
+ finally {
+ tasks.close();
+ }
+ }
+ }
+ finally {
+ Thread.currentThread().setContextClassLoader(originalClassLoader);
+ }
+ }
+ }
+}
diff --git
a/extensions-contrib/druid-iceberg-extensions/src/main/java/org/apache/druid/iceberg/input/IcebergCatalog.java
b/extensions-contrib/druid-iceberg-extensions/src/main/java/org/apache/druid/iceberg/input/IcebergCatalog.java
index 2c8de41bb38..63894e18ff2 100644
---
a/extensions-contrib/druid-iceberg-extensions/src/main/java/org/apache/druid/iceberg/input/IcebergCatalog.java
+++
b/extensions-contrib/druid-iceberg-extensions/src/main/java/org/apache/druid/iceberg/input/IcebergCatalog.java
@@ -28,6 +28,7 @@ import org.apache.druid.java.util.common.RE;
import org.apache.druid.java.util.common.StringUtils;
import org.apache.druid.java.util.common.logger.Logger;
import org.apache.iceberg.FileScanTask;
+import org.apache.iceberg.Table;
import org.apache.iceberg.TableScan;
import org.apache.iceberg.catalog.Catalog;
import org.apache.iceberg.catalog.Namespace;
@@ -37,6 +38,7 @@ import org.apache.iceberg.expressions.Expressions;
import org.apache.iceberg.io.CloseableIterable;
import org.joda.time.DateTime;
+import java.io.IOException;
import java.util.ArrayList;
import java.util.List;
@@ -59,6 +61,39 @@ public abstract class IcebergCatalog
return true;
}
+ /**
+ * Load and return the Iceberg Table object for direct use by readers that
go beyond file-path delegation.
+ */
+ public Table retrieveTable(String tableNamespace, String tableName)
+ {
+ final Catalog catalog = retrieveCatalog();
+ final Namespace namespace = Namespace.of(tableNamespace);
+ final String tableIdentifier = tableNamespace + "." + tableName;
+
+ final ClassLoader currCtxClassloader =
Thread.currentThread().getContextClassLoader();
+ try {
+
Thread.currentThread().setContextClassLoader(getClass().getClassLoader());
+ final TableIdentifier icebergTableIdentifier =
catalog.listTables(namespace).stream()
+ .filter(id ->
id.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
+ ));
+ return catalog.loadTable(icebergTableIdentifier);
+ }
+ catch (IAE e) {
+ throw e;
+ }
+ catch (Exception e) {
+ throw new RE(e, "Failed to load iceberg table with identifier [%s]",
tableIdentifier);
+ }
+ finally {
+ Thread.currentThread().setContextClassLoader(currCtxClassloader);
+ }
+ }
+
/**
* Extract the iceberg data files upto the latest snapshot associated with
the table
*
@@ -106,37 +141,11 @@ public abstract class IcebergCatalog
}
tableScan = tableScan.caseSensitive(isCaseSensitive());
- CloseableIterable<FileScanTask> tasks = tableScan.planFiles();
-
- Expression detectedResidual = null;
- for (FileScanTask task : tasks) {
- dataFilePaths.add(task.file().location());
-
- // Check for residual filters
- if (detectedResidual == null) {
- Expression residual = task.residual();
- if (residual != null && !residual.equals(Expressions.alwaysTrue())) {
- detectedResidual = residual;
- }
- }
- }
-
- // Handle residual filter based on mode
- if (detectedResidual != null) {
- String message = StringUtils.format(
- "Iceberg filter produced residual expression that requires
row-level filtering. "
- + "This typically means the filter is on a non-partition column. "
- + "Residual rows may be ingested unless filtered by transformSpec.
"
- + "Residual filter: [%s]",
- detectedResidual
- );
-
- if (residualFilterMode == ResidualFilterMode.FAIL) {
- throw DruidException.forPersona(DruidException.Persona.DEVELOPER)
-
.ofCategory(DruidException.Category.RUNTIME_FAILURE)
- .build(message);
+ enforceResidualMode(tableScan, residualFilterMode);
+ try (CloseableIterable<FileScanTask> tasks = tableScan.planFiles()) {
+ for (FileScanTask task : tasks) {
+ dataFilePaths.add(task.file().location());
}
- log.warn(message);
}
long duration = System.currentTimeMillis() - start;
@@ -153,4 +162,43 @@ public abstract class IcebergCatalog
}
return dataFilePaths;
}
+
+ /**
+ * Detects whether the planned scan carries a non-trivial residual
expression (a filter that
+ * could not be fully resolved by partition pruning) and applies {@link
ResidualFilterMode}:
+ * {@code FAIL} throws a {@link DruidException}, {@code IGNORE} logs a
warning. Shared by the
+ * path-based and Arrow reader paths.
+ */
+ public void enforceResidualMode(TableScan tableScan, ResidualFilterMode
residualFilterMode)
+ {
+ Expression detectedResidual = null;
+ try (CloseableIterable<FileScanTask> tasks = tableScan.planFiles()) {
+ for (FileScanTask task : tasks) {
+ final Expression residual = task.residual();
+ if (residual != null && !residual.equals(Expressions.alwaysTrue())) {
+ detectedResidual = residual;
+ break;
+ }
+ }
+ }
+ catch (IOException e) {
+ throw new RE(e, "Failed to plan Iceberg scan for residual detection");
+ }
+ if (detectedResidual == null) {
+ return;
+ }
+ final String message = StringUtils.format(
+ "Iceberg filter produced residual expression that requires row-level
filtering. "
+ + "This typically means the filter is on a non-partition column. "
+ + "Residual rows may be ingested unless filtered by transformSpec. "
+ + "Residual filter: [%s]",
+ detectedResidual
+ );
+ if (residualFilterMode == ResidualFilterMode.FAIL) {
+ throw DruidException.forPersona(DruidException.Persona.DEVELOPER)
+ .ofCategory(DruidException.Category.RUNTIME_FAILURE)
+ .build(message);
+ }
+ log.warn(message);
+ }
}
diff --git
a/extensions-contrib/druid-iceberg-extensions/src/main/java/org/apache/druid/iceberg/input/IcebergInputSource.java
b/extensions-contrib/druid-iceberg-extensions/src/main/java/org/apache/druid/iceberg/input/IcebergInputSource.java
index ccbb10af14d..343223602de 100644
---
a/extensions-contrib/druid-iceberg-extensions/src/main/java/org/apache/druid/iceberg/input/IcebergInputSource.java
+++
b/extensions-contrib/druid-iceberg-extensions/src/main/java/org/apache/druid/iceberg/input/IcebergInputSource.java
@@ -34,51 +34,47 @@ import org.apache.druid.data.input.InputSplit;
import org.apache.druid.data.input.InputStats;
import org.apache.druid.data.input.SplitHintSpec;
import org.apache.druid.data.input.impl.SplittableInputSource;
+import org.apache.druid.error.DruidException;
import org.apache.druid.iceberg.filter.IcebergFilter;
import org.apache.druid.java.util.common.CloseableIterators;
import org.apache.druid.java.util.common.parsers.CloseableIterator;
+import org.apache.iceberg.FileFormat;
+import org.apache.iceberg.FileScanTask;
+import org.apache.iceberg.Table;
+import org.apache.iceberg.TableScan;
+import org.apache.iceberg.io.CloseableIterable;
import org.joda.time.DateTime;
import javax.annotation.Nullable;
import java.io.File;
import java.io.IOException;
+import java.io.UncheckedIOException;
import java.util.Collections;
import java.util.List;
import java.util.stream.Stream;
/**
- * Inputsource to ingest data managed by the Iceberg table format.
- * This inputsource talks to the configured catalog, executes any configured
filters and retrieves the data file paths upto the latest snapshot associated
with the iceberg table.
- * The data file paths are then provided to a native {@link
SplittableInputSource} implementation depending on the warehouse source defined.
+ * Reads an Iceberg table. Two reader modes sit behind this single type:
+ * the default resolves the snapshot to data file paths and hands them to
{@code warehouseSource},
+ * while {@code useArrowReader} scans the table directly through Iceberg's
vectorized Arrow reader.
*/
public class IcebergInputSource implements SplittableInputSource<List<String>>
{
public static final String TYPE_KEY = "iceberg";
- @JsonProperty
private final String tableName;
-
- @JsonProperty
private final String namespace;
-
- @JsonProperty
- private IcebergCatalog icebergCatalog;
-
- @JsonProperty
- private IcebergFilter icebergFilter;
-
- @JsonProperty
- private InputSourceFactory warehouseSource;
-
- @JsonProperty
+ private final IcebergCatalog icebergCatalog;
+ private final IcebergFilter icebergFilter;
private final DateTime snapshotTime;
-
- @JsonProperty
private final ResidualFilterMode residualFilterMode;
+ private final boolean useArrowReader;
+ private final int arrowBatchSize;
- private boolean isLoaded = false;
+ @Nullable
+ private final InputSourceFactory warehouseSource;
- private SplittableInputSource delegateInputSource;
+ private final InputSourceDelegate delegate;
@JsonCreator
public IcebergInputSource(
@@ -86,24 +82,101 @@ public class IcebergInputSource implements
SplittableInputSource<List<String>>
@JsonProperty("namespace") String namespace,
@JsonProperty("icebergFilter") @Nullable IcebergFilter icebergFilter,
@JsonProperty("icebergCatalog") IcebergCatalog icebergCatalog,
- @JsonProperty("warehouseSource") InputSourceFactory warehouseSource,
+ @JsonProperty("warehouseSource") @Nullable InputSourceFactory
warehouseSource,
@JsonProperty("snapshotTime") @Nullable DateTime snapshotTime,
- @JsonProperty("residualFilterMode") @Nullable ResidualFilterMode
residualFilterMode
+ @JsonProperty("residualFilterMode") @Nullable ResidualFilterMode
residualFilterMode,
+ @JsonProperty("useArrowReader") @Nullable Boolean useArrowReader,
+ @JsonProperty("arrowBatchSize") @Nullable Integer arrowBatchSize
)
{
this.tableName = Preconditions.checkNotNull(tableName, "tableName cannot
be null");
this.namespace = Preconditions.checkNotNull(namespace, "namespace cannot
be null");
this.icebergCatalog = Preconditions.checkNotNull(icebergCatalog,
"icebergCatalog cannot be null");
this.icebergFilter = icebergFilter;
- this.warehouseSource = Preconditions.checkNotNull(warehouseSource,
"warehouseSource cannot be null");
this.snapshotTime = snapshotTime;
this.residualFilterMode = Configs.valueOrDefault(residualFilterMode,
ResidualFilterMode.IGNORE);
+ this.useArrowReader = Boolean.TRUE.equals(useArrowReader);
+ this.arrowBatchSize = arrowBatchSize != null && arrowBatchSize > 0
+ ? arrowBatchSize
+ : IcebergArrowInputSourceReader.DEFAULT_BATCH_SIZE;
+ this.warehouseSource = warehouseSource;
+
+ this.delegate = this.useArrowReader
+ ? new ArrowDelegate()
+ : new StandardDelegate(
+ Preconditions.checkNotNull(
+ warehouseSource,
+ "warehouseSource cannot be null unless
useArrowReader is true"
+ )
+ );
+ }
+
+ @JsonProperty
+ public String getTableName()
+ {
+ return tableName;
+ }
+
+ @JsonProperty
+ public String getNamespace()
+ {
+ return namespace;
+ }
+
+ @JsonProperty
+ public IcebergCatalog getIcebergCatalog()
+ {
+ return icebergCatalog;
+ }
+
+ @JsonProperty
+ public IcebergFilter getIcebergFilter()
+ {
+ return icebergFilter;
+ }
+
+ @Nullable
+ @JsonProperty
+ public DateTime getSnapshotTime()
+ {
+ return snapshotTime;
+ }
+
+ @JsonProperty
+ public ResidualFilterMode getResidualFilterMode()
+ {
+ return residualFilterMode;
+ }
+
+ @Nullable
+ @JsonProperty("warehouseSource")
+ public InputSourceFactory getWarehouseSource()
+ {
+ return warehouseSource;
+ }
+
+ @JsonProperty
+ public boolean isUseArrowReader()
+ {
+ return useArrowReader;
+ }
+
+ @JsonProperty
+ public int getArrowBatchSize()
+ {
+ return arrowBatchSize;
}
@Override
public boolean needsFormat()
{
- return true;
+ return delegate.needsFormat();
+ }
+
+ @Override
+ public boolean isSplittable()
+ {
+ return delegate.isSplittable();
}
@Override
@@ -113,10 +186,7 @@ public class IcebergInputSource implements
SplittableInputSource<List<String>>
File temporaryDirectory
)
{
- if (!isLoaded) {
- retrieveIcebergDatafiles();
- }
- return getDelegateInputSource().reader(inputRowSchema, inputFormat,
temporaryDirectory);
+ return delegate.reader(inputRowSchema, inputFormat, temporaryDirectory);
}
@Override
@@ -125,97 +195,249 @@ public class IcebergInputSource implements
SplittableInputSource<List<String>>
@Nullable SplitHintSpec splitHintSpec
) throws IOException
{
- if (!isLoaded) {
- retrieveIcebergDatafiles();
- }
- return getDelegateInputSource().createSplits(inputFormat, splitHintSpec);
+ return delegate.createSplits(inputFormat, splitHintSpec);
}
@Override
public int estimateNumSplits(InputFormat inputFormat, @Nullable
SplitHintSpec splitHintSpec) throws IOException
{
- if (!isLoaded) {
- retrieveIcebergDatafiles();
- }
- return getDelegateInputSource().estimateNumSplits(inputFormat,
splitHintSpec);
+ return delegate.estimateNumSplits(inputFormat, splitHintSpec);
}
@Override
public InputSource withSplit(InputSplit<List<String>> inputSplit)
{
- return getDelegateInputSource().withSplit(inputSplit);
+ return delegate.withSplit(inputSplit);
}
@Override
public SplitHintSpec getSplitHintSpecOrDefault(@Nullable SplitHintSpec
splitHintSpec)
{
- return getDelegateInputSource().getSplitHintSpecOrDefault(splitHintSpec);
+ return delegate.getSplitHintSpecOrDefault(splitHintSpec);
}
- @JsonProperty
- public String getTableName()
+ private Table retrieveTable()
{
- return tableName;
+ return icebergCatalog.retrieveTable(namespace, tableName);
}
- @JsonProperty
- public String getNamespace()
+ /**
+ * Mode-specific behavior. The two modes differ on more than how rows are
read: they disagree on whether
+ * an {@link InputFormat} is needed and whether the source can be split
across tasks.
+ */
+ private interface InputSourceDelegate
{
- return namespace;
- }
+ boolean needsFormat();
- @JsonProperty
- public IcebergCatalog getIcebergCatalog()
- {
- return icebergCatalog;
- }
+ boolean isSplittable();
- @JsonProperty
- public IcebergFilter getIcebergFilter()
- {
- return icebergFilter;
- }
+ InputSourceReader reader(
+ InputRowSchema inputRowSchema,
+ @Nullable InputFormat inputFormat,
+ File temporaryDirectory
+ );
- @Nullable
- @JsonProperty
- public DateTime getSnapshotTime()
- {
- return snapshotTime;
- }
+ Stream<InputSplit<List<String>>> createSplits(
+ InputFormat inputFormat,
+ @Nullable SplitHintSpec splitHintSpec
+ ) throws IOException;
- @JsonProperty
- public ResidualFilterMode getResidualFilterMode()
- {
- return residualFilterMode;
- }
+ int estimateNumSplits(InputFormat inputFormat, @Nullable SplitHintSpec
splitHintSpec) throws IOException;
- public SplittableInputSource getDelegateInputSource()
- {
- return delegateInputSource;
+ InputSource withSplit(InputSplit<List<String>> inputSplit);
+
+ SplitHintSpec getSplitHintSpecOrDefault(@Nullable SplitHintSpec
splitHintSpec);
}
- protected void retrieveIcebergDatafiles()
+ /**
+ * Resolves the snapshot to a list of data file paths and defers reading to
the warehouse input source.
+ */
+ private class StandardDelegate implements InputSourceDelegate
{
- List<String> snapshotDataFiles = icebergCatalog.extractSnapshotDataFiles(
- getNamespace(),
- getTableName(),
- getIcebergFilter(),
- getSnapshotTime(),
- getResidualFilterMode()
- );
- if (snapshotDataFiles.isEmpty()) {
- delegateInputSource = new EmptyInputSource();
- } else {
- delegateInputSource = warehouseSource.create(snapshotDataFiles);
+ private final InputSourceFactory warehouseSource;
+
+ private boolean isLoaded = false;
+ private SplittableInputSource delegateInputSource;
+
+ StandardDelegate(final InputSourceFactory warehouseSource)
+ {
+ this.warehouseSource = warehouseSource;
+ }
+
+ @Override
+ public boolean needsFormat()
+ {
+ return true;
+ }
+
+ @Override
+ public boolean isSplittable()
+ {
+ return true;
+ }
+
+ @Override
+ public InputSourceReader reader(
+ InputRowSchema inputRowSchema,
+ @Nullable InputFormat inputFormat,
+ File temporaryDirectory
+ )
+ {
+ return warehouseInputSource().reader(inputRowSchema, inputFormat,
temporaryDirectory);
+ }
+
+ @Override
+ public Stream<InputSplit<List<String>>> createSplits(
+ InputFormat inputFormat,
+ @Nullable SplitHintSpec splitHintSpec
+ ) throws IOException
+ {
+ return warehouseInputSource().createSplits(inputFormat, splitHintSpec);
+ }
+
+ @Override
+ public int estimateNumSplits(InputFormat inputFormat, @Nullable
SplitHintSpec splitHintSpec) throws IOException
+ {
+ return warehouseInputSource().estimateNumSplits(inputFormat,
splitHintSpec);
+ }
+
+ @Override
+ public InputSource withSplit(InputSplit<List<String>> inputSplit)
+ {
+ return warehouseInputSource().withSplit(inputSplit);
+ }
+
+ @Override
+ public SplitHintSpec getSplitHintSpecOrDefault(@Nullable SplitHintSpec
splitHintSpec)
+ {
+ return warehouseInputSource().getSplitHintSpecOrDefault(splitHintSpec);
+ }
+
+ private SplittableInputSource warehouseInputSource()
+ {
+ if (!isLoaded) {
+ final List<String> snapshotDataFiles =
icebergCatalog.extractSnapshotDataFiles(
+ getNamespace(),
+ getTableName(),
+ getIcebergFilter(),
+ getSnapshotTime(),
+ getResidualFilterMode()
+ );
+ if (snapshotDataFiles.isEmpty()) {
+ delegateInputSource = new EmptyInputSource();
+ } else {
+ delegateInputSource = warehouseSource.create(snapshotDataFiles);
+ }
+ isLoaded = true;
+ }
+ return delegateInputSource;
}
- isLoaded = true;
}
/**
- * This input source is used in place of a delegate input source if there
are no input file paths.
- * Certain input sources cannot be instantiated with an empty input file
list and so composing input sources such as IcebergInputSource
- * may use this input source as delegate in such cases.
+ * Scans the table through Iceberg's vectorized Arrow reader. Parquet only,
and not splittable:
+ * there is no data file list to hand out, so all rows are read by a single
task.
*/
+ private class ArrowDelegate implements InputSourceDelegate
+ {
+ @Override
+ public boolean needsFormat()
+ {
+ return false;
+ }
+
+ @Override
+ public boolean isSplittable()
+ {
+ return false;
+ }
+
+ @Override
+ public InputSourceReader reader(
+ InputRowSchema inputRowSchema,
+ @Nullable InputFormat inputFormat,
+ File temporaryDirectory
+ )
+ {
+ final Table table = retrieveTable();
+ TableScan scan =
table.newScan().caseSensitive(icebergCatalog.isCaseSensitive());
+ if (icebergFilter != null) {
+ scan = icebergFilter.filter(scan);
+ }
+ if (snapshotTime != null) {
+ scan = scan.asOfTime(snapshotTime.getMillis());
+ }
+ if (icebergFilter != null) {
+ icebergCatalog.enforceResidualMode(scan, residualFilterMode);
+ }
+ validateParquetOnly(scan);
+
+ return new IcebergArrowInputSourceReader(
+ table,
+ icebergFilter,
+ snapshotTime,
+ icebergCatalog.isCaseSensitive(),
+ inputRowSchema,
+ arrowBatchSize,
+ residualFilterMode
+ );
+ }
+
+ @Override
+ public Stream<InputSplit<List<String>>> createSplits(
+ InputFormat inputFormat,
+ @Nullable SplitHintSpec splitHintSpec
+ )
+ {
+ return Stream.of(new InputSplit<>(Collections.emptyList()));
+ }
+
+ @Override
+ public int estimateNumSplits(InputFormat inputFormat, @Nullable
SplitHintSpec splitHintSpec)
+ {
+ return 1;
+ }
+
+ @Override
+ public InputSource withSplit(InputSplit<List<String>> inputSplit)
+ {
+ return IcebergInputSource.this;
+ }
+
+ @Override
+ public SplitHintSpec getSplitHintSpecOrDefault(@Nullable SplitHintSpec
splitHintSpec)
+ {
+ return splitHintSpec == null ?
SplittableInputSource.DEFAULT_SPLIT_HINT_SPEC : splitHintSpec;
+ }
+
+ /**
+ * Iceberg's Arrow reader only supports Parquet. Checked against table
metadata so a mixed-format or
+ * ORC/Avro table fails with a clear message instead of an obscure error
inside Arrow.
+ */
+ private void validateParquetOnly(final TableScan scan)
+ {
+ try (CloseableIterable<FileScanTask> fileTasks = scan.planFiles()) {
+ for (final FileScanTask fileTask : fileTasks) {
+ final FileFormat format = fileTask.file().format();
+ if (format != FileFormat.PARQUET) {
+ throw DruidException.forPersona(DruidException.Persona.USER)
+
.ofCategory(DruidException.Category.UNSUPPORTED)
+ .build(
+ "Arrow reader supports only Parquet data
files, but table[%s.%s] has a"
+ + " data file in format[%s]. Set
useArrowReader to false for this table.",
+ namespace,
+ tableName,
+ format
+ );
+ }
+ }
+ }
+ catch (IOException e) {
+ throw new UncheckedIOException(e);
+ }
+ }
+ }
+
private static class EmptyInputSource implements SplittableInputSource
{
@Override
@@ -242,15 +464,13 @@ public class IcebergInputSource implements
SplittableInputSource<List<String>>
@Override
public CloseableIterator<InputRow> read(InputStats inputStats)
{
- return CloseableIterators.wrap(Collections.emptyIterator(), () -> {
- });
+ return CloseableIterators.wrap(Collections.emptyIterator(), () ->
{});
}
@Override
public CloseableIterator<InputRowListPlusRawValues> sample()
{
- return CloseableIterators.wrap(Collections.emptyIterator(), () -> {
- });
+ return CloseableIterators.wrap(Collections.emptyIterator(), () ->
{});
}
};
}
@@ -273,7 +493,7 @@ public class IcebergInputSource implements
SplittableInputSource<List<String>>
@Override
public InputSource withSplit(InputSplit split)
{
- return null;
+ return this;
}
}
}
diff --git
a/extensions-contrib/druid-iceberg-extensions/src/test/java/org/apache/druid/iceberg/input/IcebergArrowInputSourceReaderTest.java
b/extensions-contrib/druid-iceberg-extensions/src/test/java/org/apache/druid/iceberg/input/IcebergArrowInputSourceReaderTest.java
new file mode 100644
index 00000000000..5ebb71c7bf5
--- /dev/null
+++
b/extensions-contrib/druid-iceberg-extensions/src/test/java/org/apache/druid/iceberg/input/IcebergArrowInputSourceReaderTest.java
@@ -0,0 +1,614 @@
+/*
+ * 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 com.google.common.collect.ImmutableList;
+import com.google.common.collect.ImmutableSet;
+import org.apache.druid.data.input.ColumnsFilter;
+import org.apache.druid.data.input.InputRow;
+import org.apache.druid.data.input.InputRowSchema;
+import org.apache.druid.data.input.MapBasedInputRow;
+import org.apache.druid.data.input.impl.DimensionsSpec;
+import org.apache.druid.data.input.impl.MapInputRowParser;
+import org.apache.druid.data.input.impl.StringDimensionSchema;
+import org.apache.druid.data.input.impl.TimestampSpec;
+import org.apache.druid.error.DruidException;
+import org.apache.druid.iceberg.filter.IcebergEqualsFilter;
+import org.apache.druid.java.util.common.DateTimes;
+import org.apache.druid.java.util.common.FileUtils;
+import org.apache.druid.java.util.common.parsers.CloseableIterator;
+import org.apache.iceberg.DataFile;
+import org.apache.iceberg.DataFiles;
+import org.apache.iceberg.FileFormat;
+import org.apache.iceberg.PartitionSpec;
+import org.apache.iceberg.Schema;
+import org.apache.iceberg.Table;
+import org.apache.iceberg.catalog.Namespace;
+import org.apache.iceberg.catalog.TableIdentifier;
+import org.apache.iceberg.data.GenericRecord;
+import org.apache.iceberg.data.parquet.GenericParquetWriter;
+import org.apache.iceberg.io.DataWriter;
+import org.apache.iceberg.io.OutputFile;
+import org.apache.iceberg.parquet.Parquet;
+import org.apache.iceberg.types.Types;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+
+import java.io.File;
+import java.io.IOException;
+import java.math.BigDecimal;
+import java.util.ArrayList;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.UUID;
+
+public class IcebergArrowInputSourceReaderTest
+{
+ private static final String NAMESPACE = "default";
+ private static final String TABLE = "arrowTestTable";
+
+ private static final Schema SCHEMA = new Schema(
+ Types.NestedField.required(1, "ts", Types.LongType.get()),
+ Types.NestedField.required(2, "name", Types.StringType.get()),
+ Types.NestedField.required(3, "value", Types.DoubleType.get())
+ );
+
+ private static final InputRowSchema INPUT_SCHEMA = new InputRowSchema(
+ new TimestampSpec("ts", "millis", null),
+ DimensionsSpec.builder()
+ .setDimensions(ImmutableList.of(
+ new StringDimensionSchema("name")
+ ))
+ .build(),
+ ColumnsFilter.all()
+ );
+
+ private File warehouseDir;
+ private IcebergCatalog catalog;
+ private TableIdentifier tableId;
+
+ @BeforeEach
+ public void setup() throws IOException
+ {
+ warehouseDir = FileUtils.createTempDir();
+ catalog = new LocalCatalog(warehouseDir.getPath(), new HashMap<>(), true);
+ tableId = TableIdentifier.of(Namespace.of(NAMESPACE), TABLE);
+ }
+
+ @AfterEach
+ public void tearDown()
+ {
+ if (catalog.retrieveCatalog().tableExists(tableId)) {
+ catalog.retrieveCatalog().dropTable(tableId);
+ }
+ }
+
+ @Test
+ public void testBasicRead() throws IOException
+ {
+ final Table table = catalog.retrieveCatalog().createTable(tableId, SCHEMA);
+ writeRows(table, row(1_000L, "alice", 1.1), row(2_000L, "bob", 2.2),
row(3_000L, "carol", 3.3));
+
+ final IcebergArrowInputSourceReader reader = new
IcebergArrowInputSourceReader(
+ table,
+ null,
+ null,
+ true,
+ INPUT_SCHEMA,
+ IcebergArrowInputSourceReader.DEFAULT_BATCH_SIZE
+ );
+
+ final List<InputRow> rows = readAll(reader);
+ Assertions.assertEquals(3, rows.size());
+ Assertions.assertEquals(1_000L, rows.get(0).getTimestampFromEpoch());
+ Assertions.assertEquals("alice", rows.get(0).getDimension("name").get(0));
+ Assertions.assertEquals(2_000L, rows.get(1).getTimestampFromEpoch());
+ Assertions.assertEquals("bob", rows.get(1).getDimension("name").get(0));
+ Assertions.assertEquals(3_000L, rows.get(2).getTimestampFromEpoch());
+ Assertions.assertEquals("carol", rows.get(2).getDimension("name").get(0));
+ }
+
+ @Test
+ public void testHighPrecisionDecimalIsRejected() throws IOException
+ {
+ final Schema decimalSchema = new Schema(
+ Types.NestedField.required(1, "ts", Types.LongType.get()),
+ Types.NestedField.required(2, "amount", Types.DecimalType.of(20, 2))
+ );
+ final Table table = catalog.retrieveCatalog().createTable(tableId,
decimalSchema);
+ final BigDecimal value = new BigDecimal("123456789012345678.90");
+ final GenericRecord record = GenericRecord.create(decimalSchema);
+ record.setField("ts", 1_000L);
+ record.setField("amount", value);
+ writeRows(table, decimalSchema, record);
+
+ final InputRowSchema inputRowSchema = new InputRowSchema(
+ new TimestampSpec("ts", "millis", null),
+ DimensionsSpec.builder().build(),
+ ColumnsFilter.all()
+ );
+ final IcebergArrowInputSourceReader reader = new
IcebergArrowInputSourceReader(
+ table,
+ null,
+ null,
+ true,
+ inputRowSchema,
+ IcebergArrowInputSourceReader.DEFAULT_BATCH_SIZE
+ );
+
+ final DruidException exception =
Assertions.assertThrows(DruidException.class, reader::read);
+ Assertions.assertTrue(exception.getMessage().contains("precision greater
than 18"));
+ }
+
+ @Test
+ public void testEmptyTable() throws IOException
+ {
+ catalog.retrieveCatalog().createTable(tableId, SCHEMA);
+ final Table table = catalog.retrieveTable(NAMESPACE, TABLE);
+
+ final IcebergArrowInputSourceReader reader = new
IcebergArrowInputSourceReader(
+ table,
+ null,
+ null,
+ true,
+ INPUT_SCHEMA,
+ IcebergArrowInputSourceReader.DEFAULT_BATCH_SIZE
+ );
+
+ final List<InputRow> rows = readAll(reader);
+ Assertions.assertEquals(0, rows.size());
+ }
+
+ @Test
+ public void testWithEqualsFilter() throws IOException
+ {
+ final Table table = catalog.retrieveCatalog().createTable(tableId, SCHEMA);
+ writeRows(
+ table,
+ row(1_000L, "alice", 1.0),
+ row(2_000L, "bob", 2.0),
+ row(3_000L, "alice", 3.0)
+ );
+
+ final IcebergArrowInputSourceReader reader = new
IcebergArrowInputSourceReader(
+ table,
+ new IcebergEqualsFilter("name", "alice"),
+ null,
+ true,
+ INPUT_SCHEMA,
+ IcebergArrowInputSourceReader.DEFAULT_BATCH_SIZE
+ );
+
+ final List<InputRow> rows = readAll(reader);
+ Assertions.assertEquals(3, rows.size());
+ Assertions.assertEquals("alice", rows.get(0).getDimension("name").get(0));
+ Assertions.assertEquals("bob", rows.get(1).getDimension("name").get(0));
+ Assertions.assertEquals("alice", rows.get(2).getDimension("name").get(0));
+ }
+
+ @Test
+ public void testColumnPruning() throws IOException
+ {
+ final Table table = catalog.retrieveCatalog().createTable(tableId, SCHEMA);
+ writeRows(table, row(1_000L, "alice", 9.9));
+
+ final InputRowSchema pruned = new InputRowSchema(
+ new TimestampSpec("ts", "millis", null),
+ DimensionsSpec.builder()
+ .setDimensions(ImmutableList.of(new
StringDimensionSchema("name")))
+ .build(),
+ ColumnsFilter.all()
+ );
+
+ final IcebergArrowInputSourceReader reader = new
IcebergArrowInputSourceReader(
+ table,
+ null,
+ null,
+ true,
+ pruned,
+ IcebergArrowInputSourceReader.DEFAULT_BATCH_SIZE
+ );
+
+ final List<InputRow> rows = readAll(reader);
+ Assertions.assertEquals(1, rows.size());
+ Assertions.assertEquals(1_000L, rows.get(0).getTimestampFromEpoch());
+ Assertions.assertEquals("alice", rows.get(0).getDimension("name").get(0));
+ }
+
+ @Test
+ public void testDynamicDimensionsWithExplicitDimensions() throws IOException
+ {
+ final Table table = catalog.retrieveCatalog().createTable(tableId, SCHEMA);
+ writeRows(table, row(1_000L, "alice", 1.0));
+ final InputRowSchema inputRowSchema = new InputRowSchema(
+ new TimestampSpec("ts", "millis", null),
+ DimensionsSpec.builder()
+ .setDimensions(ImmutableList.of(new
StringDimensionSchema("name")))
+ .setIncludeAllDimensions(true)
+ .build(),
+ ColumnsFilter.all()
+ );
+
+ final IcebergArrowInputSourceReader reader = new
IcebergArrowInputSourceReader(
+ table,
+ null,
+ null,
+ true,
+ inputRowSchema,
+ IcebergArrowInputSourceReader.DEFAULT_BATCH_SIZE
+ );
+
+ final List<InputRow> rows = readAll(reader);
+ Assertions.assertEquals(1, rows.size());
+ Assertions.assertEquals(
+ MapInputRowParser.findDimensions(
+ inputRowSchema.getTimestampSpec(),
+ inputRowSchema.getDimensionsSpec(),
+ ImmutableSet.of("ts", "name", "value")
+ ),
+ rows.get(0).getDimensions()
+ );
+ }
+
+ @Test
+ public void testDynamicDimensionsHonorExclusions() throws IOException
+ {
+ final Table table = catalog.retrieveCatalog().createTable(tableId, SCHEMA);
+ writeRows(table, row(1_000L, "alice", 1.0));
+ final InputRowSchema inputRowSchema = new InputRowSchema(
+ new TimestampSpec("ts", "millis", null),
+ DimensionsSpec.builder()
+ .setDimensionExclusions(ImmutableList.of("value"))
+ .build(),
+ ColumnsFilter.all()
+ );
+
+ final IcebergArrowInputSourceReader reader = new
IcebergArrowInputSourceReader(
+ table,
+ null,
+ null,
+ true,
+ inputRowSchema,
+ IcebergArrowInputSourceReader.DEFAULT_BATCH_SIZE
+ );
+
+ final List<InputRow> rows = readAll(reader);
+ Assertions.assertEquals(1, rows.size());
+ Assertions.assertEquals(
+ MapInputRowParser.findDimensions(
+ inputRowSchema.getTimestampSpec(),
+ inputRowSchema.getDimensionsSpec(),
+ ImmutableSet.of("ts", "name", "value")
+ ),
+ rows.get(0).getDimensions()
+ );
+ }
+
+ @Test
+ public void testLargeBatch() throws IOException
+ {
+ final Table table = catalog.retrieveCatalog().createTable(tableId, SCHEMA);
+ final int count = 5_000;
+ final GenericRecord[] data = new GenericRecord[count];
+ for (int i = 0; i < count; i++) {
+ data[i] = row((long) (i + 1) * 1000, "user" + i, i * 0.1);
+ }
+ writeRows(table, data);
+
+ final IcebergArrowInputSourceReader reader = new
IcebergArrowInputSourceReader(
+ table,
+ null,
+ null,
+ true,
+ INPUT_SCHEMA,
+ IcebergArrowInputSourceReader.DEFAULT_BATCH_SIZE
+ );
+
+ final List<InputRow> rows = readAll(reader);
+ Assertions.assertEquals(count, rows.size());
+ }
+
+ @Test
+ public void testSnapshotTime() throws IOException, InterruptedException
+ {
+ final Table table = catalog.retrieveCatalog().createTable(tableId, SCHEMA);
+ writeRows(table, row(1_000L, "snap1", 1.0));
+ final long afterFirstSnapshot = System.currentTimeMillis();
+
+ Thread.sleep(10);
+ writeRows(table, row(2_000L, "snap2", 2.0));
+
+ final IcebergArrowInputSourceReader reader = new
IcebergArrowInputSourceReader(
+ table,
+ null,
+ DateTimes.utc(afterFirstSnapshot),
+ true,
+ INPUT_SCHEMA,
+ IcebergArrowInputSourceReader.DEFAULT_BATCH_SIZE
+ );
+
+ final List<InputRow> rows = readAll(reader);
+ Assertions.assertEquals(1, rows.size());
+ Assertions.assertEquals("snap1", rows.get(0).getDimension("name").get(0));
+ }
+
+ @Test
+ public void testSnapshotTimeUsesHistoricalSchema() throws IOException,
InterruptedException
+ {
+ final Table table = catalog.retrieveCatalog().createTable(tableId, SCHEMA);
+ writeRows(table, row(1_000L, "beforeRename", 1.0));
+ final long afterFirstSnapshot = System.currentTimeMillis();
+
+ Thread.sleep(10);
+ table.updateSchema().renameColumn("name", "display_name").commit();
+
+ final IcebergArrowInputSourceReader reader = new
IcebergArrowInputSourceReader(
+ table,
+ null,
+ DateTimes.utc(afterFirstSnapshot),
+ true,
+ INPUT_SCHEMA,
+ IcebergArrowInputSourceReader.DEFAULT_BATCH_SIZE
+ );
+
+ final List<InputRow> rows = readAll(reader);
+ Assertions.assertEquals(1, rows.size());
+ Assertions.assertEquals("beforeRename",
rows.get(0).getDimension("name").get(0));
+ }
+
+ @Test
+ public void
testArrowReaderWorksWhenContextClassLoaderCannotLoadArrowFormatModels() throws
IOException
+ {
+ final Table table = catalog.retrieveCatalog().createTable(tableId, SCHEMA);
+ writeRows(table, row(1_000L, "alice", 1.0));
+
+ final IcebergArrowInputSourceReader reader = new
IcebergArrowInputSourceReader(
+ table,
+ null,
+ null,
+ true,
+ INPUT_SCHEMA,
+ IcebergArrowInputSourceReader.DEFAULT_BATCH_SIZE
+ );
+
+ final ClassLoader originalClassLoader =
Thread.currentThread().getContextClassLoader();
+ final ClassLoader blockingClassLoader = new
ClassLoader(originalClassLoader)
+ {
+ @Override
+ protected Class<?> loadClass(final String name, final boolean resolve)
throws ClassNotFoundException
+ {
+ if
("org.apache.iceberg.arrow.vectorized.ArrowFormatModels".equals(name)) {
+ throw new ClassNotFoundException(name);
+ }
+ return super.loadClass(name, resolve);
+ }
+ };
+
+ try {
+ Thread.currentThread().setContextClassLoader(blockingClassLoader);
+ final List<InputRow> rows = readAll(reader);
+ Assertions.assertEquals(1, rows.size());
+ Assertions.assertEquals("alice",
rows.get(0).getDimension("name").get(0));
+ }
+ finally {
+ Thread.currentThread().setContextClassLoader(originalClassLoader);
+ }
+ }
+
+ @Test
+ public void testArrowReaderClosesTasksWhenSetupFails() throws IOException
+ {
+ final Table table = catalog.retrieveCatalog().createTable(tableId, SCHEMA);
+ final DataFile unsupportedFile = DataFiles.builder(table.spec())
+ .withPath(table.location() +
"/unsupported.orc")
+ .withFormat(FileFormat.ORC)
+ .withRecordCount(1)
+ .withFileSizeInBytes(1)
+ .build();
+ table.newAppend().appendFile(unsupportedFile).commit();
+
+ final IcebergArrowInputSourceReader reader = new
IcebergArrowInputSourceReader(
+ table,
+ null,
+ null,
+ true,
+ INPUT_SCHEMA,
+ IcebergArrowInputSourceReader.DEFAULT_BATCH_SIZE
+ );
+
+ Assertions.assertThrows(UnsupportedOperationException.class, () ->
readAll(reader));
+ }
+
+ @Test
+ public void testSchemaEvolutionRejectsFieldsMissingFromOlderFiles() throws
IOException
+ {
+ final Table table = catalog.retrieveCatalog().createTable(tableId, SCHEMA);
+ writeRows(table, row(1_000L, "old", 1.0));
+
+ table.updateSchema().addColumn("category",
Types.StringType.get()).commit();
+ final Schema evolvedSchema = table.schema();
+ final GenericRecord evolvedRow = GenericRecord.create(evolvedSchema);
+ evolvedRow.setField("ts", 2_000L);
+ evolvedRow.setField("name", "new");
+ evolvedRow.setField("value", 2.0);
+ evolvedRow.setField("category", "updated");
+ writeRows(table, evolvedSchema, evolvedRow);
+
+ final InputRowSchema inferredDimensionsSchema = new InputRowSchema(
+ new TimestampSpec("ts", "millis", null),
+ DimensionsSpec.builder().build(),
+ ColumnsFilter.all()
+ );
+ final IcebergArrowInputSourceReader reader = new
IcebergArrowInputSourceReader(
+ table,
+ null,
+ null,
+ true,
+ inferredDimensionsSchema,
+ IcebergArrowInputSourceReader.DEFAULT_BATCH_SIZE
+ );
+
+ final DruidException exception =
Assertions.assertThrows(DruidException.class, () -> readAll(reader));
+ Assertions.assertTrue(exception.getMessage().contains("different
schemas"));
+ }
+
+ @Test
+ public void testAggregatorSourceColumnSurvivesProjection() throws IOException
+ {
+ // Regression: dimensions=[name] plus ColumnsFilter inclusion of `value`
(aggregator source).
+ final Table table = catalog.retrieveCatalog().createTable(tableId, SCHEMA);
+ writeRows(table, row(1_000L, "alice", 9.0), row(2_000L, "bob", 4.5));
+
+ final InputRowSchema schemaWithAggSource = new InputRowSchema(
+ new TimestampSpec("ts", "millis", null),
+ DimensionsSpec.builder()
+ .setDimensions(ImmutableList.of(new
StringDimensionSchema("name")))
+ .build(),
+ ColumnsFilter.inclusionBased(ImmutableSet.of("ts", "name", "value"))
+ );
+
+ final IcebergArrowInputSourceReader reader = new
IcebergArrowInputSourceReader(
+ table,
+ null,
+ null,
+ true,
+ schemaWithAggSource,
+ IcebergArrowInputSourceReader.DEFAULT_BATCH_SIZE
+ );
+
+ final List<InputRow> rows = readAll(reader);
+ Assertions.assertEquals(2, rows.size());
+ final Map<String, Object> event0 = ((MapBasedInputRow)
rows.get(0)).getEvent();
+ Assertions.assertEquals("alice", event0.get("name"));
+ Assertions.assertNotNull(event0.get("value"), "aggregator source column
'value' must survive projection");
+ Assertions.assertEquals(9.0, ((Number) event0.get("value")).doubleValue(),
0.0001);
+ final Map<String, Object> event1 = ((MapBasedInputRow)
rows.get(1)).getEvent();
+ Assertions.assertEquals("bob", event1.get("name"));
+ Assertions.assertNotNull(event1.get("value"), "aggregator source column
'value' must survive projection");
+ Assertions.assertEquals(4.5, ((Number) event1.get("value")).doubleValue(),
0.0001);
+ }
+
+ @Test
+ public void testProjectionPrunesUnusedColumns() throws IOException
+ {
+ final Table table = catalog.retrieveCatalog().createTable(tableId, SCHEMA);
+ writeRows(table, row(1_000L, "alice", 7.0), row(2_000L, "bob", 8.0));
+
+ // Exclusion-based filter must push projection so excluded columns are
never read.
+ final InputRowSchema prunedSchema = new InputRowSchema(
+ new TimestampSpec("ts", "millis", null),
+ DimensionsSpec.builder()
+ .setDimensions(ImmutableList.of(new
StringDimensionSchema("name")))
+ .build(),
+ ColumnsFilter.exclusionBased(ImmutableSet.of("value"))
+ );
+
+ final IcebergArrowInputSourceReader reader = new
IcebergArrowInputSourceReader(
+ table,
+ null,
+ null,
+ true,
+ prunedSchema,
+ IcebergArrowInputSourceReader.DEFAULT_BATCH_SIZE
+ );
+
+ final List<InputRow> rows = readAll(reader);
+ Assertions.assertEquals(2, rows.size());
+ for (final InputRow r : rows) {
+ final Map<String, Object> event = ((MapBasedInputRow) r).getEvent();
+ Assertions.assertNull(
+ event.get("value"),
+ "excluded column 'value' must be pruned at scan and absent from
event"
+ );
+ Assertions.assertNotNull(event.get("name"), "included column 'name' must
be present");
+ }
+ }
+
+ // --- helpers ---
+
+ private static GenericRecord row(final long ts, final String name, final
double value)
+ {
+ final GenericRecord r = GenericRecord.create(SCHEMA);
+ r.setField("ts", ts);
+ r.setField("name", name);
+ r.setField("value", value);
+ return r;
+ }
+
+ private static void writeRows(final Table table, final GenericRecord...
records) throws IOException
+ {
+ writeRows(table, SCHEMA, records);
+ }
+
+ private static void writeRows(
+ final Table table,
+ final Schema dataSchema,
+ final GenericRecord... records
+ ) throws IOException
+ {
+ final String filepath = table.location() + "/" + UUID.randomUUID() +
".parquet";
+ final OutputFile file = table.io().newOutputFile(filepath);
+ final DataWriter<GenericRecord> writer =
+ Parquet.writeData(file)
+ .schema(dataSchema)
+ .createWriterFunc(GenericParquetWriter::create)
+ .overwrite()
+ .withSpec(PartitionSpec.unpartitioned())
+ .build();
+ try {
+ for (final GenericRecord r : records) {
+ writer.write(r);
+ }
+ }
+ finally {
+ writer.close();
+ }
+ final DataFile dataFile = writer.toDataFile();
+ table.newAppend().appendFile(dataFile).commit();
+ }
+
+ private static List<InputRow> readAll(final IcebergArrowInputSourceReader
reader) throws IOException
+ {
+ final List<InputRow> result = new ArrayList<>();
+ try (CloseableIterator<InputRow> it = reader.read(new NoopInputStats())) {
+ while (it.hasNext()) {
+ result.add(it.next());
+ }
+ }
+ return result;
+ }
+
+ private static final class NoopInputStats implements
org.apache.druid.data.input.InputStats
+ {
+ @Override
+ public void incrementProcessedBytes(final long v)
+ {
+ }
+
+ @Override
+ public long getProcessedBytes()
+ {
+ return 0;
+ }
+ }
+}
diff --git
a/extensions-contrib/druid-iceberg-extensions/src/test/java/org/apache/druid/iceberg/input/IcebergInputSourceArrowModeTest.java
b/extensions-contrib/druid-iceberg-extensions/src/test/java/org/apache/druid/iceberg/input/IcebergInputSourceArrowModeTest.java
new file mode 100644
index 00000000000..ac2af649f6b
--- /dev/null
+++
b/extensions-contrib/druid-iceberg-extensions/src/test/java/org/apache/druid/iceberg/input/IcebergInputSourceArrowModeTest.java
@@ -0,0 +1,427 @@
+/*
+ * 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 com.fasterxml.jackson.databind.ObjectMapper;
+import com.fasterxml.jackson.databind.jsontype.NamedType;
+import com.google.common.collect.ImmutableMap;
+import org.apache.druid.data.input.ColumnsFilter;
+import org.apache.druid.data.input.InputRow;
+import org.apache.druid.data.input.InputRowSchema;
+import org.apache.druid.data.input.InputSource;
+import org.apache.druid.data.input.InputSourceReader;
+import org.apache.druid.data.input.impl.DimensionsSpec;
+import org.apache.druid.data.input.impl.LocalInputSourceFactory;
+import org.apache.druid.data.input.impl.TimestampSpec;
+import org.apache.druid.error.DruidException;
+import org.apache.druid.iceberg.filter.IcebergEqualsFilter;
+import org.apache.druid.iceberg.filter.IcebergFilter;
+import org.apache.druid.jackson.DefaultObjectMapper;
+import org.apache.druid.java.util.common.DateTimes;
+import org.apache.druid.java.util.common.FileUtils;
+import org.apache.druid.java.util.common.parsers.CloseableIterator;
+import org.apache.iceberg.DataFile;
+import org.apache.iceberg.FileFormat;
+import org.apache.iceberg.FileMetadata;
+import org.apache.iceberg.FileScanTask;
+import org.apache.iceberg.Files;
+import org.apache.iceberg.PartitionKey;
+import org.apache.iceberg.PartitionSpec;
+import org.apache.iceberg.Schema;
+import org.apache.iceberg.Table;
+import org.apache.iceberg.catalog.Namespace;
+import org.apache.iceberg.catalog.TableIdentifier;
+import org.apache.iceberg.data.GenericRecord;
+import org.apache.iceberg.data.Record;
+import org.apache.iceberg.data.parquet.GenericParquetWriter;
+import org.apache.iceberg.io.CloseableIterable;
+import org.apache.iceberg.io.DataWriter;
+import org.apache.iceberg.io.OutputFile;
+import org.apache.iceberg.parquet.Parquet;
+import org.apache.iceberg.types.Types;
+import org.joda.time.DateTime;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+
+import java.io.File;
+import java.io.IOException;
+import java.util.ArrayList;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.UUID;
+
+/**
+ * Covers {@link IcebergInputSource} with {@code useArrowReader} enabled,
where the Arrow delegate reads the
+ * table scan directly instead of going through a warehouse input source.
+ */
+public class IcebergInputSourceArrowModeTest
+{
+ private IcebergCatalog testCatalog;
+ private TableIdentifier tableIdentifier;
+ private File warehouseDir;
+
+ private final Schema tableSchema = new Schema(
+ Types.NestedField.required(1, "id", Types.StringType.get()),
+ Types.NestedField.required(2, "name", Types.StringType.get())
+ );
+ private final Map<String, Object> tableData = ImmutableMap.of("id",
"123988", "name", "Foo");
+
+ private static final String NAMESPACE = "default";
+ private static final String TABLENAME = "foosTable";
+
+ @BeforeEach
+ public void setup() throws IOException
+ {
+ warehouseDir = FileUtils.createTempDir();
+ testCatalog = new LocalCatalog(warehouseDir.getPath(), new HashMap<>(),
true);
+ tableIdentifier = TableIdentifier.of(Namespace.of(NAMESPACE), TABLENAME);
+ createAndLoadTable(tableIdentifier);
+ }
+
+ @AfterEach
+ public void tearDown()
+ {
+ dropTableFromCatalog(tableIdentifier);
+ }
+
+ @Test
+ public void testReadWithNullInputStatsDoesNotNpe() throws IOException
+ {
+ final IcebergInputSource src = arrowSource(null, null, null);
+ final InputRowSchema schemaWithMissingTs = new InputRowSchema(
+ new TimestampSpec(null, null, DateTimes.utc(0L)),
+ DimensionsSpec.builder().build(),
+ ColumnsFilter.all()
+ );
+ final InputSourceReader reader = src.reader(schemaWithMissingTs, null,
FileUtils.createTempDir());
+ try (CloseableIterator<InputRow> it = reader.read()) {
+ while (it.hasNext()) {
+ Assertions.assertNotNull(it.next());
+ }
+ }
+ }
+
+ @Test
+ public void testIsNotSplittable()
+ {
+ final IcebergInputSource src = arrowSource(null, null, null);
+ Assertions.assertFalse(src.isSplittable());
+ Assertions.assertFalse(src.needsFormat());
+ }
+
+ @Test
+ public void testSingleSplitReturnsOwningSourceSoRowsAreNotDropped() throws
IOException
+ {
+ final IcebergInputSource src = arrowSource(null, null, null);
+ Assertions.assertEquals(1, src.estimateNumSplits(null, null));
+ final InputSource split = src.createSplits(null, null)
+ .map(src::withSplit)
+ .findFirst()
+ .orElseThrow(() -> new
IllegalStateException("expected one split"));
+ Assertions.assertSame(src, split);
+ }
+
+ @Test
+ public void testResidualFilterModeFail()
+ {
+ final IcebergInputSource src = arrowSource(
+ new IcebergEqualsFilter("id", "123988"),
+ null,
+ ResidualFilterMode.FAIL
+ );
+ final InputRowSchema inputRowSchema = new InputRowSchema(
+ new TimestampSpec("timestamp", "millis", null),
+ DimensionsSpec.builder().build(),
+ ColumnsFilter.all()
+ );
+ final DruidException ex = Assertions.assertThrows(
+ DruidException.class,
+ () -> {
+ final InputSourceReader reader = src.reader(inputRowSchema, null,
FileUtils.createTempDir());
+ reader.read().close();
+ }
+ );
+ Assertions.assertTrue(
+ ex.getMessage().contains("residual"),
+ "Expected residual error: " + ex.getMessage()
+ );
+ }
+
+ @Test
+ public void testResidualFilterModeFailUsesSnapshotTime() throws Exception
+ {
+ final String filterId = (String) tableData.get("id");
+ dropTableFromCatalog(tableIdentifier);
+ final PartitionSpec partitionSpec = PartitionSpec.builderFor(tableSchema)
+ .identity("id")
+ .build();
+ final Table table =
testCatalog.retrieveCatalog().createTable(tableIdentifier, tableSchema,
partitionSpec);
+ appendRow(table, partitionSpec, tableData);
+
+ final long afterPartitionedSnapshot = System.currentTimeMillis();
+ Thread.sleep(10);
+
+ table.updateSpec().removeField("id").commit();
+ appendRow(table, table.spec(), ImmutableMap.of("id", filterId, "name",
"Bar"));
+
+ final IcebergInputSource src = arrowSource(
+ new IcebergEqualsFilter("id", filterId),
+ DateTimes.utc(afterPartitionedSnapshot),
+ ResidualFilterMode.FAIL
+ );
+ final InputRowSchema inputRowSchema = new InputRowSchema(
+ new TimestampSpec(null, null, DateTimes.utc(0L)),
+ DimensionsSpec.builder().build(),
+ ColumnsFilter.all()
+ );
+
+ final InputSourceReader reader = src.reader(inputRowSchema, null,
FileUtils.createTempDir());
+ reader.read().close();
+ }
+
+ @Test
+ public void testArrowReaderRejectsSnapshotWithDeleteFiles() throws
IOException
+ {
+ final Table table = testCatalog.retrieveTable(NAMESPACE, TABLENAME);
+ final String deletePath = warehouseDir.getAbsolutePath() + "/" +
UUID.randomUUID() + ".parquet";
+ final String dataPath;
+ try (CloseableIterable<FileScanTask> tasks = table.newScan().planFiles()) {
+ dataPath = tasks.iterator().next().file().location();
+ }
+ final org.apache.iceberg.DeleteFile deleteFile =
FileMetadata.deleteFileBuilder(table.spec())
+
.ofPositionDeletes()
+
.withPath(deletePath)
+
.withFormat(FileFormat.PARQUET)
+
.withRecordCount(1)
+
.withFileSizeInBytes(1)
+
.withReferencedDataFile(dataPath)
+ .build();
+ final org.apache.iceberg.DeleteFile equalityDeleteFile =
FileMetadata.deleteFileBuilder(table.spec())
+
.ofEqualityDeletes(2)
+
.withPath(
+
warehouseDir.getAbsolutePath()
+ +
"/"
+ +
UUID.randomUUID()
+ +
".parquet"
+ )
+
.withFormat(FileFormat.PARQUET)
+
.withRecordCount(1)
+
.withFileSizeInBytes(1)
+
.build();
+
table.newRowDelta().addDeletes(deleteFile).addDeletes(equalityDeleteFile).commit();
+
+ final IcebergInputSource src = arrowSource(null, null, null);
+ final InputRowSchema inputRowSchema = new InputRowSchema(
+ new TimestampSpec(null, null, DateTimes.utc(0L)),
+ DimensionsSpec.builder().build(),
+ ColumnsFilter.all()
+ );
+
+ final DruidException ex = Assertions.assertThrows(
+ DruidException.class,
+ () -> src.reader(inputRowSchema, null,
FileUtils.createTempDir()).read()
+ );
+ Assertions.assertTrue(ex.getMessage().contains("does not support Iceberg
snapshots with delete files"));
+ }
+
+ @Test
+ public void testArrowReaderIgnoreResidualsReturnsResidualRows() throws
IOException
+ {
+ dropTableFromCatalog(tableIdentifier);
+ final PartitionSpec partitionSpec =
PartitionSpec.builderFor(tableSchema).identity("id").build();
+ final Table table =
testCatalog.retrieveCatalog().createTable(tableIdentifier, tableSchema,
partitionSpec);
+ appendRows(
+ table,
+ partitionSpec,
+ ImmutableMap.of("id", "123988", "name", "Foo"),
+ ImmutableMap.of("id", "123988", "name", "Bar")
+ );
+
+ final IcebergInputSource src = arrowSource(
+ new IcebergEqualsFilter("name", "Foo"),
+ null,
+ ResidualFilterMode.IGNORE
+ );
+ final InputRowSchema inputRowSchema = new InputRowSchema(
+ new TimestampSpec(null, null, DateTimes.utc(0L)),
+ DimensionsSpec.builder().build(),
+ ColumnsFilter.all()
+ );
+
+ final InputSourceReader reader = src.reader(inputRowSchema, null,
FileUtils.createTempDir());
+ final List<InputRow> rows = new ArrayList<>();
+ try (CloseableIterator<InputRow> iterator = reader.read()) {
+ while (iterator.hasNext()) {
+ rows.add(iterator.next());
+ }
+ }
+ Assertions.assertEquals(2, rows.size());
+ Assertions.assertEquals("Foo", rows.get(0).getDimension("name").get(0));
+ Assertions.assertEquals("Bar", rows.get(1).getDimension("name").get(0));
+ }
+
+ @Test
+ public void testWarehouseSourceNotRequiredInArrowMode()
+ {
+ Assertions.assertFalse(arrowSource(null, null, null).isSplittable());
+ }
+
+ @Test
+ public void testWarehouseSourceRequiredWhenArrowDisabled()
+ {
+ Assertions.assertThrows(
+ NullPointerException.class,
+ () -> new IcebergInputSource(
+ TABLENAME,
+ NAMESPACE,
+ null,
+ testCatalog,
+ null,
+ null,
+ null,
+ false,
+ null
+ )
+ );
+ }
+
+ @Test
+ public void testArrowPropertiesSurviveSerde() throws IOException
+ {
+ final ObjectMapper mapper = new DefaultObjectMapper();
+ mapper.registerSubtypes(
+ new NamedType(LocalCatalog.class, LocalCatalog.TYPE_KEY),
+ new NamedType(IcebergInputSource.class, IcebergInputSource.TYPE_KEY)
+ );
+ final IcebergInputSource src = arrowSource(null, null, null);
+ final IcebergInputSource roundTripped = (IcebergInputSource)
mapper.readValue(
+ mapper.writeValueAsBytes(src),
+ InputSource.class
+ );
+ Assertions.assertTrue(roundTripped.isUseArrowReader());
+ Assertions.assertEquals(512, roundTripped.getArrowBatchSize());
+ Assertions.assertFalse(roundTripped.needsFormat());
+ Assertions.assertFalse(roundTripped.isSplittable());
+ }
+
+ @Test
+ public void testDefaultsToStandardModeWhenArrowUnset()
+ {
+ final IcebergInputSource src = new IcebergInputSource(
+ TABLENAME,
+ NAMESPACE,
+ null,
+ testCatalog,
+ new LocalInputSourceFactory(),
+ null,
+ null,
+ null,
+ null
+ );
+ Assertions.assertFalse(src.isUseArrowReader());
+ Assertions.assertTrue(src.needsFormat());
+ Assertions.assertTrue(src.isSplittable());
+ }
+
+ private IcebergInputSource arrowSource(
+ final IcebergFilter icebergFilter,
+ final DateTime snapshotTime,
+ final ResidualFilterMode residualFilterMode
+ )
+ {
+ return new IcebergInputSource(
+ TABLENAME,
+ NAMESPACE,
+ icebergFilter,
+ testCatalog,
+ null,
+ snapshotTime,
+ residualFilterMode,
+ true,
+ 512
+ );
+ }
+
+ private void createAndLoadTable(TableIdentifier id) throws IOException
+ {
+ final Table table = testCatalog.retrieveCatalog().createTable(id,
tableSchema, PartitionSpec.unpartitioned());
+ appendRow(table, PartitionSpec.unpartitioned(), tableData);
+ }
+
+ private void appendRow(Table table, PartitionSpec partitionSpec, Map<String,
Object> rowData) throws IOException
+ {
+ appendRows(table, partitionSpec, rowData);
+ }
+
+ private void appendRows(
+ final Table table,
+ final PartitionSpec partitionSpec,
+ final Map<String, Object>... rows
+ ) throws IOException
+ {
+ final String fname = UUID.randomUUID() + ".parquet";
+ final File dataFile = new File(warehouseDir.getAbsolutePath() + "/" +
fname);
+ Assertions.assertTrue(dataFile.createNewFile());
+ final OutputFile out = Files.localOutput(dataFile);
+ final DataWriter<Record> writer;
+ if (partitionSpec.isUnpartitioned()) {
+ writer = Parquet.writeData(out)
+ .schema(tableSchema)
+ .createWriterFunc(GenericParquetWriter::create)
+ .overwrite()
+ .withSpec(partitionSpec)
+ .build();
+ } else {
+ final PartitionKey partitionKey = new PartitionKey(partitionSpec,
tableSchema);
+ final GenericRecord partitionRow = GenericRecord.create(tableSchema);
+ partitionRow.setField("id", rows[0].get("id"));
+ partitionRow.setField("name", rows[0].get("name"));
+ partitionKey.partition(partitionRow);
+ writer = Parquet.writeData(out)
+ .schema(tableSchema)
+ .createWriterFunc(GenericParquetWriter::create)
+ .overwrite()
+ .withSpec(partitionSpec)
+ .withPartition(partitionKey)
+ .build();
+ }
+ try {
+ for (Map<String, Object> rowData : rows) {
+ final GenericRecord row = GenericRecord.create(tableSchema);
+ row.setField("id", rowData.get("id"));
+ row.setField("name", rowData.get("name"));
+ writer.write(row);
+ }
+ }
+ finally {
+ writer.close();
+ }
+ final DataFile df = writer.toDataFile();
+ table.newAppend().appendFile(df).commit();
+ }
+
+ private void dropTableFromCatalog(TableIdentifier id)
+ {
+ testCatalog.retrieveCatalog().dropTable(id);
+ }
+}
diff --git
a/extensions-contrib/druid-iceberg-extensions/src/test/java/org/apache/druid/iceberg/input/IcebergInputSourceTest.java
b/extensions-contrib/druid-iceberg-extensions/src/test/java/org/apache/druid/iceberg/input/IcebergInputSourceTest.java
index 35d03298e82..a31735b4c30 100644
---
a/extensions-contrib/druid-iceberg-extensions/src/test/java/org/apache/druid/iceberg/input/IcebergInputSourceTest.java
+++
b/extensions-contrib/druid-iceberg-extensions/src/test/java/org/apache/druid/iceberg/input/IcebergInputSourceTest.java
@@ -99,6 +99,8 @@ public class IcebergInputSourceTest
testCatalog,
new LocalInputSourceFactory(),
null,
+ null,
+ null,
null
);
Stream<InputSplit<List<String>>> splits = inputSource.createSplits(null,
new MaxSizeSplitHintSpec(null, null));
@@ -135,6 +137,8 @@ public class IcebergInputSourceTest
testCatalog,
new LocalInputSourceFactory(),
null,
+ null,
+ null,
null
);
Stream<InputSplit<List<String>>> splits = inputSource.createSplits(null,
new MaxSizeSplitHintSpec(null, null));
@@ -151,6 +155,8 @@ public class IcebergInputSourceTest
testCatalog,
new LocalInputSourceFactory(),
null,
+ null,
+ null,
null
);
Stream<InputSplit<List<String>>> splits = inputSource.createSplits(null,
new MaxSizeSplitHintSpec(null, null));
@@ -187,6 +193,8 @@ public class IcebergInputSourceTest
testCatalog,
new LocalInputSourceFactory(),
DateTimes.nowUtc(),
+ null,
+ null,
null
);
Stream<InputSplit<List<String>>> splits = inputSource.createSplits(null,
new MaxSizeSplitHintSpec(null, null));
@@ -207,6 +215,8 @@ public class IcebergInputSourceTest
caseInsensitiveCatalog,
new LocalInputSourceFactory(),
null,
+ null,
+ null,
null
);
@@ -232,7 +242,9 @@ public class IcebergInputSourceTest
testCatalog,
new LocalInputSourceFactory(),
null,
- ResidualFilterMode.IGNORE
+ ResidualFilterMode.IGNORE,
+ null,
+ null
);
Stream<InputSplit<List<String>>> splits = inputSource.createSplits(null,
new MaxSizeSplitHintSpec(null, null));
Assertions.assertEquals(1, splits.count());
@@ -249,7 +261,9 @@ public class IcebergInputSourceTest
testCatalog,
new LocalInputSourceFactory(),
null,
- ResidualFilterMode.FAIL
+ ResidualFilterMode.FAIL,
+ null,
+ null
);
DruidException exception = Assertions.assertThrows(
DruidException.class,
@@ -277,7 +291,9 @@ public class IcebergInputSourceTest
testCatalog,
new LocalInputSourceFactory(),
null,
- ResidualFilterMode.FAIL
+ ResidualFilterMode.FAIL,
+ null,
+ null
);
Stream<InputSplit<List<String>>> splits = inputSource.createSplits(null,
new MaxSizeSplitHintSpec(null, null));
Assertions.assertEquals(1, splits.count());
@@ -300,7 +316,9 @@ public class IcebergInputSourceTest
testCatalog,
new LocalInputSourceFactory(),
null,
- ResidualFilterMode.FAIL
+ ResidualFilterMode.FAIL,
+ null,
+ null
);
DruidException exception = Assertions.assertThrows(
DruidException.class,
diff --git a/licenses.yaml b/licenses.yaml
index bc6fcdb1a72..1862e52d1b2 100644
--- a/licenses.yaml
+++ b/licenses.yaml
@@ -7338,6 +7338,75 @@ notice: |
---
+name: Apache Iceberg Arrow
+license_category: binary
+module: extensions-contrib/druid-iceberg-extensions
+license_name: Apache License version 2.0
+version: 1.11.0
+libraries:
+ - org.apache.iceberg: iceberg-arrow
+
+---
+
+name: Apache Arrow
+license_category: binary
+module: extensions-contrib/druid-iceberg-extensions
+license_name: Apache License version 2.0
+version: 15.0.2
+libraries:
+ - org.apache.arrow: arrow-format
+ - org.apache.arrow: arrow-memory-core
+ - org.apache.arrow: arrow-memory-unsafe
+ - org.apache.arrow: arrow-vector
+notices:
+ - arrow-format: |
+ Arrow Format
+ Copyright 2024 The Apache Software Foundation
+
+ This product includes software developed at
+ The Apache Software Foundation (http://www.apache.org/).
+ - arrow-memory-core: |
+ Arrow Memory - Core
+ Copyright 2024 The Apache Software Foundation
+
+ This product includes software developed at
+ The Apache Software Foundation (http://www.apache.org/).
+ - arrow-memory-unsafe: |
+ Arrow Memory - Unsafe
+ Copyright 2024 The Apache Software Foundation
+
+ This product includes software developed at
+ The Apache Software Foundation (http://www.apache.org/).
+ - arrow-vector: |
+ Arrow Vectors
+ Copyright 2024 The Apache Software Foundation
+
+ This product includes software developed at
+ The Apache Software Foundation (http://www.apache.org/).
+
+---
+
+name: FlatBuffers Java API
+license_category: binary
+module: extensions-contrib/druid-iceberg-extensions
+license_name: Apache License version 2.0
+version: 23.5.26
+libraries:
+ - com.google.flatbuffers: flatbuffers-java
+
+---
+
+name: Eclipse Collections
+license_category: binary
+module: extensions-contrib/druid-iceberg-extensions
+license_name: Eclipse Public License 1.0
+version: 11.1.0
+libraries:
+ - org.eclipse.collections: eclipse-collections
+ - org.eclipse.collections: eclipse-collections-api
+
+---
+
name: Kafka Schema Registry Client
version: 8.3.1
license_category: binary
diff --git a/pom.xml b/pom.xml
index 0fea8db9cc1..e6b00405abc 100644
--- a/pom.xml
+++ b/pom.xml
@@ -104,6 +104,7 @@
<guice.version>6.0.0</guice.version>
<hamcrest.version>3.0</hamcrest.version>
<iceberg.core.version>1.11.0</iceberg.core.version>
+ <arrow.version>15.0.2</arrow.version>
<jetty.version>12.1.13</jetty.version>
<jspecify.version>1.0.1</jspecify.version>
<jersey.version>1.19.4</jersey.version>
@@ -429,6 +430,22 @@
<version>3.8.7</version>
</dependency>
+ <!-- Arrow + Iceberg-Arrow -->
+ <dependency>
+ <groupId>org.apache.arrow</groupId>
+ <artifactId>arrow-vector</artifactId>
+ <version>${arrow.version}</version>
+ </dependency>
+ <dependency>
+ <groupId>org.apache.arrow</groupId>
+ <artifactId>arrow-memory-unsafe</artifactId>
+ <version>${arrow.version}</version>
+ </dependency>
+ <dependency>
+ <groupId>org.apache.iceberg</groupId>
+ <artifactId>iceberg-arrow</artifactId>
+ <version>${iceberg.core.version}</version>
+ </dependency>
<!-- transitive dependency of calcite-core and avatica-core, which
require different
versions; pinned to the higher of the two -->
<dependency>
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]