Skip to main content

CompactionExecutor

Struct CompactionExecutor 

Source
pub struct CompactionExecutor { /* private fields */ }
Expand description

Executes compaction plans: reads N small files, merges them into a single AI-Lake file with a rebuilt index, and commits to the catalog.

The index algorithm is chosen via CompactionIndexStrategy (default: Auto, which detects GPU / CPU cores at compaction time — the same heuristic used by write_batch_auto).

For large tables use compact_deferred / run_deferred: the merged Parquet is persisted immediately and the HNSW build runs in a background Tokio task, decoupling I/O cost from CPU cost.

Implementations§

Source§

impl CompactionExecutor

Source

pub fn new(store: Arc<dyn Store>, policy: VectorStoragePolicy) -> Self

Source

pub fn with_index_strategy(self, strategy: CompactionIndexStrategy) -> Self

Override the default (Auto) index strategy for this executor.

Source

pub fn with_fts_config(self, cfg: FtsConfig) -> Self

Rebuild and embed a Tantivy FTS index in the compacted output file.

Source

pub async fn compact( &self, files: &[DataFileEntry], output_path: &str, ) -> AilakeResult<DataFileEntry>

Merge files into a single new file at output_path.

Reads all input files in parallel to minimise S3 latency, then rebuilds the HNSW / IVF-PQ index synchronously. For very large merges (N > 100 000 vectors) prefer compact_deferred, which offloads the index build to a background Tokio task.

Returns the DataFileEntry for the merged file.

Source

pub async fn compact_incremental( &self, files: &[DataFileEntry], output_path: &str, ) -> AilakeResult<DataFileEntry>

Merge files into a single new file using incremental HNSW insertion.

Identifies the dominant file — the file holding >= 40 % of the total row count — loads its existing HNSW graph from the AILK section, then calls HnswIndex::insert_node for every vector from the remaining files.

Complexity vs compact:

  • Full rebuild: O(N log N), N = total rows.
  • Incremental (this method): O(N_dom) deserialization + O(N_small × log N_dom). For a 90 / 10 split (N = 1 M, N_dom = 900 k) the speedup is ~7×.

Fallbacks (all degrade gracefully to compact):

  • No file holds >= 40 % of rows.
  • Dominant file’s HNSW cannot be loaded (IVF-PQ, IndexStatus::Indexing, corrupt).

RowId contract: dominant file’s vectors are placed first in the merged Parquet (positions 0..N_dom-1); other files follow. The existing RowIds from the dominant HNSW remain valid; new nodes receive RowIds N_dom..N-1.

Source

pub async fn compact_deferred( &self, files: &[DataFileEntry], output_path: &str, catalog: Arc<dyn CatalogProvider>, table: &TableIdent, ) -> AilakeResult<DataFileEntry>

Merge files into a single new file at output_path, writing Parquet immediately and building the HNSW / IVF-PQ index in a background Tokio task.

The merged file appears in the catalog as IndexStatus::Indexing until the background task completes; queries fall back to flat scan during that window (same behaviour as write_batch_deferred).

Returns the DataFileEntry with IndexStatus::Indexing. The entry transitions to Ready automatically when the background build finishes.

Source

pub async fn run( &self, planner: &CompactionPlanner, table: &TableIdent, catalog: Arc<dyn CatalogProvider>, output_prefix: &str, ) -> AilakeResult<Option<DataFileEntry>>

Full compaction workflow: plan, compact (synchronous HNSW rebuild), drop old files from catalog, commit.

Source

pub async fn run_deferred( &self, planner: &CompactionPlanner, table: &TableIdent, catalog: Arc<dyn CatalogProvider>, output_prefix: &str, ) -> AilakeResult<Option<DataFileEntry>>

Full compaction workflow with deferred HNSW build: plan, write merged Parquet immediately, commit as Indexing, spawn background index build.

Use for large tables where inline HNSW rebuild blocks too long.

Note: FTS index is not rebuilt in deferred mode — compact_deferred writes Parquet-only immediately and the background task (build_and_patch_index) only builds the HNSW/IVF-PQ index. Use run (synchronous) when FTS preservation on compaction is required.

Trait Implementations§

Source§

impl Clone for CompactionExecutor

Source§

fn clone(&self) -> CompactionExecutor

Returns a duplicate of the value. Read more
1.0.0 (const: unstable) · Source§

fn clone_from(&mut self, source: &Self)

Performs copy-assignment from source. Read more

Auto Trait Implementations§

Blanket Implementations§

Source§

impl<T> Any for T
where T: 'static + ?Sized,

Source§

fn type_id(&self) -> TypeId

Gets the TypeId of self. Read more
Source§

impl<T> Borrow<T> for T
where T: ?Sized,

Source§

fn borrow(&self) -> &T

Immutably borrows from an owned value. Read more
Source§

impl<T> BorrowMut<T> for T
where T: ?Sized,

Source§

fn borrow_mut(&mut self) -> &mut T

Mutably borrows from an owned value. Read more
Source§

impl<ST, DT> CastableFrom<ST, Initialized, Initialized> for DT
where ST: ?Sized, DT: ?Sized,

Source§

