Compare commits

..

2 Commits

Author SHA1 Message Date
Gatefixer 88d8a69a99 fix(python): scope Instructor compatibility shim 2026-08-06 01:30:35 +00:00
Gatefixer bd779bb7d5 fix(python): support legacy InstructorEmbedding downloads 2026-08-06 01:12:01 +00:00
11 changed files with 134 additions and 166 deletions
-3
View File
@@ -587,9 +587,6 @@ Modeled after ``VACUUM`` in PostgreSQL.
you have added or modified 100,000 or more records or run more than 20 data
modification operations.
Cleanup retention must exceed the longest expected write. Retention shorter
than 10 minutes requires `deleteUnverified: true` and exclusive write access.
#### Parameters
* **options?**: `Partial`<[`OptimizeOptions`](../interfaces/OptimizeOptions.md)>
+3 -6
View File
@@ -16,10 +16,7 @@ cleanupOlderThan: Date;
If set then all versions older than the given date
be removed. The current version will never be removed.
The default is 7 days. The resulting retention period must be longer than
the longest expected write. Values shorter than 10 minutes require
`deleteUnverified: true` and are only safe when no other process can write
to the dataset.
The default is 7 days
#### Example
@@ -29,8 +26,8 @@ const olderThan = new Date();
olderThan.setDate(olderThan.getDate() - 1));
tbl.optimize({cleanupOlderThan: olderThan});
// With exclusive access, delete all versions except the current version
tbl.optimize({cleanupOlderThan: new Date(), deleteUnverified: true});
// Delete all versions except the current version
tbl.optimize({cleanupOlderThan: new Date()});
```
***
+1 -9
View File
@@ -2196,15 +2196,7 @@ describe("when optimizing a dataset", () => {
});
it("cleanups old versions", async () => {
await expect(
table.optimize({ cleanupOlderThan: new Date() }),
).rejects.toThrow("at least 10 minutes");
expect(await table.version()).toBe(2);
const stats = await table.optimize({
cleanupOlderThan: new Date(),
deleteUnverified: true,
});
const stats = await table.optimize({ cleanupOlderThan: new Date() });
expect(stats.prune.bytesRemoved).toBeGreaterThan(0);
expect(stats.prune.oldVersionsRemoved).toBe(3);
});
+3 -9
View File
@@ -129,18 +129,15 @@ export interface OptimizeOptions {
/**
* If set then all versions older than the given date
* be removed. The current version will never be removed.
* The default is 7 days. The resulting retention period must be longer than
* the longest expected write. Values shorter than 10 minutes require
* `deleteUnverified: true` and are only safe when no other process can write
* to the dataset.
* The default is 7 days
* @example
* // Delete all versions older than 1 day
* const olderThan = new Date();
* olderThan.setDate(olderThan.getDate() - 1));
* tbl.optimize({cleanupOlderThan: olderThan});
*
* // With exclusive access, delete all versions except the current version
* tbl.optimize({cleanupOlderThan: new Date(), deleteUnverified: true});
* // Delete all versions except the current version
* tbl.optimize({cleanupOlderThan: new Date()});
*/
cleanupOlderThan: Date;
/**
@@ -747,9 +744,6 @@ export abstract class Table {
* optimize should be run frequently. A good rule of thumb is to run optimize if
* you have added or modified 100,000 or more records or run more than 20 data
* modification operations.
*
* Cleanup retention must exceed the longest expected write. Retention shorter
* than 10 minutes requires `deleteUnverified: true` and exclusive write access.
*/
abstract optimize(options?: Partial<OptimizeOptions>): Promise<OptimizeStats>;
/** List all indices that have been created with {@link Table.createIndex} */
-2
View File
@@ -6,7 +6,6 @@ use std::collections::HashMap;
use chrono::{DateTime, Utc};
use lancedb::ipc::{ipc_file_to_batches, ipc_file_to_schema};
use lancedb::table::optimize::validate_cleanup_options;
use lancedb::table::{
AddDataMode, ColumnAlteration as LanceColumnAlteration, Duration,
FieldMetadataUpdate as LanceFieldMetadataUpdate, FtsToken as LanceDbFtsToken,
@@ -560,7 +559,6 @@ impl Table {
} else {
None
};
validate_cleanup_options(older_than, delete_unverified).default_error()?;
let compaction_stats = inner
.optimize(OptimizeAction::Compact {
+56 -3
View File
@@ -3,6 +3,7 @@
from typing import List
from urllib.parse import unquote, urlparse
import numpy as np
@@ -125,9 +126,20 @@ class InstructorEmbeddingFunction(TextEmbeddingFunction):
@weak_lru(maxsize=1)
def get_model(self):
instructor_embedding = attempt_import_or_raise(
"InstructorEmbedding", "InstructorEmbedding"
)
huggingface_hub = attempt_import_or_raise("huggingface_hub", "huggingface-hub")
missing = object()
original_cached_download = getattr(huggingface_hub, "cached_download", missing)
if original_cached_download is missing:
huggingface_hub.cached_download = _cached_download(huggingface_hub)
try:
instructor_embedding = attempt_import_or_raise(
"InstructorEmbedding", "InstructorEmbedding"
)
finally:
if original_cached_download is missing:
del huggingface_hub.cached_download
torch = attempt_import_or_raise("torch", "torch")
model = instructor_embedding.INSTRUCTOR(self.name)
@@ -140,3 +152,44 @@ class InstructorEmbeddingFunction(TextEmbeddingFunction):
model, {torch.nn.Linear}, dtype=torch.qint8
)
return model
def _cached_download(huggingface_hub):
"""Provide the legacy download API used by sentence-transformers 2.2.x."""
def cached_download(
*,
url,
cache_dir=None,
force_filename=None,
library_name=None,
library_version=None,
user_agent=None,
use_auth_token=None,
**_,
):
path = urlparse(url).path.lstrip("/")
try:
repo_id, resolved_path = path.split("/resolve/", maxsplit=1)
revision, filename = resolved_path.split("/", maxsplit=1)
except ValueError as err:
raise ValueError(f"Unsupported Hugging Face Hub URL: {url}") from err
repo_id = unquote(repo_id)
revision = unquote(revision)
filename = unquote(filename)
# sentence-transformers derives force_filename from this Hub path with
# os.path.join. Using the URL path beneath local_dir produces the same
# local destination without sending Windows separators to the Hub.
return huggingface_hub.hf_hub_download(
repo_id=repo_id,
filename=filename,
revision=revision,
local_dir=cache_dir,
library_name=library_name,
library_version=library_version,
user_agent=user_agent,
token=use_auth_token,
)
return cached_download
+11 -39
View File
@@ -117,23 +117,6 @@ _MODEL_BACKED_TOKENIZER_ERRORS = (
"Failed to initialize default tokenizer",
)
_MIN_SAFE_CLEANUP_AGE = timedelta(minutes=10)
def _validate_cleanup_options(
older_than: Optional[timedelta], delete_unverified: bool
) -> None:
if (
older_than is not None
and older_than < _MIN_SAFE_CLEANUP_AGE
and not delete_unverified
):
raise ValueError(
"cleanup age must be at least 10 minutes unless delete_unverified is "
"true; short cleanup windows can remove a manifest still needed by an "
"in-progress write"
)
def _add_unique_note(exception: BaseException, note: str) -> None:
existing_notes = getattr(exception, "__notes__", ()) or ()
@@ -1784,9 +1767,7 @@ class Table(ABC):
----------
older_than: timedelta, default None
The minimum age of the version to delete. If None, then this defaults
to two weeks. This must be longer than the longest expected write.
Values shorter than 10 minutes require `delete_unverified=True` and
are only safe when no other process can write to the dataset.
to two weeks.
delete_unverified: bool, default False
Because they may be part of an in-progress transaction, files newer
than 7 days old are not deleted by default. If you are sure that
@@ -1854,11 +1835,9 @@ class Table(ABC):
Parameters
----------
cleanup_older_than: timedelta, optional default 7 days
All files belonging to versions older than this will be removed. The
latest version is never removed. This must be longer than the longest
expected write. Values shorter than 10 minutes require
`delete_unverified=True` and are only safe when no other process can
write to the dataset.
All files belonging to versions older than this will be removed. Set
to 0 days to remove all versions except the latest. The latest version
is never removed.
delete_unverified: bool, default False
Files leftover from a failed transaction may appear to be part of an
in-progress operation (e.g. appending new data) and these files will not
@@ -3807,9 +3786,7 @@ class LanceTable(Table):
----------
older_than: timedelta, default None
The minimum age of the version to delete. If None, then this defaults
to two weeks. This must be longer than the longest expected write.
Values shorter than 10 minutes require `delete_unverified=True` and
are only safe when no other process can write to the dataset.
to two weeks.
delete_unverified: bool, default False
Because they may be part of an in-progress transaction, files newer
than 7 days old are not deleted by default. If you are sure that
@@ -3822,7 +3799,6 @@ class LanceTable(Table):
The stats of the cleanup operation, including how many bytes were
freed.
"""
_validate_cleanup_options(older_than, delete_unverified)
return self.to_lance().cleanup_old_versions(
older_than, delete_unverified=delete_unverified
)
@@ -3868,11 +3844,9 @@ class LanceTable(Table):
Parameters
----------
cleanup_older_than: timedelta, optional default 7 days
All files belonging to versions older than this will be removed. The
latest version is never removed. This must be longer than the longest
expected write. Values shorter than 10 minutes require
`delete_unverified=True` and are only safe when no other process can
write to the dataset.
All files belonging to versions older than this will be removed. Set
to 0 days to remove all versions except the latest. The latest version
is never removed.
delete_unverified: bool, default False
Files leftover from a failed transaction may appear to be part of an
in-progress operation (e.g. appending new data) and these files will not
@@ -6064,11 +6038,9 @@ class AsyncTable:
Parameters
----------
cleanup_older_than: timedelta, optional default 7 days
All files belonging to versions older than this will be removed. The
latest version is never removed. This must be longer than the longest
expected write. Values shorter than 10 minutes require
`delete_unverified=True` and are only safe when no other process can
write to the dataset.
All files belonging to versions older than this will be removed. Set
to 0 days to remove all versions except the latest. The latest version
is never removed.
delete_unverified: bool, default False
Files leftover from a failed transaction may appear to be part of an
in-progress operation (e.g. appending new data) and these files will not
+56
View File
@@ -1,8 +1,11 @@
# SPDX-License-Identifier: Apache-2.0
# SPDX-FileCopyrightText: Copyright The LanceDB Authors
import ntpath
import os
import pickle
import sys
from types import ModuleType
from typing import List, Optional, Union
from unittest.mock import MagicMock, patch
@@ -522,6 +525,59 @@ def test_embedding_function_safe_model_dump(embedding_type):
)
def test_instructor_embedding_supports_huggingface_hub_without_cached_download(
tmp_path, monkeypatch
):
from lancedb.embeddings.instructor import InstructorEmbeddingFunction
hub_download = MagicMock(return_value="/cache/1_Pooling/config.json")
huggingface_hub = ModuleType("huggingface_hub")
huggingface_hub.hf_hub_download = hub_download
torch = ModuleType("torch")
monkeypatch.setitem(sys.modules, "huggingface_hub", huggingface_hub)
monkeypatch.setitem(sys.modules, "torch", torch)
monkeypatch.delitem(sys.modules, "InstructorEmbedding", raising=False)
monkeypatch.syspath_prepend(str(tmp_path))
(tmp_path / "InstructorEmbedding.py").write_text(
"from huggingface_hub import cached_download\n\n"
"class INSTRUCTOR:\n"
" def __init__(self, name):\n"
" self.name = name\n"
)
embedding = InstructorEmbeddingFunction.create(show_progress_bar=False)
instructor_model = embedding.get_model()
assert instructor_model.name == "hkunlp/instructor-base"
assert not hasattr(huggingface_hub, "cached_download")
instructor_embedding = sys.modules["InstructorEmbedding"]
path = instructor_embedding.cached_download(
url=(
"https://huggingface.co/hkunlp/instructor-base/resolve/abc123/"
"1_Pooling/config.json"
),
cache_dir="/cache",
force_filename=ntpath.join("1_Pooling", "config.json"),
library_name="sentence-transformers",
library_version="2.2.2",
use_auth_token="token",
)
assert path == "/cache/1_Pooling/config.json"
hub_download.assert_called_once_with(
repo_id="hkunlp/instructor-base",
filename="1_Pooling/config.json",
revision="abc123",
local_dir="/cache",
library_name="sentence-transformers",
library_version="2.2.2",
user_agent=None,
token="token",
)
@patch("time.sleep")
def test_retry(mock_sleep):
test_function = MagicMock(side_effect=[Exception] * 9 + ["result"])
+1 -11
View File
@@ -2980,9 +2980,6 @@ def test_compact_cleanup(tmp_db: DBConnection):
stats = table.cleanup_old_versions()
assert stats.bytes_removed == 0
with pytest.raises(ValueError, match="at least 10 minutes"):
table.cleanup_old_versions(older_than=timedelta(0))
stats = table.cleanup_old_versions(older_than=timedelta(0), delete_unverified=True)
assert stats.bytes_removed > 0
assert table.version == 4
@@ -3399,14 +3396,7 @@ async def test_optimize(mem_db_async: AsyncConnection):
assert stats.prune.bytes_removed == 0
assert stats.prune.old_versions_removed == 0
version_before_rejected_cleanup = await table.version()
with pytest.raises(ValueError, match="at least 10 minutes"):
await table.optimize(cleanup_older_than=timedelta(seconds=0))
assert await table.version() == version_before_rejected_cleanup
stats = await table.optimize(
cleanup_older_than=timedelta(seconds=0), delete_unverified=True
)
stats = await table.optimize(cleanup_older_than=timedelta(seconds=0))
assert stats.prune.bytes_removed > 0
assert stats.prune.old_versions_removed == 3
-2
View File
@@ -19,7 +19,6 @@ use arrow::{
};
use lancedb::blob::{BlobFile, BlobRangeRequest};
use lancedb::index::scalar::FtsIndexBuilder;
use lancedb::table::optimize::validate_cleanup_options;
use lancedb::table::{
AddDataMode, ColumnAlteration, Duration, FieldMetadataUpdate, FtsToken as LanceDbFtsToken,
NewColumnTransform, OptimizeAction, OptimizeOptions, Ref, Table as LanceDbTable,
@@ -1218,7 +1217,6 @@ impl Table {
} else {
None
};
validate_cleanup_options(older_than, delete_unverified).infer_error()?;
future_into_py(self_.py(), async move {
let compaction_stats = inner
.optimize(OptimizeAction::Compact {
+3 -82
View File
@@ -18,33 +18,7 @@ pub use chrono::Duration;
pub use lance::dataset::optimize::CompactionOptions;
use super::NativeTable;
use crate::error::{Error, Result};
const MIN_SAFE_CLEANUP_AGE_MINUTES: i64 = 10;
/// Validate version-cleanup options before starting an optimization operation.
///
/// This is public for the language bindings, which run compaction and cleanup as
/// separate operations and must reject unsafe cleanup options before compaction
/// changes the table.
#[doc(hidden)]
pub fn validate_cleanup_options(
older_than: Option<Duration>,
delete_unverified: Option<bool>,
) -> Result<()> {
let minimum_age =
Duration::try_minutes(MIN_SAFE_CLEANUP_AGE_MINUTES).expect("minimum cleanup age is valid");
if older_than.is_some_and(|age| age < minimum_age) && delete_unverified != Some(true) {
return Err(Error::InvalidInput {
message: format!(
"cleanup age must be at least {MIN_SAFE_CLEANUP_AGE_MINUTES} minutes unless \
delete_unverified is true; short cleanup windows can remove a manifest still \
needed by an in-progress write"
),
});
}
Ok(())
}
use crate::error::Result;
/// Optimize the dataset.
///
@@ -86,9 +60,7 @@ pub enum OptimizeAction {
///
/// Once a version is pruned it can no longer be checked out.
Prune {
/// The duration of time to keep versions of the dataset. This should be longer than the
/// longest expected write. Values shorter than 10 minutes require `delete_unverified` to
/// be true and are only safe when no other process can write to the dataset.
/// The duration of time to keep versions of the dataset.
older_than: Option<Duration>,
/// Because they may be part of an in-progress transaction, files newer than 7 days old are not deleted by default.
/// If you are sure that there are no in-progress transactions, then you can set this to True to delete all files older than `older_than`.
@@ -192,15 +164,6 @@ pub(crate) async fn execute_optimize(
table: &NativeTable,
action: OptimizeAction,
) -> Result<OptimizeStats> {
if let OptimizeAction::Prune {
older_than,
delete_unverified,
..
} = &action
{
validate_cleanup_options(*older_than, *delete_unverified)?;
}
let mut stats = OptimizeStats {
compaction: None,
prune: None,
@@ -259,7 +222,7 @@ mod tests {
use crate::connect;
use crate::index::{Index, scalar::BTreeIndexBuilder};
use crate::query::ExecutableQuery;
use crate::table::{CompactionOptions, Duration, OptimizeAction, OptimizeStats};
use crate::table::{CompactionOptions, OptimizeAction, OptimizeStats};
use futures::TryStreamExt;
#[tokio::test]
@@ -417,48 +380,6 @@ mod tests {
assert_eq!(all_values, expected);
}
#[tokio::test]
async fn test_optimize_rejects_unsafe_cleanup_age() {
let conn = connect("memory://").execute().await.unwrap();
let batch = RecordBatch::try_new(
Arc::new(Schema::new(vec![Field::new("i", DataType::Int32, false)])),
vec![Arc::new(Int32Array::from_iter_values(0..10))],
)
.unwrap();
let table = conn
.create_table("test_unsafe_prune", batch.clone())
.execute()
.await
.unwrap();
table.add(batch).execute().await.unwrap();
let versions_before = table
.list_versions()
.await
.unwrap()
.into_iter()
.map(|version| version.version)
.collect::<Vec<_>>();
let err = table
.optimize(OptimizeAction::Prune {
older_than: Some(Duration::zero()),
delete_unverified: Some(false),
error_if_tagged_old_versions: None,
})
.await
.unwrap_err();
assert!(err.to_string().contains("at least 10 minutes"));
let versions_after = table
.list_versions()
.await
.unwrap()
.into_iter()
.map(|version| version.version)
.collect::<Vec<_>>();
assert_eq!(versions_after, versions_before);
}
#[tokio::test]
async fn test_optimize_index() {
let conn = connect("memory://").execute().await.unwrap();