This is an automated email from the ASF dual-hosted git repository.
FrankChen021 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 e35ff894868 build(deps): bump delta-kernel.version from 3.2.1 to 4.3.1
(#20012)
e35ff894868 is described below
commit e35ff894868f7d4082ff921b93ca1d4d723d4b68
Author: dependabot[bot] <49699333+dependabot[bot]@users.noreply.github.com>
AuthorDate: Sun Aug 16 13:07:09 2026 +0800
build(deps): bump delta-kernel.version from 3.2.1 to 4.3.1 (#20012)
* build(deps): bump delta-kernel.version from 3.2.1 to 4.3.1
Bumps `delta-kernel.version` from 3.2.1 to 4.3.1.
Updates `io.delta:delta-kernel-api` from 3.2.1 to 4.3.1
- [Release notes](https://github.com/delta-io/delta/releases)
- [Commits](https://github.com/delta-io/delta/compare/v3.2.1...v4.3.1)
Updates `io.delta:delta-kernel-defaults` from 3.2.1 to 4.3.1
- [Release notes](https://github.com/delta-io/delta/releases)
- [Commits](https://github.com/delta-io/delta/compare/v3.2.1...v4.3.1)
---
updated-dependencies:
- dependency-name: io.delta:delta-kernel-api
dependency-version: 4.3.1
dependency-type: direct:production
update-type: version-update:semver-major
- dependency-name: io.delta:delta-kernel-defaults
dependency-version: 4.3.1
dependency-type: direct:production
update-type: version-update:semver-major
...
Signed-off-by: dependabot[bot] <[email protected]>
* fix(deltalake): adapt extension to Delta Kernel 4.3.1
---------
Signed-off-by: dependabot[bot] <[email protected]>
Co-authored-by: dependabot[bot]
<49699333+dependabot[bot]@users.noreply.github.com>
Co-authored-by: Frank Chen <[email protected]>
---
.../druid-deltalake-extensions/pom.xml | 2 +-
.../apache/druid/delta/input/DeltaInputSource.java | 33 ++++++++++++----------
.../org/apache/druid/delta/input/RowSerde.java | 3 +-
.../druid/delta/input/DeltaInputRowTest.java | 15 ++++++----
.../apache/druid/delta/input/DeltaTestUtils.java | 6 ++--
5 files changed, 33 insertions(+), 26 deletions(-)
diff --git a/extensions-contrib/druid-deltalake-extensions/pom.xml
b/extensions-contrib/druid-deltalake-extensions/pom.xml
index 3149728adfb..617a2eea543 100644
--- a/extensions-contrib/druid-deltalake-extensions/pom.xml
+++ b/extensions-contrib/druid-deltalake-extensions/pom.xml
@@ -35,7 +35,7 @@
<modelVersion>4.0.0</modelVersion>
<properties>
- <delta-kernel.version>3.2.1</delta-kernel.version>
+ <delta-kernel.version>4.3.1</delta-kernel.version>
</properties>
<dependencies>
diff --git
a/extensions-contrib/druid-deltalake-extensions/src/main/java/org/apache/druid/delta/input/DeltaInputSource.java
b/extensions-contrib/druid-deltalake-extensions/src/main/java/org/apache/druid/delta/input/DeltaInputSource.java
index 4f255d020f7..44dc34ca3d2 100644
---
a/extensions-contrib/druid-deltalake-extensions/src/main/java/org/apache/druid/delta/input/DeltaInputSource.java
+++
b/extensions-contrib/druid-deltalake-extensions/src/main/java/org/apache/druid/delta/input/DeltaInputSource.java
@@ -33,6 +33,7 @@ import io.delta.kernel.data.FilteredColumnarBatch;
import io.delta.kernel.data.Row;
import io.delta.kernel.defaults.engine.DefaultEngine;
import io.delta.kernel.engine.Engine;
+import io.delta.kernel.engine.FileReadResult;
import io.delta.kernel.exceptions.TableNotFoundException;
import io.delta.kernel.expressions.Predicate;
import io.delta.kernel.internal.InternalScanFileUtils;
@@ -145,7 +146,7 @@ public class DeltaInputSource implements
SplittableInputSource<DeltaSplit>
if (deltaSplit != null) {
final Row scanState = deserialize(engine, deltaSplit.getStateRow());
final StructType physicalReadSchema =
- ScanStateRow.getPhysicalDataReadSchema(engine, scanState);
+ ScanStateRow.getPhysicalDataReadSchema(scanState);
for (String file : deltaSplit.getFiles()) {
final Row scanFile = deserialize(engine, file);
@@ -157,22 +158,22 @@ public class DeltaInputSource implements
SplittableInputSource<DeltaSplit>
final Table table = Table.forPath(engine, tablePath);
final Snapshot snapshot = getSnapshotForTable(table, engine);
- final StructType fullSnapshotSchema = snapshot.getSchema(engine);
+ final StructType fullSnapshotSchema = snapshot.getSchema();
final StructType prunedSchema = pruneSchema(
fullSnapshotSchema,
inputRowSchema.getColumnsFilter()
);
- final ScanBuilder scanBuilder = snapshot.getScanBuilder(engine);
+ final ScanBuilder scanBuilder = snapshot.getScanBuilder();
if (filter != null) {
- scanBuilder.withFilter(engine,
filter.getFilterPredicate(fullSnapshotSchema));
+
scanBuilder.withFilter(filter.getFilterPredicate(fullSnapshotSchema));
}
- final Scan scan = scanBuilder.withReadSchema(engine,
prunedSchema).build();
+ final Scan scan = scanBuilder.withReadSchema(prunedSchema).build();
final CloseableIterator<FilteredColumnarBatch> scanFilesIter =
scan.getScanFiles(engine);
final Row scanState = scan.getScanState(engine);
final StructType physicalReadSchema =
- ScanStateRow.getPhysicalDataReadSchema(engine, scanState);
+ ScanStateRow.getPhysicalDataReadSchema(scanState);
while (scanFilesIter.hasNext()) {
final FilteredColumnarBatch scanFileBatch = scanFilesIter.next();
@@ -217,13 +218,13 @@ public class DeltaInputSource implements
SplittableInputSource<DeltaSplit>
catch (TableNotFoundException e) {
throw InvalidInput.exception(e, "tablePath[%s] not found.", tablePath);
}
- final StructType fullSnapshotSchema = snapshot.getSchema(engine);
+ final StructType fullSnapshotSchema = snapshot.getSchema();
- final ScanBuilder scanBuilder = snapshot.getScanBuilder(engine);
+ final ScanBuilder scanBuilder = snapshot.getScanBuilder();
if (filter != null) {
- scanBuilder.withFilter(engine,
filter.getFilterPredicate(fullSnapshotSchema));
+ scanBuilder.withFilter(filter.getFilterPredicate(fullSnapshotSchema));
}
- final Scan scan = scanBuilder.withReadSchema(engine,
fullSnapshotSchema).build();
+ final Scan scan = scanBuilder.withReadSchema(fullSnapshotSchema).build();
// scan files iterator for the current snapshot
final CloseableIterator<FilteredColumnarBatch> scanFilesIterator =
scan.getScanFiles(engine);
@@ -323,11 +324,13 @@ public class DeltaInputSource implements
SplittableInputSource<DeltaSplit>
{
final FileStatus fileStatus =
InternalScanFileUtils.getAddFileStatus(scanFile);
- final CloseableIterator<ColumnarBatch> physicalDataIter =
engine.getParquetHandler().readParquetFiles(
- Utils.singletonCloseableIterator(fileStatus),
- physicalReadSchema,
- optionalPredicate
- );
+ final CloseableIterator<ColumnarBatch> physicalDataIter =
engine.getParquetHandler()
+ .readParquetFiles(
+ Utils.singletonCloseableIterator(fileStatus),
+ physicalReadSchema,
+ optionalPredicate
+ )
+ .map(FileReadResult::getData);
return Scan.transformPhysicalData(
engine,
diff --git
a/extensions-contrib/druid-deltalake-extensions/src/main/java/org/apache/druid/delta/input/RowSerde.java
b/extensions-contrib/druid-deltalake-extensions/src/main/java/org/apache/druid/delta/input/RowSerde.java
index d7c6fcccdba..0b0c62abef8 100644
---
a/extensions-contrib/druid-deltalake-extensions/src/main/java/org/apache/druid/delta/input/RowSerde.java
+++
b/extensions-contrib/druid-deltalake-extensions/src/main/java/org/apache/druid/delta/input/RowSerde.java
@@ -26,6 +26,7 @@ import com.fasterxml.jackson.databind.node.ObjectNode;
import io.delta.kernel.data.Row;
import io.delta.kernel.defaults.internal.data.DefaultJsonRow;
import io.delta.kernel.engine.Engine;
+import io.delta.kernel.internal.types.DataTypeJsonSerDe;
import io.delta.kernel.internal.util.VectorUtils;
import io.delta.kernel.types.ArrayType;
import io.delta.kernel.types.BooleanType;
@@ -90,7 +91,7 @@ public class RowSerde
try {
JsonNode jsonNode = OBJECT_MAPPER.readTree(jsonRowWithSchema);
JsonNode schemaNode = jsonNode.get("schema");
- StructType schema =
engine.getJsonHandler().deserializeStructType(schemaNode.asText());
+ StructType schema =
DataTypeJsonSerDe.deserializeStructType(schemaNode.asText());
return parseRowFromJsonWithSchema((ObjectNode) jsonNode.get("row"),
schema);
}
catch (JsonProcessingException e) {
diff --git
a/extensions-contrib/druid-deltalake-extensions/src/test/java/org/apache/druid/delta/input/DeltaInputRowTest.java
b/extensions-contrib/druid-deltalake-extensions/src/test/java/org/apache/druid/delta/input/DeltaInputRowTest.java
index 9e270bdab10..3acd795e31f 100644
---
a/extensions-contrib/druid-deltalake-extensions/src/test/java/org/apache/druid/delta/input/DeltaInputRowTest.java
+++
b/extensions-contrib/druid-deltalake-extensions/src/test/java/org/apache/druid/delta/input/DeltaInputRowTest.java
@@ -25,6 +25,7 @@ import io.delta.kernel.data.FilteredColumnarBatch;
import io.delta.kernel.data.Row;
import io.delta.kernel.defaults.engine.DefaultEngine;
import io.delta.kernel.engine.Engine;
+import io.delta.kernel.engine.FileReadResult;
import io.delta.kernel.exceptions.TableNotFoundException;
import io.delta.kernel.internal.InternalScanFileUtils;
import io.delta.kernel.internal.data.ScanStateRow;
@@ -73,7 +74,7 @@ public class DeltaInputRowTest
final Scan scan = DeltaTestUtils.getScan(engine, deltaTablePath);
final Row scanState = scan.getScanState(engine);
- final StructType physicalReadSchema =
ScanStateRow.getPhysicalDataReadSchema(engine, scanState);
+ final StructType physicalReadSchema =
ScanStateRow.getPhysicalDataReadSchema(scanState);
final CloseableIterator<FilteredColumnarBatch> scanFileIter =
scan.getScanFiles(engine);
int totalRecordCount = 0;
@@ -85,11 +86,13 @@ public class DeltaInputRowTest
final Row scanFile = scanFileRows.next();
final FileStatus fileStatus =
InternalScanFileUtils.getAddFileStatus(scanFile);
- final CloseableIterator<ColumnarBatch> physicalDataIter =
engine.getParquetHandler().readParquetFiles(
- Utils.singletonCloseableIterator(fileStatus),
- physicalReadSchema,
- Optional.empty()
- );
+ final CloseableIterator<ColumnarBatch> physicalDataIter =
engine.getParquetHandler()
+ .readParquetFiles(
+ Utils.singletonCloseableIterator(fileStatus),
+ physicalReadSchema,
+ Optional.empty()
+ )
+ .map(FileReadResult::getData);
final CloseableIterator<FilteredColumnarBatch> dataIter =
Scan.transformPhysicalData(
engine,
scanState,
diff --git
a/extensions-contrib/druid-deltalake-extensions/src/test/java/org/apache/druid/delta/input/DeltaTestUtils.java
b/extensions-contrib/druid-deltalake-extensions/src/test/java/org/apache/druid/delta/input/DeltaTestUtils.java
index 6d49428ce02..e7154360109 100644
---
a/extensions-contrib/druid-deltalake-extensions/src/test/java/org/apache/druid/delta/input/DeltaTestUtils.java
+++
b/extensions-contrib/druid-deltalake-extensions/src/test/java/org/apache/druid/delta/input/DeltaTestUtils.java
@@ -33,9 +33,9 @@ public class DeltaTestUtils
{
final Table table = Table.forPath(engine, deltaTablePath);
final Snapshot snapshot = table.getLatestSnapshot(engine);
- final StructType readSchema = snapshot.getSchema(engine);
- final ScanBuilder scanBuilder = snapshot.getScanBuilder(engine)
- .withReadSchema(engine,
readSchema);
+ final StructType readSchema = snapshot.getSchema();
+ final ScanBuilder scanBuilder = snapshot.getScanBuilder()
+ .withReadSchema(readSchema);
return scanBuilder.build();
}
}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]