Skip to content
Open
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
2 changes: 1 addition & 1 deletion java/lance-jni/src/optimize.rs
Original file line number Diff line number Diff line change
Expand Up @@ -122,7 +122,7 @@ fn inner_plan_compaction<'local>(
}

#[unsafe(no_mangle)]
pub extern "system" fn Java_org_lance_compaction_Compaction_nativeCommitCompaction<'local>(
pub extern "system" fn Java_org_lance_compaction_Compaction_commitCompactionNative<'local>(
mut env: JNIEnv<'local>,
_obj: JObject,
java_dataset: JObject, // Dataset
Expand Down
36 changes: 19 additions & 17 deletions java/src/main/java/org/lance/CommitBuilder.java
Original file line number Diff line number Diff line change
Expand Up @@ -279,23 +279,25 @@ public CommitBuilder commitTimeout(Duration timeout) {
public Dataset execute(Transaction transaction) {
Preconditions.checkNotNull(transaction, "Transaction must not be null");
if (dataset != null) {
Dataset result =
nativeCommitToDataset(
dataset,
transaction,
detached,
enableV2ManifestPaths,
writeParams,
useStableRowIds,
storageFormat,
maxRetries,
skipAutoCleanup,
namespaceClient,
tableId,
namespaceClientManagedVersioning,
commitTimeoutNanos);
result.setAllocator(dataset.allocator());
return result;
try (LockManager.ReadLock readLock = dataset.acquireReadLock()) {
Dataset result =
nativeCommitToDataset(
dataset,
transaction,
detached,
enableV2ManifestPaths,
writeParams,
useStableRowIds,
storageFormat,
maxRetries,
skipAutoCleanup,
namespaceClient,
tableId,
namespaceClientManagedVersioning,
commitTimeoutNanos);
result.setAllocator(dataset.allocator());
return result;
}
}
if (uri != null) {
Dataset result =
Expand Down
20 changes: 20 additions & 0 deletions java/src/main/java/org/lance/Dataset.java
Original file line number Diff line number Diff line change
Expand Up @@ -1736,6 +1736,26 @@ private void updateToNewDataset(Dataset newDataset) {
newDataset.nativeDatasetHandle = 0;
}

/**
* Acquires a shared read lock that pins the native dataset handle, blocking a concurrent {@link
* #close()} until the lock is released.
*
* <p>Any code that passes this {@link Dataset} into a native method must hold this lock for the
* whole native call; otherwise {@code close()} can release the native dataset mid-call and crash
* the JVM. The lock is reentrant and intended for try-with-resources use.
*
* @return the acquired read lock
* @throws IllegalArgumentException if the dataset is already closed
*/
public LockManager.ReadLock acquireReadLock() {
LockManager.ReadLock readLock = lockManager.acquireReadLock();
if (nativeDatasetHandle == 0) {
readLock.close();
throw new IllegalArgumentException("Dataset is closed");
}
return readLock;
}

/**
* Closes this dataset and releases any system resources associated with it. If the dataset is
* already closed, then invoking this method has no effect.
Expand Down
20 changes: 14 additions & 6 deletions java/src/main/java/org/lance/Fragment.java
Original file line number Diff line number Diff line change
Expand Up @@ -113,7 +113,9 @@ public LanceScanner newScan(ScanOptions options) {
* returns a new fragment with the updated deletion vector.
*/
public FragmentMetadata deleteRows(List<Integer> rowIndexes) {
return nativeDeleteRows(dataset, fragmentMetadata.getId(), rowIndexes);
try (LockManager.ReadLock readLock = dataset.acquireReadLock()) {
return nativeDeleteRows(dataset, fragmentMetadata.getId(), rowIndexes);
}
}

private static native FragmentMetadata nativeDeleteRows(
Expand All @@ -129,7 +131,9 @@ public int getId() {
* @return row counts in this Fragment
*/
public int countRows() {
return countRowsNative(dataset, fragmentMetadata.getId());
try (LockManager.ReadLock readLock = dataset.acquireReadLock()) {
return countRowsNative(dataset, fragmentMetadata.getId());
}
}

/**
Expand All @@ -153,8 +157,10 @@ public int countRows() {
* @return the fragment metadata and new schema.
*/
public FragmentMergeResult mergeColumns(ArrowArrayStream stream, String leftOn, String rightOn) {
return nativeMergeColumns(
dataset, fragmentMetadata.getId(), stream.memoryAddress(), leftOn, rightOn);
try (LockManager.ReadLock readLock = dataset.acquireReadLock()) {
return nativeMergeColumns(
dataset, fragmentMetadata.getId(), stream.memoryAddress(), leftOn, rightOn);
}
}

private native FragmentMergeResult nativeMergeColumns(
Expand Down Expand Up @@ -186,8 +192,10 @@ private native FragmentMergeResult nativeMergeColumns(
*/
public FragmentUpdateResult updateColumns(
ArrowArrayStream stream, String leftOn, String rightOn) {
return nativeUpdateColumns(
dataset, fragmentMetadata.getId(), stream.memoryAddress(), leftOn, rightOn);
try (LockManager.ReadLock readLock = dataset.acquireReadLock()) {
return nativeUpdateColumns(
dataset, fragmentMetadata.getId(), stream.memoryAddress(), leftOn, rightOn);
}
}

public FragmentUpdateResult updateColumns(ArrowArrayStream stream) {
Expand Down
3 changes: 2 additions & 1 deletion java/src/main/java/org/lance/SqlQuery.java
Original file line number Diff line number Diff line change
Expand Up @@ -51,7 +51,8 @@ public SqlQuery withRowAddr(boolean withAddr) {
}

public ArrowReader intoBatchRecords() throws IOException {
try (ArrowArrayStream s = ArrowArrayStream.allocateNew(dataset.allocator())) {
try (LockManager.ReadLock readLock = dataset.acquireReadLock();
ArrowArrayStream s = ArrowArrayStream.allocateNew(dataset.allocator())) {
intoBatchRecords(
dataset, sql, Optional.ofNullable(table), withRowId, withRowAddr, s.memoryAddress());
return Data.importArrayStream(dataset.allocator(), s);
Expand Down
76 changes: 60 additions & 16 deletions java/src/main/java/org/lance/compaction/Compaction.java
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@

import org.lance.Dataset;
import org.lance.JniLoader;
import org.lance.LockManager;

import com.google.common.base.Preconditions;

Expand All @@ -32,21 +33,23 @@ public static CompactionPlan planCompaction(
Preconditions.checkNotNull(dataset);
Preconditions.checkNotNull(compactionOptions);

return nativePlanCompaction(
dataset,
compactionOptions.getTargetRowsPerFragment(),
compactionOptions.getMaxRowsPerGroup(),
compactionOptions.getMaxBytesPerFile(),
compactionOptions.getMaterializeDeletions(),
compactionOptions.getMaterializeDeletionsThreshold(),
compactionOptions.getNumThreads(),
compactionOptions.getBatchSize(),
compactionOptions.getDeferIndexRemap(),
compactionOptions.getCompactionMode(),
compactionOptions.getBinaryCopyReadBatchBytes(),
compactionOptions.getMaxSourceFragments(),
compactionOptions.getMaxSourceRows(),
compactionOptions.getMaxSourceBytes());
try (LockManager.ReadLock readLock = dataset.acquireReadLock()) {
return nativePlanCompaction(
dataset,
compactionOptions.getTargetRowsPerFragment(),
compactionOptions.getMaxRowsPerGroup(),
compactionOptions.getMaxBytesPerFile(),
compactionOptions.getMaterializeDeletions(),
compactionOptions.getMaterializeDeletionsThreshold(),
compactionOptions.getNumThreads(),
compactionOptions.getBatchSize(),
compactionOptions.getDeferIndexRemap(),
compactionOptions.getCompactionMode(),
compactionOptions.getBinaryCopyReadBatchBytes(),
compactionOptions.getMaxSourceFragments(),
compactionOptions.getMaxSourceRows(),
compactionOptions.getMaxSourceBytes());
}
}

public static CompactionMetrics commitCompaction(
Expand All @@ -72,7 +75,48 @@ public static CompactionMetrics commitCompaction(
compactionOptions.getMaxSourceBytes());
}

public static native CompactionMetrics nativeCommitCompaction(
/**
* Java wrapper around the raw commit-compaction JNI call. It acquires the dataset read lock so
* the native call cannot race with {@link Dataset#close()}; keep the raw native method private so
* no caller can bypass this lock.
*/
public static CompactionMetrics nativeCommitCompaction(
Dataset dataset,
List<RewriteResult> rewriteResults,
Optional<Long> targetRowsPerFragment,
Optional<Long> maxRowsPerGroup,
Optional<Long> maxBytesPerFile,
Optional<Boolean> materializeDeletions,
Optional<Float> materializeDeletionsThreshold,
Optional<Long> numThreads,
Optional<Long> batchSize,
Optional<Boolean> deferIndexRemap,
Optional<String> compactionMode,
Optional<Long> binaryCopyReadBatchBytes,
Optional<Long> maxSourceFragments,
Optional<Long> maxSourceRows,
Optional<Long> maxSourceBytes) {
try (LockManager.ReadLock readLock = dataset.acquireReadLock()) {
return commitCompactionNative(
dataset,
rewriteResults,
targetRowsPerFragment,
maxRowsPerGroup,
maxBytesPerFile,
materializeDeletions,
materializeDeletionsThreshold,
numThreads,
batchSize,
deferIndexRemap,
compactionMode,
binaryCopyReadBatchBytes,
maxSourceFragments,
maxSourceRows,
maxSourceBytes);
}
}

private static native CompactionMetrics commitCompactionNative(
Dataset dataset,
List<RewriteResult> rewriteResults,
Optional<Long> targetRowsPerFragment,
Expand Down
37 changes: 20 additions & 17 deletions java/src/main/java/org/lance/compaction/CompactionTask.java
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,7 @@
package org.lance.compaction;

import org.lance.Dataset;
import org.lance.LockManager;

import com.google.common.base.MoreObjects;

Expand Down Expand Up @@ -42,23 +43,25 @@ public String toString() {
}

public RewriteResult execute(Dataset dataset) {
return nativeExecute(
dataset,
taskData,
readVersion,
compactionOptions.getTargetRowsPerFragment(),
compactionOptions.getMaxRowsPerGroup(),
compactionOptions.getMaxBytesPerFile(),
compactionOptions.getMaterializeDeletions(),
compactionOptions.getMaterializeDeletionsThreshold(),
compactionOptions.getNumThreads(),
compactionOptions.getBatchSize(),
compactionOptions.getDeferIndexRemap(),
compactionOptions.getCompactionMode(),
compactionOptions.getBinaryCopyReadBatchBytes(),
compactionOptions.getMaxSourceFragments(),
compactionOptions.getMaxSourceRows(),
compactionOptions.getMaxSourceBytes());
try (LockManager.ReadLock readLock = dataset.acquireReadLock()) {
return nativeExecute(
dataset,
taskData,
readVersion,
compactionOptions.getTargetRowsPerFragment(),
compactionOptions.getMaxRowsPerGroup(),
compactionOptions.getMaxBytesPerFile(),
compactionOptions.getMaterializeDeletions(),
compactionOptions.getMaterializeDeletionsThreshold(),
compactionOptions.getNumThreads(),
compactionOptions.getBatchSize(),
compactionOptions.getDeferIndexRemap(),
compactionOptions.getCompactionMode(),
compactionOptions.getBinaryCopyReadBatchBytes(),
compactionOptions.getMaxSourceFragments(),
compactionOptions.getMaxSourceRows(),
compactionOptions.getMaxSourceBytes());
}
}

private native RewriteResult nativeExecute(
Expand Down
5 changes: 4 additions & 1 deletion java/src/main/java/org/lance/delta/DatasetDeltaBuilder.java
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@

import org.lance.Dataset;
import org.lance.JniLoader;
import org.lance.LockManager;

import java.util.Optional;

Expand Down Expand Up @@ -71,7 +72,9 @@ public DatasetDeltaBuilder withEndVersion(long version) {

/** Build the DatasetDelta after validating builder state. */
public DatasetDelta build() {
return nativeBuild(dataset, comparedAgainst, beginVersion, endVersion);
try (LockManager.ReadLock readLock = dataset.acquireReadLock()) {
return nativeBuild(dataset, comparedAgainst, beginVersion, endVersion);
}
}

private static native DatasetDelta nativeBuild(
Expand Down
9 changes: 7 additions & 2 deletions java/src/main/java/org/lance/index/vector/VectorTrainer.java
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@

import org.lance.Dataset;
import org.lance.JniLoader;
import org.lance.LockManager;
import org.lance.index.DistanceType;

import org.apache.arrow.util.Preconditions;
Expand Down Expand Up @@ -64,7 +65,9 @@ public static float[] trainIvfCentroids(
column != null && !column.isEmpty(), "column cannot be null or empty");
Preconditions.checkArgument(params != null, "params cannot be null");
Preconditions.checkArgument(distanceType != null, "distanceType cannot be null");
return nativeTrainIvfCentroids(dataset, column, params, distanceType.toString());
try (LockManager.ReadLock readLock = dataset.acquireReadLock()) {
return nativeTrainIvfCentroids(dataset, column, params, distanceType.toString());
}
}

/**
Expand Down Expand Up @@ -98,7 +101,9 @@ public static float[] trainPqCodebook(
column != null && !column.isEmpty(), "column cannot be null or empty");
Preconditions.checkArgument(params != null, "params cannot be null");
Preconditions.checkArgument(distanceType != null, "distanceType cannot be null");
return nativeTrainPqCodebook(dataset, column, params, distanceType.toString());
try (LockManager.ReadLock readLock = dataset.acquireReadLock()) {
return nativeTrainPqCodebook(dataset, column, params, distanceType.toString());
}
}

private static native float[] nativeTrainIvfCentroids(
Expand Down
59 changes: 31 additions & 28 deletions java/src/main/java/org/lance/ipc/AsyncScanner.java
Original file line number Diff line number Diff line change
Expand Up @@ -62,34 +62,37 @@ public static AsyncScanner create(
Preconditions.checkNotNull(dataset);
Preconditions.checkNotNull(options);
Preconditions.checkNotNull(allocator);
AsyncScanner scanner =
createAsyncScanner(
dataset,
options.getFragmentIds(),
options.getColumns(),
options.getSubstraitFilter(),
options.getFilter(),
options.getBatchSize(),
options.getBatchSizeBytes(),
options.getIoBufferSize(),
options.getLimit(),
options.getOffset(),
options.getNearest(),
options.getFullTextQuery(),
options.isPrefilter(),
options.isWithRowId(),
options.isWithRowAddress(),
options.getBatchReadahead(),
options.getFragmentReadahead(),
options.isScanInOrder(),
options.getLateMaterialization(),
options.getColumnOrderings(),
options.isUseScalarIndex(),
options.isFastSearch(),
options.getSubstraitAggregate(),
options.isIncludeDeletedRows(),
options.isStrictBatchSize(),
options.isDisableScoringAutoprojection());
AsyncScanner scanner;
try (LockManager.ReadLock readLock = dataset.acquireReadLock()) {
scanner =
createAsyncScanner(
dataset,
options.getFragmentIds(),
options.getColumns(),
options.getSubstraitFilter(),
options.getFilter(),
options.getBatchSize(),
options.getBatchSizeBytes(),
options.getIoBufferSize(),
options.getLimit(),
options.getOffset(),
options.getNearest(),
options.getFullTextQuery(),
options.isPrefilter(),
options.isWithRowId(),
options.isWithRowAddress(),
options.getBatchReadahead(),
options.getFragmentReadahead(),
options.isScanInOrder(),
options.getLateMaterialization(),
options.getColumnOrderings(),
options.isUseScalarIndex(),
options.isFastSearch(),
options.getSubstraitAggregate(),
options.isIncludeDeletedRows(),
options.isStrictBatchSize(),
options.isDisableScoringAutoprojection());
}
scanner.allocator = allocator;
return scanner;
}
Expand Down
Loading
Loading