2020-10-23 14:11:00 +02:00
|
|
|
use std::borrow::Cow;
|
2024-06-12 14:05:52 +02:00
|
|
|
use std::collections::{BTreeMap, HashMap, HashSet};
|
2020-10-23 14:11:00 +02:00
|
|
|
use std::fs::File;
|
|
|
|
|
2024-08-30 11:49:47 +02:00
|
|
|
use either::Either;
|
2022-03-23 17:28:41 +01:00
|
|
|
use fxhash::FxHashMap;
|
2021-08-31 11:44:15 +02:00
|
|
|
use itertools::Itertools;
|
2024-11-18 11:59:18 +01:00
|
|
|
use obkv::{KvReader, KvWriter};
|
2020-10-23 14:11:00 +02:00
|
|
|
use roaring::RoaringBitmap;
|
2022-06-21 11:12:51 +02:00
|
|
|
use serde_json::Value;
|
2022-04-11 15:43:18 +02:00
|
|
|
use smartstring::SmartString;
|
2020-10-23 14:11:00 +02:00
|
|
|
|
2023-02-16 18:42:47 +01:00
|
|
|
use super::helpers::{
|
2024-11-18 11:59:18 +01:00
|
|
|
create_sorter, sorter_into_reader, EitherObkvMerge, ObkvsKeepLastAdditionMergeDeletions,
|
|
|
|
ObkvsMergeAdditionsAndDeletions,
|
2023-02-16 18:42:47 +01:00
|
|
|
};
|
2024-08-30 11:49:47 +02:00
|
|
|
use super::{IndexDocumentsMethod, IndexerConfig, KeepFirst};
|
2024-11-18 11:59:18 +01:00
|
|
|
use crate::documents::DocumentsBatchIndex;
|
2021-08-31 11:44:15 +02:00
|
|
|
use crate::error::{Error, InternalError, UserError};
|
2022-12-19 15:59:22 +01:00
|
|
|
use crate::index::{db_name, main_key};
|
2024-05-28 14:53:45 +02:00
|
|
|
use crate::update::del_add::{
|
2024-11-18 11:59:18 +01:00
|
|
|
into_del_add_obkv, into_del_add_obkv_conditional_operation, DelAddOperation,
|
2024-05-28 14:53:45 +02:00
|
|
|
};
|
2023-11-02 13:37:54 +01:00
|
|
|
use crate::update::index_documents::GrenadParameters;
|
2024-11-18 11:59:18 +01:00
|
|
|
use crate::update::settings::InnerIndexSettingsDiff;
|
|
|
|
use crate::update::AvailableIds;
|
2024-06-12 17:10:19 +02:00
|
|
|
use crate::vector::parsed_vectors::{ExplicitVectors, VectorOrArrayOfVectors};
|
2024-09-18 18:13:37 +02:00
|
|
|
use crate::vector::settings::WriteBackToDocuments;
|
2024-09-19 10:35:17 +02:00
|
|
|
use crate::vector::ArroyWrapper;
|
2024-05-22 16:05:55 +02:00
|
|
|
use crate::{
|
|
|
|
is_faceted_by, FieldDistribution, FieldId, FieldIdMapMissingEntry, FieldsIdsMap, Index, Result,
|
|
|
|
};
|
2020-10-24 16:23:08 +02:00
|
|
|
|
2020-10-23 14:11:00 +02:00
|
|
|
pub struct TransformOutput {
|
2021-01-20 17:27:43 +01:00
|
|
|
pub primary_key: String,
|
2024-03-26 13:27:43 +01:00
|
|
|
pub settings_diff: InnerIndexSettingsDiff,
|
2021-06-21 15:57:41 +02:00
|
|
|
pub field_distribution: FieldDistribution,
|
2020-10-23 14:11:00 +02:00
|
|
|
pub documents_count: usize,
|
2024-05-21 14:53:26 +02:00
|
|
|
pub original_documents: Option<File>,
|
2024-05-21 16:16:36 +02:00
|
|
|
pub flattened_documents: Option<File>,
|
2020-10-23 14:11:00 +02:00
|
|
|
}
|
|
|
|
|
2020-11-22 11:54:04 +01:00
|
|
|
/// Extract the external ids, deduplicate and compute the new internal documents ids
|
2020-11-01 11:50:10 +01:00
|
|
|
/// and fields ids, writing all the documents under their internal ids into a final file.
|
|
|
|
///
|
|
|
|
/// Outputs the new `FieldsIdsMap`, the new `UsersIdsDocumentsIds` map, the new documents ids,
|
|
|
|
/// the replaced documents ids, the number of documents in this update and the file
|
|
|
|
/// containing all those documents.
|
2021-12-08 14:12:07 +01:00
|
|
|
pub struct Transform<'a, 'i> {
|
2020-10-26 20:18:10 +01:00
|
|
|
pub index: &'i Index,
|
2021-12-08 14:12:07 +01:00
|
|
|
indexer_settings: &'a IndexerConfig,
|
|
|
|
pub index_documents_method: IndexDocumentsMethod,
|
2020-10-23 14:11:00 +02:00
|
|
|
}
|
|
|
|
|
2023-02-14 17:55:26 +01:00
|
|
|
/// This enum is specific to the grenad sorter stored in the transform.
|
|
|
|
/// It's used as the first byte of the grenads and tells you if the document id was an addition or a deletion.
|
2023-02-08 12:53:38 +01:00
|
|
|
#[repr(u8)]
|
2023-02-16 18:42:47 +01:00
|
|
|
pub enum Operation {
|
2023-02-08 12:53:38 +01:00
|
|
|
Addition,
|
|
|
|
Deletion,
|
|
|
|
}
|
|
|
|
|
2021-12-08 14:12:07 +01:00
|
|
|
impl<'a, 'i> Transform<'a, 'i> {
|
|
|
|
pub fn new(
|
2024-07-09 17:25:39 +02:00
|
|
|
wtxn: &mut heed::RwTxn<'_>,
|
2021-12-08 14:12:07 +01:00
|
|
|
index: &'i Index,
|
|
|
|
indexer_settings: &'a IndexerConfig,
|
|
|
|
index_documents_method: IndexDocumentsMethod,
|
2024-06-13 17:47:44 +02:00
|
|
|
_autogenerate_docids: bool,
|
2022-03-23 17:28:41 +01:00
|
|
|
) -> Result<Self> {
|
2024-08-30 11:49:47 +02:00
|
|
|
use IndexDocumentsMethod::{ReplaceDocuments, UpdateDocuments};
|
|
|
|
|
2021-12-08 14:12:07 +01:00
|
|
|
// We must choose the appropriate merge function for when two or more documents
|
|
|
|
// with the same user id must be merged or fully replaced in the same batch.
|
|
|
|
let merge_function = match index_documents_method {
|
2024-08-30 11:49:47 +02:00
|
|
|
ReplaceDocuments => Either::Left(ObkvsKeepLastAdditionMergeDeletions),
|
|
|
|
UpdateDocuments => Either::Right(ObkvsMergeAdditionsAndDeletions),
|
2021-12-08 14:12:07 +01:00
|
|
|
};
|
|
|
|
|
|
|
|
// We initialize the sorter with the user indexing settings.
|
2023-11-02 14:47:43 +01:00
|
|
|
let original_sorter = create_sorter(
|
|
|
|
grenad::SortAlgorithm::Stable,
|
2024-09-30 16:08:29 +02:00
|
|
|
merge_function,
|
2023-11-02 14:47:43 +01:00
|
|
|
indexer_settings.chunk_compression_type,
|
|
|
|
indexer_settings.chunk_compression_level,
|
|
|
|
indexer_settings.max_nb_chunks,
|
|
|
|
indexer_settings.max_memory.map(|mem| mem / 2),
|
2024-10-17 09:30:18 +02:00
|
|
|
true,
|
2023-11-02 14:47:43 +01:00
|
|
|
);
|
2021-12-08 14:12:07 +01:00
|
|
|
|
2022-03-23 17:28:41 +01:00
|
|
|
// We initialize the sorter with the user indexing settings.
|
2023-11-02 14:47:43 +01:00
|
|
|
let flattened_sorter = create_sorter(
|
|
|
|
grenad::SortAlgorithm::Stable,
|
|
|
|
merge_function,
|
|
|
|
indexer_settings.chunk_compression_type,
|
|
|
|
indexer_settings.chunk_compression_level,
|
|
|
|
indexer_settings.max_nb_chunks,
|
|
|
|
indexer_settings.max_memory.map(|mem| mem / 2),
|
2024-10-17 09:30:18 +02:00
|
|
|
true,
|
2023-11-02 14:47:43 +01:00
|
|
|
);
|
2022-06-07 15:44:55 +02:00
|
|
|
let documents_ids = index.documents_ids(wtxn)?;
|
2022-03-23 17:28:41 +01:00
|
|
|
|
2024-11-18 17:39:55 +01:00
|
|
|
Ok(Transform { index, indexer_settings, index_documents_method })
|
2021-12-08 14:12:07 +01:00
|
|
|
}
|
|
|
|
|
2022-03-23 17:28:41 +01:00
|
|
|
// Flatten a document from the fields ids map contained in self and insert the new
|
2022-04-12 11:22:36 +02:00
|
|
|
// created fields. Returns `None` if the document doesn't need to be flattened.
|
2024-03-26 13:27:43 +01:00
|
|
|
#[tracing::instrument(
|
|
|
|
level = "trace",
|
|
|
|
skip(obkv, fields_ids_map),
|
|
|
|
target = "indexing::transform"
|
|
|
|
)]
|
|
|
|
fn flatten_from_fields_ids_map(
|
2024-08-29 19:20:10 +02:00
|
|
|
obkv: &KvReader<FieldId>,
|
2024-03-26 13:27:43 +01:00
|
|
|
fields_ids_map: &mut FieldsIdsMap,
|
|
|
|
) -> Result<Option<Vec<u8>>> {
|
2022-04-12 11:22:36 +02:00
|
|
|
if obkv
|
|
|
|
.iter()
|
|
|
|
.all(|(_, value)| !json_depth_checker::should_flatten_from_unchecked_slice(value))
|
|
|
|
{
|
|
|
|
return Ok(None);
|
|
|
|
}
|
|
|
|
|
2022-04-25 14:09:52 +02:00
|
|
|
// store the keys and values the original obkv + the flattened json
|
|
|
|
// We first extract all the key+value out of the obkv. If a value is not nested
|
|
|
|
// we keep a reference on its value. If the value is nested we'll get its value
|
|
|
|
// as an owned `Vec<u8>` after flattening it.
|
2024-07-09 17:25:39 +02:00
|
|
|
let mut key_value: Vec<(FieldId, Cow<'_, [u8]>)> = Vec::new();
|
2022-04-25 14:09:52 +02:00
|
|
|
|
|
|
|
// the object we're going to use to store the fields that need to be flattened.
|
2022-03-23 17:28:41 +01:00
|
|
|
let mut doc = serde_json::Map::new();
|
|
|
|
|
2022-04-25 14:09:52 +02:00
|
|
|
// we recreate a json containing only the fields that needs to be flattened.
|
|
|
|
// all the raw values get inserted directly in the `key_value` vec.
|
|
|
|
for (key, value) in obkv.iter() {
|
|
|
|
if json_depth_checker::should_flatten_from_unchecked_slice(value) {
|
2024-03-26 13:27:43 +01:00
|
|
|
let key = fields_ids_map.name(key).ok_or(FieldIdMapMissingEntry::FieldId {
|
2022-04-25 14:09:52 +02:00
|
|
|
field_id: key,
|
|
|
|
process: "Flatten from fields ids map.",
|
|
|
|
})?;
|
|
|
|
|
|
|
|
let value = serde_json::from_slice::<Value>(value)
|
|
|
|
.map_err(crate::error::InternalError::SerdeJson)?;
|
|
|
|
doc.insert(key.to_string(), value);
|
|
|
|
} else {
|
|
|
|
key_value.push((key, value.into()));
|
|
|
|
}
|
2022-03-23 17:28:41 +01:00
|
|
|
}
|
|
|
|
|
|
|
|
let flattened = flatten_serde_json::flatten(&doc);
|
|
|
|
|
2022-04-25 14:09:52 +02:00
|
|
|
// Once we have the flattened version we insert all the new generated fields_ids
|
|
|
|
// (if any) in the fields ids map and serialize the value.
|
|
|
|
for (key, value) in flattened.into_iter() {
|
2024-03-26 13:27:43 +01:00
|
|
|
let fid = fields_ids_map.insert(&key).ok_or(UserError::AttributeLimitReached)?;
|
2022-03-23 17:28:41 +01:00
|
|
|
let value = serde_json::to_vec(&value).map_err(InternalError::SerdeJson)?;
|
2022-04-25 14:09:52 +02:00
|
|
|
key_value.push((fid, value.into()));
|
2022-03-23 17:28:41 +01:00
|
|
|
}
|
|
|
|
|
2022-04-25 14:09:52 +02:00
|
|
|
// we sort the key. If there was a conflict between the obkv and the new generated value the
|
|
|
|
// keys will be consecutive.
|
|
|
|
key_value.sort_unstable_by_key(|(key, _)| *key);
|
|
|
|
|
|
|
|
let mut buffer = Vec::new();
|
|
|
|
Self::create_obkv_from_key_value(&mut key_value, &mut buffer)?;
|
2022-04-12 11:22:36 +02:00
|
|
|
Ok(Some(buffer))
|
2022-03-23 17:28:41 +01:00
|
|
|
}
|
|
|
|
|
2022-04-25 14:09:52 +02:00
|
|
|
/// Generate an obkv from a slice of key / value sorted by key.
|
|
|
|
fn create_obkv_from_key_value(
|
2024-07-09 17:25:39 +02:00
|
|
|
key_value: &mut [(FieldId, Cow<'_, [u8]>)],
|
2022-04-25 14:09:52 +02:00
|
|
|
output_buffer: &mut Vec<u8>,
|
|
|
|
) -> Result<()> {
|
|
|
|
debug_assert!(
|
|
|
|
key_value.windows(2).all(|vec| vec[0].0 <= vec[1].0),
|
|
|
|
"The slice of key / value pair must be sorted."
|
|
|
|
);
|
|
|
|
|
|
|
|
output_buffer.clear();
|
|
|
|
let mut writer = KvWriter::new(output_buffer);
|
|
|
|
|
|
|
|
let mut skip_next_value = false;
|
|
|
|
for things in key_value.windows(2) {
|
|
|
|
if skip_next_value {
|
|
|
|
skip_next_value = false;
|
|
|
|
continue;
|
|
|
|
}
|
|
|
|
let (key1, value1) = &things[0];
|
|
|
|
let (key2, value2) = &things[1];
|
|
|
|
|
|
|
|
// now we're going to look for conflicts between the keys. For example the following documents would cause a conflict:
|
|
|
|
// { "doggo.name": "jean", "doggo": { "name": "paul" } }
|
|
|
|
// we should find a first "doggo.name" from the obkv and a second one from the flattening.
|
|
|
|
// but we must generate the following document:
|
|
|
|
// { "doggo.name": ["jean", "paul"] }
|
|
|
|
// thus we're going to merge the value from the obkv and the flattened document in a single array and skip the next
|
|
|
|
// iteration.
|
|
|
|
if key1 == key2 {
|
|
|
|
skip_next_value = true;
|
|
|
|
|
|
|
|
let value1 = serde_json::from_slice(value1)
|
|
|
|
.map_err(crate::error::InternalError::SerdeJson)?;
|
|
|
|
let value2 = serde_json::from_slice(value2)
|
|
|
|
.map_err(crate::error::InternalError::SerdeJson)?;
|
|
|
|
let value = match (value1, value2) {
|
|
|
|
(Value::Array(mut left), Value::Array(mut right)) => {
|
|
|
|
left.append(&mut right);
|
|
|
|
Value::Array(left)
|
|
|
|
}
|
|
|
|
(Value::Array(mut array), value) | (value, Value::Array(mut array)) => {
|
|
|
|
array.push(value);
|
|
|
|
Value::Array(array)
|
|
|
|
}
|
|
|
|
(left, right) => Value::Array(vec![left, right]),
|
|
|
|
};
|
|
|
|
|
|
|
|
let value = serde_json::to_vec(&value).map_err(InternalError::SerdeJson)?;
|
|
|
|
writer.insert(*key1, value)?;
|
|
|
|
} else {
|
|
|
|
writer.insert(*key1, value1)?;
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
if !skip_next_value {
|
|
|
|
// the unwrap is safe here, we know there was at least one value in the document
|
|
|
|
let (key, value) = key_value.last().unwrap();
|
|
|
|
writer.insert(*key, value)?;
|
|
|
|
}
|
|
|
|
|
2022-03-23 17:28:41 +01:00
|
|
|
Ok(())
|
|
|
|
}
|
|
|
|
|
2024-04-17 10:54:48 +02:00
|
|
|
/// Rebind the field_ids of the provided document to their values
|
|
|
|
/// based on the field_ids_maps difference between the old and the new settings,
|
|
|
|
/// then fill the provided buffers with delta documents using KvWritterDelAdd.
|
2024-06-12 14:05:52 +02:00
|
|
|
#[allow(clippy::too_many_arguments)] // need the vectors + fid, feel free to create a struct xo xo
|
2024-03-26 13:27:43 +01:00
|
|
|
fn rebind_existing_document(
|
2024-08-29 19:20:10 +02:00
|
|
|
old_obkv: &KvReader<FieldId>,
|
2024-03-26 13:27:43 +01:00
|
|
|
settings_diff: &InnerIndexSettingsDiff,
|
2024-05-22 16:05:55 +02:00
|
|
|
modified_faceted_fields: &HashSet<String>,
|
2024-06-12 14:05:52 +02:00
|
|
|
mut injected_vectors: serde_json::Map<String, serde_json::Value>,
|
|
|
|
old_vectors_fid: Option<FieldId>,
|
2024-05-22 16:05:55 +02:00
|
|
|
original_obkv_buffer: Option<&mut Vec<u8>>,
|
|
|
|
flattened_obkv_buffer: Option<&mut Vec<u8>>,
|
2024-03-26 13:27:43 +01:00
|
|
|
) -> Result<()> {
|
2024-05-22 16:05:55 +02:00
|
|
|
// Always keep the primary key.
|
|
|
|
let is_primary_key = |id: FieldId| -> bool { settings_diff.primary_key_id == Some(id) };
|
|
|
|
|
|
|
|
// If only a faceted field has been added, keep only this field.
|
|
|
|
let must_reindex_facets = settings_diff.reindex_facets();
|
|
|
|
let necessary_faceted_field = |id: FieldId| -> bool {
|
|
|
|
let field_name = settings_diff.new.fields_ids_map.name(id).unwrap();
|
|
|
|
must_reindex_facets
|
|
|
|
&& modified_faceted_fields
|
|
|
|
.iter()
|
|
|
|
.any(|long| is_faceted_by(long, field_name) || is_faceted_by(field_name, long))
|
|
|
|
};
|
|
|
|
|
|
|
|
// Alway provide all fields when vectors are involved because
|
|
|
|
// we need the fields for the prompt/templating.
|
|
|
|
let reindex_vectors = settings_diff.reindex_vectors();
|
|
|
|
|
2024-05-29 17:46:28 +02:00
|
|
|
// The operations that we must perform on the different fields.
|
|
|
|
let mut operations = HashMap::new();
|
2024-06-13 14:20:48 +02:00
|
|
|
let mut error_seen = false;
|
2024-05-28 14:53:45 +02:00
|
|
|
|
2024-03-26 13:27:43 +01:00
|
|
|
let mut obkv_writer = KvWriter::<_, FieldId>::memory();
|
2024-06-12 14:05:52 +02:00
|
|
|
'write_fid: for (id, val) in old_obkv.iter() {
|
|
|
|
if !injected_vectors.is_empty() {
|
|
|
|
'inject_vectors: {
|
|
|
|
let Some(vectors_fid) = old_vectors_fid else { break 'inject_vectors };
|
|
|
|
|
2024-06-12 17:10:19 +02:00
|
|
|
if id < vectors_fid {
|
2024-06-12 14:05:52 +02:00
|
|
|
break 'inject_vectors;
|
|
|
|
}
|
|
|
|
|
2024-06-12 17:10:19 +02:00
|
|
|
let mut existing_vectors = if id == vectors_fid {
|
|
|
|
let existing_vectors: std::result::Result<
|
|
|
|
serde_json::Map<String, serde_json::Value>,
|
|
|
|
serde_json::Error,
|
|
|
|
> = serde_json::from_slice(val);
|
|
|
|
|
|
|
|
match existing_vectors {
|
|
|
|
Ok(existing_vectors) => existing_vectors,
|
|
|
|
Err(error) => {
|
2024-06-13 14:20:48 +02:00
|
|
|
if !error_seen {
|
|
|
|
tracing::error!(%error, "Unexpected `_vectors` field that is not a map. Treating as an empty map");
|
|
|
|
error_seen = true;
|
|
|
|
}
|
2024-06-12 17:10:19 +02:00
|
|
|
Default::default()
|
|
|
|
}
|
2024-06-12 14:05:52 +02:00
|
|
|
}
|
2024-06-12 17:10:19 +02:00
|
|
|
} else {
|
|
|
|
Default::default()
|
2024-06-12 14:05:52 +02:00
|
|
|
};
|
|
|
|
|
|
|
|
existing_vectors.append(&mut injected_vectors);
|
|
|
|
|
2024-06-12 17:10:19 +02:00
|
|
|
operations.insert(vectors_fid, DelAddOperation::DeletionAndAddition);
|
|
|
|
obkv_writer
|
|
|
|
.insert(vectors_fid, serde_json::to_vec(&existing_vectors).unwrap())?;
|
|
|
|
if id == vectors_fid {
|
|
|
|
continue 'write_fid;
|
|
|
|
}
|
2024-06-12 14:05:52 +02:00
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2024-05-29 17:46:28 +02:00
|
|
|
if is_primary_key(id) || necessary_faceted_field(id) || reindex_vectors {
|
|
|
|
operations.insert(id, DelAddOperation::DeletionAndAddition);
|
2024-03-26 13:27:43 +01:00
|
|
|
obkv_writer.insert(id, val)?;
|
2024-05-29 17:46:28 +02:00
|
|
|
} else if let Some(operation) = settings_diff.reindex_searchable_id(id) {
|
|
|
|
operations.insert(id, operation);
|
2024-05-28 14:53:45 +02:00
|
|
|
obkv_writer.insert(id, val)?;
|
2024-03-26 13:27:43 +01:00
|
|
|
}
|
|
|
|
}
|
2024-06-12 17:10:19 +02:00
|
|
|
if !injected_vectors.is_empty() {
|
|
|
|
'inject_vectors: {
|
|
|
|
let Some(vectors_fid) = old_vectors_fid else { break 'inject_vectors };
|
|
|
|
|
|
|
|
operations.insert(vectors_fid, DelAddOperation::DeletionAndAddition);
|
|
|
|
obkv_writer.insert(vectors_fid, serde_json::to_vec(&injected_vectors).unwrap())?;
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2024-04-03 11:19:45 +02:00
|
|
|
let data = obkv_writer.into_inner()?;
|
2024-08-29 19:20:10 +02:00
|
|
|
let obkv = KvReader::<FieldId>::from_slice(&data);
|
2024-03-26 13:27:43 +01:00
|
|
|
|
2024-05-22 16:05:55 +02:00
|
|
|
if let Some(original_obkv_buffer) = original_obkv_buffer {
|
|
|
|
original_obkv_buffer.clear();
|
|
|
|
into_del_add_obkv(obkv, DelAddOperation::DeletionAndAddition, original_obkv_buffer)?;
|
|
|
|
}
|
2024-03-26 13:27:43 +01:00
|
|
|
|
2024-05-22 16:05:55 +02:00
|
|
|
if let Some(flattened_obkv_buffer) = flattened_obkv_buffer {
|
|
|
|
// take the non-flattened version if flatten_from_fields_ids_map returns None.
|
|
|
|
let mut fields_ids_map = settings_diff.new.fields_ids_map.clone();
|
2024-09-24 17:24:50 +02:00
|
|
|
let flattened = Self::flatten_from_fields_ids_map(obkv, &mut fields_ids_map)?;
|
2024-08-29 19:20:10 +02:00
|
|
|
let flattened = flattened.as_deref().map_or(obkv, KvReader::from_slice);
|
2024-03-26 13:27:43 +01:00
|
|
|
|
2024-05-22 16:05:55 +02:00
|
|
|
flattened_obkv_buffer.clear();
|
2024-05-28 14:53:45 +02:00
|
|
|
into_del_add_obkv_conditional_operation(flattened, flattened_obkv_buffer, |id| {
|
2024-05-29 17:46:28 +02:00
|
|
|
operations.get(&id).copied().unwrap_or(DelAddOperation::DeletionAndAddition)
|
2024-05-28 14:53:45 +02:00
|
|
|
})?;
|
2024-05-22 16:05:55 +02:00
|
|
|
}
|
2024-03-26 13:27:43 +01:00
|
|
|
|
|
|
|
Ok(())
|
|
|
|
}
|
|
|
|
|
2022-12-06 11:38:15 +01:00
|
|
|
/// Clear all databases. Returns a `TransformOutput` with a file that contains the documents
|
|
|
|
/// of the index with the attributes reordered accordingly to the `FieldsIdsMap` given as argument.
|
|
|
|
///
|
2020-11-03 13:20:11 +01:00
|
|
|
// TODO this can be done in parallel by using the rayon `ThreadPool`.
|
2024-05-30 12:14:22 +02:00
|
|
|
#[tracing::instrument(
|
|
|
|
level = "trace"
|
|
|
|
skip(self, wtxn, settings_diff),
|
|
|
|
target = "indexing::documents"
|
|
|
|
)]
|
2022-12-06 11:38:15 +01:00
|
|
|
pub fn prepare_for_documents_reindexing(
|
2020-11-03 13:20:11 +01:00
|
|
|
self,
|
2023-11-22 18:21:19 +01:00
|
|
|
wtxn: &mut heed::RwTxn<'i>,
|
2024-03-26 13:27:43 +01:00
|
|
|
settings_diff: InnerIndexSettingsDiff,
|
2021-06-16 18:33:33 +02:00
|
|
|
) -> Result<TransformOutput> {
|
2021-12-08 14:12:07 +01:00
|
|
|
// There already has been a document addition, the primary key should be set by now.
|
2022-12-19 15:59:22 +01:00
|
|
|
let primary_key = self
|
|
|
|
.index
|
|
|
|
.primary_key(wtxn)?
|
|
|
|
.ok_or(InternalError::DatabaseMissingEntry {
|
|
|
|
db_name: db_name::MAIN,
|
|
|
|
key: Some(main_key::PRIMARY_KEY_KEY),
|
|
|
|
})?
|
|
|
|
.to_string();
|
2021-12-08 14:12:07 +01:00
|
|
|
let field_distribution = self.index.field_distribution(wtxn)?;
|
2022-12-06 11:38:15 +01:00
|
|
|
|
2021-12-08 14:12:07 +01:00
|
|
|
let documents_ids = self.index.documents_ids(wtxn)?;
|
2020-11-03 13:20:11 +01:00
|
|
|
let documents_count = documents_ids.len() as usize;
|
|
|
|
|
2023-11-02 13:37:54 +01:00
|
|
|
// We initialize the sorter with the user indexing settings.
|
2024-05-21 14:53:26 +02:00
|
|
|
let mut original_sorter = if settings_diff.reindex_vectors() {
|
|
|
|
Some(create_sorter(
|
|
|
|
grenad::SortAlgorithm::Stable,
|
2024-08-30 11:49:47 +02:00
|
|
|
KeepFirst,
|
2024-05-21 14:53:26 +02:00
|
|
|
self.indexer_settings.chunk_compression_type,
|
|
|
|
self.indexer_settings.chunk_compression_level,
|
|
|
|
self.indexer_settings.max_nb_chunks,
|
|
|
|
self.indexer_settings.max_memory.map(|mem| mem / 2),
|
2024-10-17 09:30:18 +02:00
|
|
|
true,
|
2024-05-21 14:53:26 +02:00
|
|
|
))
|
|
|
|
} else {
|
|
|
|
None
|
|
|
|
};
|
2022-03-23 17:28:41 +01:00
|
|
|
|
2024-09-23 18:56:15 +02:00
|
|
|
let readers: BTreeMap<&str, (ArroyWrapper, &RoaringBitmap)> = settings_diff
|
2024-06-12 14:05:52 +02:00
|
|
|
.embedding_config_updates
|
|
|
|
.iter()
|
|
|
|
.filter_map(|(name, action)| {
|
2024-09-18 18:13:37 +02:00
|
|
|
if let Some(WriteBackToDocuments { embedder_id, user_provided }) =
|
|
|
|
action.write_back()
|
2024-06-12 14:05:52 +02:00
|
|
|
{
|
2024-09-23 18:56:15 +02:00
|
|
|
let reader = ArroyWrapper::new(
|
|
|
|
self.index.vector_arroy,
|
|
|
|
*embedder_id,
|
|
|
|
action.was_quantized,
|
|
|
|
);
|
|
|
|
Some((name.as_str(), (reader, user_provided)))
|
2024-06-12 14:05:52 +02:00
|
|
|
} else {
|
|
|
|
None
|
|
|
|
}
|
|
|
|
})
|
|
|
|
.collect();
|
|
|
|
|
|
|
|
let old_vectors_fid = settings_diff
|
|
|
|
.old
|
|
|
|
.fields_ids_map
|
|
|
|
.id(crate::vector::parsed_vectors::RESERVED_VECTORS_FIELD_NAME);
|
|
|
|
|
2023-11-02 13:37:54 +01:00
|
|
|
// We initialize the sorter with the user indexing settings.
|
2024-05-21 16:16:36 +02:00
|
|
|
let mut flattened_sorter =
|
|
|
|
if settings_diff.reindex_searchable() || settings_diff.reindex_facets() {
|
|
|
|
Some(create_sorter(
|
|
|
|
grenad::SortAlgorithm::Stable,
|
2024-08-30 11:49:47 +02:00
|
|
|
KeepFirst,
|
2024-05-21 16:16:36 +02:00
|
|
|
self.indexer_settings.chunk_compression_type,
|
|
|
|
self.indexer_settings.chunk_compression_level,
|
|
|
|
self.indexer_settings.max_nb_chunks,
|
|
|
|
self.indexer_settings.max_memory.map(|mem| mem / 2),
|
2024-10-17 09:30:18 +02:00
|
|
|
true,
|
2024-05-21 16:16:36 +02:00
|
|
|
))
|
|
|
|
} else {
|
|
|
|
None
|
|
|
|
};
|
2020-11-03 13:20:11 +01:00
|
|
|
|
2024-05-22 16:05:55 +02:00
|
|
|
if original_sorter.is_some() || flattened_sorter.is_some() {
|
|
|
|
let modified_faceted_fields = settings_diff.modified_faceted_fields();
|
|
|
|
let mut original_obkv_buffer = Vec::new();
|
|
|
|
let mut flattened_obkv_buffer = Vec::new();
|
|
|
|
let mut document_sorter_key_buffer = Vec::new();
|
|
|
|
for result in self.index.external_documents_ids().iter(wtxn)? {
|
|
|
|
let (external_id, docid) = result?;
|
|
|
|
let old_obkv = self.index.documents.get(wtxn, &docid)?.ok_or(
|
|
|
|
InternalError::DatabaseMissingEntry { db_name: db_name::DOCUMENTS, key: None },
|
|
|
|
)?;
|
2020-11-03 13:20:11 +01:00
|
|
|
|
2024-06-12 14:05:52 +02:00
|
|
|
let injected_vectors: std::result::Result<
|
|
|
|
serde_json::Map<String, serde_json::Value>,
|
|
|
|
arroy::Error,
|
|
|
|
> = readers
|
|
|
|
.iter()
|
2024-09-23 18:56:15 +02:00
|
|
|
.filter_map(|(name, (reader, user_provided))| {
|
2024-06-12 14:05:52 +02:00
|
|
|
if !user_provided.contains(docid) {
|
|
|
|
return None;
|
|
|
|
}
|
2024-09-23 18:56:15 +02:00
|
|
|
match reader.item_vectors(wtxn, docid) {
|
|
|
|
Ok(vectors) if vectors.is_empty() => None,
|
|
|
|
Ok(vectors) => Some(Ok((
|
|
|
|
name.to_string(),
|
|
|
|
serde_json::to_value(ExplicitVectors {
|
|
|
|
embeddings: Some(
|
|
|
|
VectorOrArrayOfVectors::from_array_of_vectors(vectors),
|
|
|
|
),
|
|
|
|
regenerate: false,
|
|
|
|
})
|
|
|
|
.unwrap(),
|
|
|
|
))),
|
|
|
|
Err(e) => Some(Err(e)),
|
2024-06-12 14:05:52 +02:00
|
|
|
}
|
|
|
|
})
|
|
|
|
.collect();
|
|
|
|
|
|
|
|
let injected_vectors = injected_vectors?;
|
|
|
|
|
2024-05-22 16:05:55 +02:00
|
|
|
Self::rebind_existing_document(
|
|
|
|
old_obkv,
|
|
|
|
&settings_diff,
|
|
|
|
&modified_faceted_fields,
|
2024-06-12 14:05:52 +02:00
|
|
|
injected_vectors,
|
|
|
|
old_vectors_fid,
|
2024-05-22 16:05:55 +02:00
|
|
|
Some(&mut original_obkv_buffer).filter(|_| original_sorter.is_some()),
|
|
|
|
Some(&mut flattened_obkv_buffer).filter(|_| flattened_sorter.is_some()),
|
|
|
|
)?;
|
2023-11-02 13:37:54 +01:00
|
|
|
|
2024-05-22 16:05:55 +02:00
|
|
|
if let Some(original_sorter) = original_sorter.as_mut() {
|
|
|
|
document_sorter_key_buffer.clear();
|
|
|
|
document_sorter_key_buffer.extend_from_slice(&docid.to_be_bytes());
|
|
|
|
document_sorter_key_buffer.extend_from_slice(external_id.as_bytes());
|
|
|
|
original_sorter.insert(&document_sorter_key_buffer, &original_obkv_buffer)?;
|
|
|
|
}
|
|
|
|
if let Some(flattened_sorter) = flattened_sorter.as_mut() {
|
|
|
|
flattened_sorter.insert(docid.to_be_bytes(), &flattened_obkv_buffer)?;
|
|
|
|
}
|
2024-05-21 16:16:36 +02:00
|
|
|
}
|
2020-11-03 13:20:11 +01:00
|
|
|
}
|
|
|
|
|
2024-06-12 14:05:52 +02:00
|
|
|
// delete all vectors from the embedders that need removal
|
2024-09-23 18:56:15 +02:00
|
|
|
for (_, (reader, _)) in readers {
|
|
|
|
let dimensions = reader.dimensions(wtxn)?;
|
|
|
|
reader.clear(wtxn, dimensions)?;
|
2024-06-12 14:05:52 +02:00
|
|
|
}
|
|
|
|
|
2023-11-02 13:37:54 +01:00
|
|
|
let grenad_params = GrenadParameters {
|
|
|
|
chunk_compression_type: self.indexer_settings.chunk_compression_type,
|
|
|
|
chunk_compression_level: self.indexer_settings.chunk_compression_level,
|
|
|
|
max_memory: self.indexer_settings.max_memory,
|
|
|
|
max_nb_chunks: self.indexer_settings.max_nb_chunks, // default value, may be chosen.
|
|
|
|
};
|
2022-03-23 17:28:41 +01:00
|
|
|
|
2023-11-02 13:37:54 +01:00
|
|
|
// Once we have written all the documents, we merge everything into a Reader.
|
2024-05-21 16:16:36 +02:00
|
|
|
let flattened_documents = match flattened_sorter {
|
|
|
|
Some(flattened_sorter) => Some(sorter_into_reader(flattened_sorter, grenad_params)?),
|
|
|
|
None => None,
|
|
|
|
};
|
2024-05-21 14:53:26 +02:00
|
|
|
let original_documents = match original_sorter {
|
|
|
|
Some(original_sorter) => Some(sorter_into_reader(original_sorter, grenad_params)?),
|
|
|
|
None => None,
|
|
|
|
};
|
2020-11-03 13:20:11 +01:00
|
|
|
|
2024-03-26 13:27:43 +01:00
|
|
|
Ok(TransformOutput {
|
2020-11-03 13:20:11 +01:00
|
|
|
primary_key,
|
2021-06-17 15:16:20 +02:00
|
|
|
field_distribution,
|
2024-03-26 13:27:43 +01:00
|
|
|
settings_diff,
|
2020-11-03 13:20:11 +01:00
|
|
|
documents_count,
|
2024-05-21 14:53:26 +02:00
|
|
|
original_documents: original_documents.map(|od| od.into_inner().into_inner()),
|
2024-05-21 16:16:36 +02:00
|
|
|
flattened_documents: flattened_documents.map(|fd| fd.into_inner().into_inner()),
|
2024-03-26 13:27:43 +01:00
|
|
|
})
|
2020-10-23 14:11:00 +02:00
|
|
|
}
|
|
|
|
}
|
2020-10-31 16:10:15 +01:00
|
|
|
|
2021-08-31 11:44:15 +02:00
|
|
|
/// Drops all the value of type `U` in vec, and reuses the allocation to create a `Vec<T>`.
|
|
|
|
///
|
|
|
|
/// The size and alignment of T and U must match.
|
|
|
|
fn drop_and_reuse<U, T>(mut vec: Vec<U>) -> Vec<T> {
|
|
|
|
debug_assert_eq!(std::mem::align_of::<U>(), std::mem::align_of::<T>());
|
|
|
|
debug_assert_eq!(std::mem::size_of::<U>(), std::mem::size_of::<T>());
|
|
|
|
vec.clear();
|
|
|
|
debug_assert!(vec.is_empty());
|
|
|
|
vec.into_iter().map(|_| unreachable!()).collect()
|
|
|
|
}
|
|
|
|
|
2023-02-14 18:23:57 +01:00
|
|
|
#[cfg(test)]
|
|
|
|
mod test {
|
2024-08-30 11:49:47 +02:00
|
|
|
use grenad::MergeFunction;
|
2024-11-18 11:59:18 +01:00
|
|
|
use obkv::KvReaderU16;
|
2024-08-30 11:49:47 +02:00
|
|
|
|
2023-02-14 18:23:57 +01:00
|
|
|
use super::*;
|
|
|
|
|
|
|
|
#[test]
|
|
|
|
fn merge_obkvs() {
|
2023-10-12 11:46:56 +02:00
|
|
|
let mut additive_doc_0 = Vec::new();
|
|
|
|
let mut deletive_doc_0 = Vec::new();
|
|
|
|
let mut del_add_doc_0 = Vec::new();
|
|
|
|
let mut kv_writer = KvWriter::memory();
|
2023-02-14 18:23:57 +01:00
|
|
|
kv_writer.insert(0_u8, [0]).unwrap();
|
2023-10-12 11:46:56 +02:00
|
|
|
let buffer = kv_writer.into_inner().unwrap();
|
2023-11-20 10:53:40 +01:00
|
|
|
into_del_add_obkv(
|
2024-08-29 19:20:10 +02:00
|
|
|
KvReaderU16::from_slice(&buffer),
|
2023-11-20 10:53:40 +01:00
|
|
|
DelAddOperation::Addition,
|
|
|
|
&mut additive_doc_0,
|
|
|
|
)
|
|
|
|
.unwrap();
|
2023-10-12 11:46:56 +02:00
|
|
|
additive_doc_0.insert(0, Operation::Addition as u8);
|
2023-11-20 10:53:40 +01:00
|
|
|
into_del_add_obkv(
|
2024-08-29 19:20:10 +02:00
|
|
|
KvReaderU16::from_slice(&buffer),
|
2023-11-20 10:53:40 +01:00
|
|
|
DelAddOperation::Deletion,
|
|
|
|
&mut deletive_doc_0,
|
|
|
|
)
|
|
|
|
.unwrap();
|
2023-10-12 11:46:56 +02:00
|
|
|
deletive_doc_0.insert(0, Operation::Deletion as u8);
|
2023-11-20 10:53:40 +01:00
|
|
|
into_del_add_obkv(
|
2024-08-29 19:20:10 +02:00
|
|
|
KvReaderU16::from_slice(&buffer),
|
2023-11-20 10:53:40 +01:00
|
|
|
DelAddOperation::DeletionAndAddition,
|
|
|
|
&mut del_add_doc_0,
|
|
|
|
)
|
|
|
|
.unwrap();
|
2023-10-12 11:46:56 +02:00
|
|
|
del_add_doc_0.insert(0, Operation::Addition as u8);
|
|
|
|
|
|
|
|
let mut additive_doc_1 = Vec::new();
|
|
|
|
let mut kv_writer = KvWriter::memory();
|
|
|
|
kv_writer.insert(1_u8, [1]).unwrap();
|
|
|
|
let buffer = kv_writer.into_inner().unwrap();
|
2023-11-20 10:53:40 +01:00
|
|
|
into_del_add_obkv(
|
2024-08-29 19:20:10 +02:00
|
|
|
KvReaderU16::from_slice(&buffer),
|
2023-11-20 10:53:40 +01:00
|
|
|
DelAddOperation::Addition,
|
|
|
|
&mut additive_doc_1,
|
|
|
|
)
|
|
|
|
.unwrap();
|
2023-10-12 11:46:56 +02:00
|
|
|
additive_doc_1.insert(0, Operation::Addition as u8);
|
|
|
|
|
|
|
|
let mut additive_doc_0_1 = Vec::new();
|
|
|
|
let mut kv_writer = KvWriter::memory();
|
|
|
|
kv_writer.insert(0_u8, [0]).unwrap();
|
|
|
|
kv_writer.insert(1_u8, [1]).unwrap();
|
|
|
|
let buffer = kv_writer.into_inner().unwrap();
|
2023-11-20 10:53:40 +01:00
|
|
|
into_del_add_obkv(
|
2024-08-29 19:20:10 +02:00
|
|
|
KvReaderU16::from_slice(&buffer),
|
2023-11-20 10:53:40 +01:00
|
|
|
DelAddOperation::Addition,
|
|
|
|
&mut additive_doc_0_1,
|
|
|
|
)
|
|
|
|
.unwrap();
|
2023-10-12 11:46:56 +02:00
|
|
|
additive_doc_0_1.insert(0, Operation::Addition as u8);
|
|
|
|
|
2024-08-30 11:49:47 +02:00
|
|
|
let ret = MergeFunction::merge(
|
|
|
|
&ObkvsMergeAdditionsAndDeletions,
|
|
|
|
&[],
|
|
|
|
&[Cow::from(additive_doc_0.as_slice())],
|
|
|
|
)
|
|
|
|
.unwrap();
|
2023-10-12 11:46:56 +02:00
|
|
|
assert_eq!(*ret, additive_doc_0);
|
2023-02-14 18:23:57 +01:00
|
|
|
|
2024-08-30 11:49:47 +02:00
|
|
|
let ret = MergeFunction::merge(
|
|
|
|
&ObkvsMergeAdditionsAndDeletions,
|
2023-10-12 11:46:56 +02:00
|
|
|
&[],
|
|
|
|
&[Cow::from(deletive_doc_0.as_slice()), Cow::from(additive_doc_0.as_slice())],
|
|
|
|
)
|
|
|
|
.unwrap();
|
|
|
|
assert_eq!(*ret, del_add_doc_0);
|
|
|
|
|
2024-08-30 11:49:47 +02:00
|
|
|
let ret = MergeFunction::merge(
|
|
|
|
&ObkvsMergeAdditionsAndDeletions,
|
2023-10-12 11:46:56 +02:00
|
|
|
&[],
|
|
|
|
&[Cow::from(additive_doc_0.as_slice()), Cow::from(deletive_doc_0.as_slice())],
|
|
|
|
)
|
|
|
|
.unwrap();
|
|
|
|
assert_eq!(*ret, deletive_doc_0);
|
|
|
|
|
2024-08-30 11:49:47 +02:00
|
|
|
let ret = MergeFunction::merge(
|
|
|
|
&ObkvsMergeAdditionsAndDeletions,
|
2023-10-12 11:46:56 +02:00
|
|
|
&[],
|
|
|
|
&[
|
|
|
|
Cow::from(additive_doc_1.as_slice()),
|
|
|
|
Cow::from(deletive_doc_0.as_slice()),
|
|
|
|
Cow::from(additive_doc_0.as_slice()),
|
|
|
|
],
|
|
|
|
)
|
|
|
|
.unwrap();
|
|
|
|
assert_eq!(*ret, del_add_doc_0);
|
2023-02-14 18:23:57 +01:00
|
|
|
|
2024-08-30 11:49:47 +02:00
|
|
|
let ret = MergeFunction::merge(
|
|
|
|
&ObkvsMergeAdditionsAndDeletions,
|
2023-02-14 18:23:57 +01:00
|
|
|
&[],
|
2023-10-12 11:46:56 +02:00
|
|
|
&[Cow::from(additive_doc_1.as_slice()), Cow::from(additive_doc_0.as_slice())],
|
2023-02-14 18:23:57 +01:00
|
|
|
)
|
|
|
|
.unwrap();
|
2023-10-12 11:46:56 +02:00
|
|
|
assert_eq!(*ret, additive_doc_0_1);
|
2023-02-14 18:23:57 +01:00
|
|
|
|
2024-08-30 11:49:47 +02:00
|
|
|
let ret = MergeFunction::merge(
|
|
|
|
&ObkvsKeepLastAdditionMergeDeletions,
|
2023-02-14 18:23:57 +01:00
|
|
|
&[],
|
2023-10-12 11:46:56 +02:00
|
|
|
&[Cow::from(additive_doc_1.as_slice()), Cow::from(additive_doc_0.as_slice())],
|
2023-02-14 18:23:57 +01:00
|
|
|
)
|
|
|
|
.unwrap();
|
2023-10-12 11:46:56 +02:00
|
|
|
assert_eq!(*ret, additive_doc_0);
|
2023-02-14 18:23:57 +01:00
|
|
|
|
2024-08-30 11:49:47 +02:00
|
|
|
let ret = MergeFunction::merge(
|
|
|
|
&ObkvsKeepLastAdditionMergeDeletions,
|
2023-02-14 18:23:57 +01:00
|
|
|
&[],
|
|
|
|
&[
|
2023-10-12 11:46:56 +02:00
|
|
|
Cow::from(deletive_doc_0.as_slice()),
|
|
|
|
Cow::from(additive_doc_1.as_slice()),
|
|
|
|
Cow::from(additive_doc_0.as_slice()),
|
2023-02-14 18:23:57 +01:00
|
|
|
],
|
|
|
|
)
|
|
|
|
.unwrap();
|
2023-10-12 11:46:56 +02:00
|
|
|
assert_eq!(*ret, del_add_doc_0);
|
2023-02-14 18:23:57 +01:00
|
|
|
}
|
|
|
|
}
|