From 1358bd6b337984e4b112cfa50cc92a035f6b76cd Mon Sep 17 00:00:00 2001 From: rockyyin Date: Fri, 17 Jul 2026 00:05:33 +0800 Subject: [PATCH] feat: bump lance to 7.0.0, arrow to 18.3.0, and migrate to org.lance API MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Upgrade lance-core from 0.23.3 to 7.0.0 - Upgrade Apache Arrow from 14.0.0 to 18.3.0 - Upgrade JUnit from 5.9.3 to 5.10.1 - Migrate groupId from com.lancedb to org.lance - Add slf4j-api exclusion for lance-core - Bump Java compiler source/target from 1.8 to 11 (required by lance 7.0.0) - Update Fragment API: Fragment.create() → Fragment.write().execute() - Update commit API: FragmentOperation → Transaction + CommitBuilder - Update IndexParams: Builder() → builder() static factory - Refresh import paths across all source files --- .gitignore | 14 +++++- pom.xml | 18 ++++--- .../connector/lance/LanceAggregateSource.java | 8 ++-- .../connector/lance/LanceIndexBuilder.java | 25 +++++----- .../connector/lance/LanceInputFormat.java | 8 ++-- .../flink/connector/lance/LanceSink.java | 48 ++++++++++++------- .../flink/connector/lance/LanceSource.java | 8 ++-- .../connector/lance/LanceVectorSearch.java | 10 ++-- .../connector/lance/table/LanceCatalog.java | 2 +- 9 files changed, 84 insertions(+), 57 deletions(-) diff --git a/.gitignore b/.gitignore index c0c6c6b..cc2a8eb 100644 --- a/.gitignore +++ b/.gitignore @@ -61,4 +61,16 @@ build/ .agents/ ### Lance test data (generated at runtime) ### -test-data/ \ No newline at end of file +test-data/ + +### Local planning & research (not for upstream PR) ### +.research/ +task_plan.md +findings.md +progress.md + +### Local planning & research (not for upstream PR) ### +.research/ +task_plan.md +findings.md +progress.md \ No newline at end of file diff --git a/pom.xml b/pom.xml index 5547152..8cb4161 100644 --- a/pom.xml +++ b/pom.xml @@ -15,14 +15,14 @@ UTF-8 UTF-8 - 1.8 - 1.8 + 11 + 11 1.16.1 - 0.23.3 - 14.0.0 - 5.9.3 + 7.0.0 + 18.3.0 + 5.10.1 1.7.36 2.20.0 5.3.1 @@ -73,9 +73,15 @@ - com.lancedb + org.lance lance-core ${lance.version} + + + org.slf4j + slf4j-api + + diff --git a/src/main/java/org/apache/flink/connector/lance/LanceAggregateSource.java b/src/main/java/org/apache/flink/connector/lance/LanceAggregateSource.java index ed512e5..9b37f6f 100644 --- a/src/main/java/org/apache/flink/connector/lance/LanceAggregateSource.java +++ b/src/main/java/org/apache/flink/connector/lance/LanceAggregateSource.java @@ -28,10 +28,10 @@ import org.apache.flink.table.data.RowData; import org.apache.flink.table.types.logical.RowType; -import com.lancedb.lance.Dataset; -import com.lancedb.lance.Fragment; -import com.lancedb.lance.ipc.LanceScanner; -import com.lancedb.lance.ipc.ScanOptions; +import org.lance.Dataset; +import org.lance.Fragment; +import org.lance.ipc.LanceScanner; +import org.lance.ipc.ScanOptions; import org.apache.arrow.memory.BufferAllocator; import org.apache.arrow.memory.RootAllocator; import org.apache.arrow.vector.VectorSchemaRoot; diff --git a/src/main/java/org/apache/flink/connector/lance/LanceIndexBuilder.java b/src/main/java/org/apache/flink/connector/lance/LanceIndexBuilder.java index 3f93897..2c5639b 100644 --- a/src/main/java/org/apache/flink/connector/lance/LanceIndexBuilder.java +++ b/src/main/java/org/apache/flink/connector/lance/LanceIndexBuilder.java @@ -20,14 +20,14 @@ import org.apache.flink.connector.lance.config.LanceOptions; -import com.lancedb.lance.Dataset; -import com.lancedb.lance.index.DistanceType; -import com.lancedb.lance.index.IndexParams; -import com.lancedb.lance.index.IndexType; -import com.lancedb.lance.index.vector.HnswBuildParams; -import com.lancedb.lance.index.vector.IvfBuildParams; -import com.lancedb.lance.index.vector.PQBuildParams; -import com.lancedb.lance.index.vector.VectorIndexParams; +import org.lance.Dataset; +import org.lance.index.DistanceType; +import org.lance.index.IndexParams; +import org.lance.index.IndexType; +import org.lance.index.vector.HnswBuildParams; +import org.lance.index.vector.IvfBuildParams; +import org.lance.index.vector.PQBuildParams; +import org.lance.index.vector.VectorIndexParams; import org.apache.arrow.memory.BufferAllocator; import org.apache.arrow.memory.RootAllocator; import org.slf4j.Logger; @@ -131,8 +131,7 @@ public IndexBuildResult buildIndex() throws IOException { .build(); VectorIndexParams ivfPqParams = VectorIndexParams.withIvfPqParams( distanceType, ivfParams, pqParams); - indexParams = new IndexParams.Builder() - .setDistanceType(distanceType) + indexParams = IndexParams.builder() .setVectorIndexParams(ivfPqParams) .build(); break; @@ -150,8 +149,7 @@ public IndexBuildResult buildIndex() throws IOException { .build(); VectorIndexParams ivfHnswParams = VectorIndexParams.withIvfHnswPqParams( distanceType, ivfParams, hnswParams, hnswPqParams); - indexParams = new IndexParams.Builder() - .setDistanceType(distanceType) + indexParams = IndexParams.builder() .setVectorIndexParams(ivfHnswParams) .build(); break; @@ -159,8 +157,7 @@ public IndexBuildResult buildIndex() throws IOException { case IVF_FLAT: lanceIndexType = IndexType.IVF_FLAT; VectorIndexParams ivfFlatParams = VectorIndexParams.ivfFlat(numPartitions, distanceType); - indexParams = new IndexParams.Builder() - .setDistanceType(distanceType) + indexParams = IndexParams.builder() .setVectorIndexParams(ivfFlatParams) .build(); break; diff --git a/src/main/java/org/apache/flink/connector/lance/LanceInputFormat.java b/src/main/java/org/apache/flink/connector/lance/LanceInputFormat.java index 2785a00..5638680 100644 --- a/src/main/java/org/apache/flink/connector/lance/LanceInputFormat.java +++ b/src/main/java/org/apache/flink/connector/lance/LanceInputFormat.java @@ -28,10 +28,10 @@ import org.apache.flink.table.data.RowData; import org.apache.flink.table.types.logical.RowType; -import com.lancedb.lance.Dataset; -import com.lancedb.lance.Fragment; -import com.lancedb.lance.ipc.LanceScanner; -import com.lancedb.lance.ipc.ScanOptions; +import org.lance.Dataset; +import org.lance.Fragment; +import org.lance.ipc.LanceScanner; +import org.lance.ipc.ScanOptions; import org.apache.arrow.memory.BufferAllocator; import org.apache.arrow.memory.RootAllocator; import org.apache.arrow.vector.VectorSchemaRoot; diff --git a/src/main/java/org/apache/flink/connector/lance/LanceSink.java b/src/main/java/org/apache/flink/connector/lance/LanceSink.java index feeec90..84e3069 100644 --- a/src/main/java/org/apache/flink/connector/lance/LanceSink.java +++ b/src/main/java/org/apache/flink/connector/lance/LanceSink.java @@ -29,11 +29,14 @@ import org.apache.flink.table.data.RowData; import org.apache.flink.table.types.logical.RowType; -import com.lancedb.lance.Dataset; -import com.lancedb.lance.Fragment; -import com.lancedb.lance.FragmentMetadata; -import com.lancedb.lance.FragmentOperation; -import com.lancedb.lance.WriteParams; +import org.lance.Dataset; +import org.lance.Fragment; +import org.lance.FragmentMetadata; +import org.lance.WriteParams; +import org.lance.CommitBuilder; +import org.lance.Transaction; +import org.lance.operation.Append; +import org.lance.operation.Overwrite; import org.apache.arrow.memory.BufferAllocator; import org.apache.arrow.memory.RootAllocator; import org.apache.arrow.vector.VectorSchemaRoot; @@ -48,7 +51,6 @@ import java.util.ArrayList; import java.util.Collections; import java.util.List; -import java.util.Optional; /** * Lance Sink implementation. @@ -161,17 +163,21 @@ public void flush() throws IOException { .build(); // Create Fragment - List fragments = Fragment.create( - datasetPath, - allocator, - root, - writeParams - ); + List fragments = Fragment.write() + .datasetUri(datasetPath) + .allocator(allocator) + .data(root) + .writeParams(writeParams) + .execute(); if (!datasetExists) { // Create new dataset (using Overwrite operation) - FragmentOperation.Overwrite overwrite = new FragmentOperation.Overwrite(fragments, arrowSchema); - dataset = overwrite.commit(allocator, datasetPath, Optional.empty(), Collections.emptyMap()); + Overwrite operation = Overwrite.builder().fragments(fragments).schema(arrowSchema).build(); + final CommitBuilder builder = + new CommitBuilder(datasetPath, allocator).writeParams(Collections.emptyMap()); + try (Transaction txn = new Transaction.Builder().operation(operation).build()) { + dataset = builder.execute(txn); + } datasetExists = true; isFirstWrite = false; LOG.info("Created new dataset: {}", datasetPath); @@ -179,13 +185,19 @@ public void flush() throws IOException { // Append data if (isFirstWrite && options.getWriteMode() == LanceOptions.WriteMode.OVERWRITE) { // First write and overwrite mode - FragmentOperation.Overwrite overwrite = new FragmentOperation.Overwrite(fragments, arrowSchema); - dataset = overwrite.commit(allocator, datasetPath, Optional.empty(), Collections.emptyMap()); + Overwrite operation = Overwrite.builder().fragments(fragments).schema(arrowSchema).build(); + final CommitBuilder builder = new CommitBuilder(datasetPath, allocator); + try (Transaction txn = new Transaction.Builder().operation(operation).build()) { + dataset = builder.execute(txn); + } isFirstWrite = false; } else { // Append mode - FragmentOperation.Append append = new FragmentOperation.Append(fragments); - dataset = append.commit(allocator, datasetPath, Optional.empty(), Collections.emptyMap()); + Append operation = Append.builder().fragments(fragments).build(); + final CommitBuilder builder = new CommitBuilder(datasetPath, allocator); + try (Transaction txn = new Transaction.Builder().operation(operation).build()) { + dataset = builder.execute(txn); + } } } diff --git a/src/main/java/org/apache/flink/connector/lance/LanceSource.java b/src/main/java/org/apache/flink/connector/lance/LanceSource.java index ade00a0..2292b2b 100644 --- a/src/main/java/org/apache/flink/connector/lance/LanceSource.java +++ b/src/main/java/org/apache/flink/connector/lance/LanceSource.java @@ -27,10 +27,10 @@ import org.apache.flink.table.data.RowData; import org.apache.flink.table.types.logical.RowType; -import com.lancedb.lance.Dataset; -import com.lancedb.lance.Fragment; -import com.lancedb.lance.ipc.LanceScanner; -import com.lancedb.lance.ipc.ScanOptions; +import org.lance.Dataset; +import org.lance.Fragment; +import org.lance.ipc.LanceScanner; +import org.lance.ipc.ScanOptions; import org.apache.arrow.memory.BufferAllocator; import org.apache.arrow.memory.RootAllocator; import org.apache.arrow.vector.VectorSchemaRoot; diff --git a/src/main/java/org/apache/flink/connector/lance/LanceVectorSearch.java b/src/main/java/org/apache/flink/connector/lance/LanceVectorSearch.java index faf0c5a..d51477f 100644 --- a/src/main/java/org/apache/flink/connector/lance/LanceVectorSearch.java +++ b/src/main/java/org/apache/flink/connector/lance/LanceVectorSearch.java @@ -26,11 +26,11 @@ import org.apache.flink.table.data.RowData; import org.apache.flink.table.types.logical.RowType; -import com.lancedb.lance.Dataset; -import com.lancedb.lance.index.DistanceType; -import com.lancedb.lance.ipc.LanceScanner; -import com.lancedb.lance.ipc.Query; -import com.lancedb.lance.ipc.ScanOptions; +import org.lance.Dataset; +import org.lance.index.DistanceType; +import org.lance.ipc.LanceScanner; +import org.lance.ipc.Query; +import org.lance.ipc.ScanOptions; import org.apache.arrow.memory.BufferAllocator; import org.apache.arrow.memory.RootAllocator; import org.apache.arrow.vector.Float8Vector; diff --git a/src/main/java/org/apache/flink/connector/lance/table/LanceCatalog.java b/src/main/java/org/apache/flink/connector/lance/table/LanceCatalog.java index 74fed94..600c60a 100644 --- a/src/main/java/org/apache/flink/connector/lance/table/LanceCatalog.java +++ b/src/main/java/org/apache/flink/connector/lance/table/LanceCatalog.java @@ -49,7 +49,7 @@ import org.apache.flink.table.types.DataType; import org.apache.flink.table.types.logical.RowType; -import com.lancedb.lance.Dataset; +import org.lance.Dataset; import org.apache.arrow.memory.BufferAllocator; import org.apache.arrow.memory.RootAllocator; import org.slf4j.Logger;