Skip to content
Draft
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
1,602 changes: 1,574 additions & 28 deletions Cargo.lock

Large diffs are not rendered by default.

3 changes: 3 additions & 0 deletions Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -35,17 +35,20 @@ datafusion = { version = "54.0.0", default-features = false }
arrow = { version = "58.0.0", features = ["prettyprint", "ffi"] }
arrow-array = "58.0.0"
arrow-schema = "58.0.0"
bytes = "1"
# Direct to name `chrono::TimeDelta` (the field type of lance's public
# `AutoCleanupParams`) and `chrono::DateTime`/`Utc` (index metadata
# timestamps); already in the graph transitively via lance.
chrono = { version = "0.4", default-features = false }
half = "2"
tokio = { version = "1", features = ["rt-multi-thread", "sync"] }
futures = "0.3"
foyer = "=0.22.6"
log = "0.4"
libc = "0.2"
# Explicitly install the HTTP transport when embedded in a static C/C++ executable.
opendal = { version = "=0.59.2", default-features = false, features = ["http-transport-reqwest"] }
object_store = "0.14.2"
pin-project = "1.0"
prost = "0.14"
snafu = "0.9"
Expand Down
28 changes: 28 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -68,6 +68,7 @@ Based on the [liblance RFC](https://github.com/lance-format/lance/discussions/60
| [x] | Async scan | Callback-based `lance_scanner_scan_async()` for non-blocking scans |
| [x] | Dataset metadata | `lance_dataset_version()`, `lance_dataset_count_rows()`, `lance_dataset_latest_version()` |
| [x] | Filter pushdown | `lance_scanner_set_substrait_filter()` accepts a serialized Substrait `ExtendedExpression`; `lance_scanner_additional_sql_filter()` adds SQL predicates with AND before scanning starts |
| [x] | Data-file cache | Optional Foyer memory/disk cache for immutable `data/*.lance` reads |

## Multi-vector search

Expand Down Expand Up @@ -231,6 +232,33 @@ auto ds = lance::Dataset::open_with_session(session, "data.lance");
auto stats = session.cache_stats();
```

To add a process-local memory/disk cache for remote Lance data-file reads,
create the session with Foyer configuration. The cache is deliberately narrow:
whole-object, single-range, and batched range reads of direct `data/*.lance`
children are cached. Conditional and versioned reads, plus manifests, deletion
files, and index files, keep using Lance's normal paths. Use one shared session
for datasets that share the cache directory.
Cache entries are isolated by the underlying object-store instance because the
wrapper interface does not expose a complete backend identity. Datasets sharing
the same live store can reuse entries; a new store instance or process restart
starts a new cache namespace. Reopening the disk tier while that same store is
still alive can recover its entries. Identical bucket/path names on different
endpoints never share metadata, sizes, or data blocks.

```cpp
lance::DataCacheOptions data_cache{
"/var/cache/my-service/lance",
512ULL * 1024 * 1024, // memory tier
100ULL * 1024 * 1024 * 1024, // disk tier
1ULL * 1024 * 1024, // range-cache block
};
lance::Session session(
6ULL * 1024 * 1024 * 1024,
1ULL * 1024 * 1024 * 1024,
data_cache);
auto ds = lance::Dataset::open_with_session(session, "s3://bucket/data.lance");
```

### Open at a specific version

`lance_dataset_open` takes a `version` argument — `0` means the latest, any
Expand Down
63 changes: 63 additions & 0 deletions include/lance/lance.h
Original file line number Diff line number Diff line change
Expand Up @@ -214,6 +214,36 @@ typedef struct LanceSessionCacheStats {
uint64_t metadata_cache_size_bytes;
} LanceSessionCacheStats;

/**
* Configuration for the optional Foyer cache of immutable Lance data files.
*
* Whole-object, single-range, and batched range reads of direct
* `data/<file>.lance` children are cached. Conditional and versioned reads,
* plus metadata, deletion files, and index files, continue to use Lance's normal
* paths.
*/
typedef struct LanceDataCacheOptions {
const char* directory;
/** Maximum raw data bytes retained by Foyer's in-memory tier. */
uint64_t memory_capacity_bytes;
/** Maximum bytes allocated to Foyer's disk tier. */
uint64_t disk_capacity_bytes;
/** Data-file range cache unit, in bytes. */
uint64_t read_block_size_bytes;
} LanceDataCacheOptions;

/**
* Cumulative Foyer data-cache statistics for one opened dataset handle.
*
* Successful reads are accumulated. Both fields measure bytes returned to the
* dataset reader. Their sum is the logical data-file range bytes observed by
* the Foyer wrapper; block-aligned origin read amplification is not included.
*/
typedef struct LanceDataCacheStatistics {
uint64_t bytes_read_from_cache;
uint64_t bytes_read_from_remote;
} LanceDataCacheStatistics;

/**
* Create a session that can share metadata and index caches across datasets.
*
Expand All @@ -225,6 +255,25 @@ LanceSession* lance_session_new(
uint64_t metadata_cache_size_bytes
);

/**
* Create a shared Lance session with a Foyer data-file cache.
*
* `data_cache_options` and its `directory` field must not be NULL. The cache
* directory and all capacities are process configuration and remain owned by
* the caller; their values are copied during this call.
*
* `read_block_size_bytes` must be a non-zero multiple of 4096. The memory
* capacity must hold at least one read block. The disk capacity must be a
* multiple of 4096 and hold at least two read blocks.
*
* @return Session handle, or NULL on error
*/
LanceSession* lance_session_new_with_data_cache(
uint64_t index_cache_size_bytes,
uint64_t metadata_cache_size_bytes,
const LanceDataCacheOptions* data_cache_options
);

/**
* Close a session handle. Safe to call with NULL. Datasets previously opened
* with the session remain valid and retain the shared cache state.
Expand Down Expand Up @@ -281,6 +330,20 @@ LanceDataset* lance_dataset_open_with_session(
const LanceSession* session
);

/**
* Copy this dataset handle's cumulative data-cache statistics.
*
* A dataset not opened with a data cache reports all-zero statistics. The
* snapshot belongs only to this dataset handle; the underlying cache may
* still be shared by other datasets through a session.
*
* @return 0 on success, -1 on error
*/
int32_t lance_dataset_get_data_cache_statistics(
const LanceDataset* dataset,
LanceDataCacheStatistics* out_statistics
);

/** Close and free a dataset handle. Safe to call with NULL. */
void lance_dataset_close(LanceDataset* dataset);

Expand Down
29 changes: 29 additions & 0 deletions include/lance/lance.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -176,6 +176,13 @@ struct SqlColumn {

// ─── Shared Session ──────────────────────────────────────────────────────────

struct DataCacheOptions {
std::string directory;
uint64_t memory_capacity_bytes;
uint64_t disk_capacity_bytes;
uint64_t read_block_size_bytes;
};

class Session {
Handle<LanceSession, lance_session_close> handle_;

Expand All @@ -185,6 +192,21 @@ class Session {
if (!handle_) check_error();
}

Session(uint64_t index_cache_size_bytes,
uint64_t metadata_cache_size_bytes,
const DataCacheOptions& data_cache_options) {
LanceDataCacheOptions options{
data_cache_options.directory.c_str(),
data_cache_options.memory_capacity_bytes,
data_cache_options.disk_capacity_bytes,
data_cache_options.read_block_size_bytes,
};
handle_ = Handle<LanceSession, lance_session_close>(
lance_session_new_with_data_cache(
index_cache_size_bytes, metadata_cache_size_bytes, &options));
if (!handle_) check_error();
}

LanceSessionCacheStats cache_stats() const {
LanceSessionCacheStats stats{};
if (lance_session_get_cache_stats(handle_.get(), &stats) != 0)
Expand Down Expand Up @@ -350,6 +372,13 @@ class Dataset {
return Dataset(ds);
}

LanceDataCacheStatistics data_cache_statistics() const {
LanceDataCacheStatistics statistics{};
if (lance_dataset_get_data_cache_statistics(handle_.get(), &statistics) != 0)
check_error();
return statistics;
}

/// Write an Arrow record batch stream to a Lance dataset and return the
/// open dataset at the committed version.
///
Expand Down
68 changes: 68 additions & 0 deletions src/data_cache.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,68 @@
// SPDX-License-Identifier: Apache-2.0
// SPDX-FileCopyrightText: Copyright The Lance Authors

//! Private bridge between a shared data-cache backend and dataset handles.

use std::fmt::Debug;
use std::sync::Arc;

use lance::Dataset;
use lance_core::Result;

use crate::dataset::LanceDataset;
use crate::error::ffi_try;

/// Data-cache statistics owned by one opened dataset.
#[repr(C)]
#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
pub struct LanceDataCacheStatistics {
/// Requested bytes returned from usable data-cache entries.
pub bytes_read_from_cache: u64,
/// Requested bytes returned after a data-cache miss or fallback.
pub bytes_read_from_remote: u64,
}

pub(crate) trait DatasetDataCache: Debug + Send + Sync {
fn snapshot(&self) -> LanceDataCacheStatistics;

fn attach_fresh(&self, dataset: Dataset) -> (Dataset, Arc<dyn DatasetDataCache>);
}

pub(crate) trait DataCacheFactory: Debug + Send + Sync {
fn attach(&self, dataset: Dataset) -> (Dataset, Arc<dyn DatasetDataCache>);
}

/// Copy this dataset handle's cumulative data-cache statistics into
/// `out_statistics`.
///
/// A dataset not opened with a data cache reports all-zero statistics.
#[unsafe(no_mangle)]
pub unsafe extern "C" fn lance_dataset_get_data_cache_statistics(
dataset: *const LanceDataset,
out_statistics: *mut LanceDataCacheStatistics,
) -> i32 {
ffi_try!(
unsafe { dataset_get_data_cache_statistics_inner(dataset, out_statistics) },
neg
)
}

unsafe fn dataset_get_data_cache_statistics_inner(
dataset: *const LanceDataset,
out_statistics: *mut LanceDataCacheStatistics,
) -> Result<i32> {
if dataset.is_null() || out_statistics.is_null() {
return Err(lance_core::Error::invalid_input_source(
"dataset and out_statistics must not be NULL".into(),
));
}
let dataset = unsafe { &*dataset };
let statistics = dataset
.data_cache
.as_ref()
.map_or_else(LanceDataCacheStatistics::default, |cache| cache.snapshot());
unsafe {
std::ptr::write_unaligned(out_statistics, statistics);
}
Ok(0)
}
11 changes: 11 additions & 0 deletions src/dataset.rs
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,7 @@ use lance::Dataset;
use lance::dataset::builder::DatasetBuilder;
use lance_core::Result;

use crate::data_cache::DatasetDataCache;
use crate::error::{ffi_try, swallow_unwind};
use crate::helpers;
use crate::runtime::block_on;
Expand All @@ -23,6 +24,7 @@ use crate::stream_guard::guarded_ffi_stream_from_reader;
/// Opaque handle representing an opened Lance dataset.
pub struct LanceDataset {
pub(crate) inner: RwLock<Arc<Dataset>>,
pub(crate) data_cache: Option<Arc<dyn DatasetDataCache>>,
}

impl LanceDataset {
Expand Down Expand Up @@ -182,8 +184,16 @@ unsafe fn open_dataset_inner(
}

let dataset = block_on(builder.load())?;
let (dataset, data_cache) =
if let Some(factory) = session.and_then(|session| session.data_cache_factory.clone()) {
let (dataset, data_cache) = factory.attach(dataset);
(dataset, Some(data_cache))
} else {
(dataset, None)
};
let handle = LanceDataset {
inner: RwLock::new(Arc::new(dataset)),
data_cache,
};
Ok(Box::into_raw(Box::new(handle)))
}
Expand Down Expand Up @@ -519,6 +529,7 @@ mod tests {
.unwrap();
let handle = LanceDataset {
inner: RwLock::new(Arc::new(dataset)),
data_cache: None,
};
(tmp, handle)
}
Expand Down
Loading
Loading