Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
14 changes: 13 additions & 1 deletion .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -61,4 +61,16 @@ build/
.agents/

### Lance test data (generated at runtime) ###
test-data/
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
18 changes: 12 additions & 6 deletions pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -15,14 +15,14 @@
<properties>
<project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
<project.reporting.outputEncoding>UTF-8</project.reporting.outputEncoding>
<maven.compiler.source>1.8</maven.compiler.source>
<maven.compiler.target>1.8</maven.compiler.target>
<maven.compiler.source>11</maven.compiler.source>
<maven.compiler.target>11</maven.compiler.target>

<!-- 版本管理 -->
<flink.version>1.16.1</flink.version>
<lance.version>0.23.3</lance.version>
<arrow.version>14.0.0</arrow.version>
<junit.version>5.9.3</junit.version>
<lance.version>7.0.0</lance.version>
<arrow.version>18.3.0</arrow.version>
<junit.version>5.10.1</junit.version>
<slf4j.version>1.7.36</slf4j.version>
<log4j.version>2.20.0</log4j.version>
<mockito.version>5.3.1</mockito.version>
Expand Down Expand Up @@ -73,9 +73,15 @@

<!-- ==================== Lance Java SDK ==================== -->
<dependency>
<groupId>com.lancedb</groupId>
<groupId>org.lance</groupId>
<artifactId>lance-core</artifactId>
<version>${lance.version}</version>
<exclusions>
<exclusion>
<groupId>org.slf4j</groupId>
<artifactId>slf4j-api</artifactId>
</exclusion>
</exclusions>
</dependency>

<!-- ==================== Apache Arrow ==================== -->
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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;
Expand All @@ -150,17 +149,15 @@ 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;

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;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
48 changes: 30 additions & 18 deletions src/main/java/org/apache/flink/connector/lance/LanceSink.java
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -48,7 +51,6 @@
import java.util.ArrayList;
import java.util.Collections;
import java.util.List;
import java.util.Optional;

/**
* Lance Sink implementation.
Expand Down Expand Up @@ -161,31 +163,41 @@ public void flush() throws IOException {
.build();

// Create Fragment
List<FragmentMetadata> fragments = Fragment.create(
datasetPath,
allocator,
root,
writeParams
);
List<FragmentMetadata> 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);
} else {
// 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);
}
}
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down