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]

Reply via email to