impl<ST, DT> CastableFrom<ST, Uninit, Uninit> for DT
where ST: ?Sized, DT: ?Sized,

Source§

impl<T> CloneToUninit for T
where T: Clone,

Source§

unsafe fn clone_to_uninit(&self, dest: *mut u8)

🔬This is a nightly-only experimental API. (clone_to_uninit)
Performs copy-assignment from self to dest. Read more
Source§

impl<T> Downcast for T
where T: Any,

Source§

fn into_any(self: Box<T>) -> Box<dyn Any>

Convert Box<dyn Trait> (where Trait: Downcast) to Box<dyn Any>. Box<dyn Any> can then be further downcast into Box<ConcreteType> where ConcreteType implements Trait.
Source§

fn into_any_rc(self: Rc<T>) -> Rc<dyn Any>

Convert Rc<Trait> (where Trait: Downcast) to Rc<Any>. Rc<Any> can then be further downcast into Rc<ConcreteType> where ConcreteType implements Trait.
Source§

fn as_any(&self) -> &(dyn Any + 'static)

Convert &Trait (where Trait: Downcast) to &Any. This is needed since Rust cannot generate &Any’s vtable from &Trait’s.
Source§

fn as_any_mut(&mut self) -> &mut (dyn Any + 'static)

Convert &mut Trait (where Trait: Downcast) to &Any. This is needed since Rust cannot generate &mut Any’s vtable from &mut Trait’s.
Source§

impl<T> DowncastSync for T
where T: Any + Send + Sync,

Source§

fn into_any_arc(self: Arc<T>) -> Arc<dyn Any + Send + Sync>

Convert Arc<Trait> (where Trait: Downcast) to Arc<Any>. Arc<Any> can then be further downcast into Arc<ConcreteType> where ConcreteType implements Trait.
Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

Source§

impl<T> Fruit for T
where T: Send + Downcast,

Source§

impl<T> Instrument for T

Source§

fn instrument(self, span: Span) -> Instrumented<Self>

Instruments this type with the provided Span, returning an Instrumented wrapper. Read more
Source§

fn in_current_span(self) -> Instrumented<Self>

Instruments this type with the current Span, returning an Instrumented wrapper. Read more
Source§

impl<T, U> Into<U> for T
where U: From<T>,

Source§

fn into(self) -> U

Calls U::from(self).

That is, this conversion is whatever the implementation of From<T> for U chooses to do.

Source§

impl<T> IntoEither for T

Source§

fn into_either(self, into_left: bool) -> Either<Self, Self>

Converts self into a Left variant of Either<Self, Self> if into_left is true. Converts self into a Right variant of Either<Self, Self> otherwise. Read more
Source§

fn into_either_with<F>(self, into_left: F) -> Either<Self, Self>
where F: FnOnce(&Self) -> bool,

Converts self into a Left variant of Either<Self, Self> if into_left(&self) returns true. Converts self into a Right variant of Either<Self, Self> otherwise. Read more
Source§

impl<T> Pointable for T

Source§

const ALIGN: usize

The alignment of pointer.
Source§

type Init = T

The type for initializers.
Source§

unsafe fn init(init: <T as Pointable>::Init) -> usize

Initializes a with the given initializer. Read more
Source§

unsafe fn deref<'a>(ptr: usize) -> &'a T

Dereferences the given pointer. Read more
Source§

unsafe fn deref_mut<'a>(ptr: usize) -> &'a mut T

Mutably dereferences the given pointer. Read more
Source§

unsafe fn drop(ptr: usize)

Drops the object pointed to by the given pointer. Read more
Source§

impl<T> Read<Exclusive, BecauseExclusive> for T
where T: ?Sized,

Source§

impl<T> Same for T

Source§

type Output = T

Should always be Self
Source§

impl<T> ToOwned for T
where T: Clone,

Source§

type Owned = T

The resulting type after obtaining ownership.
Source§

fn to_owned(&self) -> T

Creates owned data from borrowed data, usually by cloning. Read more
Source§

fn clone_into(&self, target: &mut T)

Uses borrowed data to replace owned data, usually by cloning. Read more
Source§

impl<T, U> TryFrom<U> for T
where U: Into<T>,

Source§

type Error = Infallible

The type returned in the event of a conversion error.
Source§

fn try_from(value: U) -> Result<T, <T as TryFrom<U>>::Error>

Performs the conversion.
Source§

impl<T, U> TryInto<U> for T
where U: TryFrom<T>,

Source§

type Error = <U as TryFrom<T>>::Error

The type returned in the event of a conversion error.
Source§

fn try_into(self) -> Result<U, <U as TryFrom<T>>::Error>

Performs the conversion.
Source§

impl<V, T> VZip<V> for T
where V: MultiLane<T>,

Source§

fn vzip(self) -> V

Source§

impl<T> WithSubscriber for T

Source§

fn with_subscriber<S>(self, subscriber: S) -> WithDispatch<Self>
where S: Into<Dispatch>,

Attaches the provided Subscriber to this type, returning a WithDispatch wrapper. Read more
Source§

fn with_current_subscriber(self) -> WithDispatch<Self>

Attaches the current default Subscriber to this type, returning a WithDispatch wrapper. Read more