mirror of
https://github.com/meilisearch/meilisearch.git
synced 2024-11-22 10:07:40 +08:00
Change the way the FST is built
This commit is contained in:
parent
2d1caf27df
commit
6e87332410
@ -2,7 +2,7 @@ use std::fs::File;
|
|||||||
use std::io::{self, BufWriter};
|
use std::io::{self, BufWriter};
|
||||||
|
|
||||||
use bincode::ErrorKind;
|
use bincode::ErrorKind;
|
||||||
use fst::{Set, SetBuilder};
|
use fst::{Set, SetBuilder, Streamer};
|
||||||
use grenad::Merger;
|
use grenad::Merger;
|
||||||
use heed::types::Bytes;
|
use heed::types::Bytes;
|
||||||
use heed::{BoxedError, Database, RoTxn};
|
use heed::{BoxedError, Database, RoTxn};
|
||||||
@ -11,8 +11,8 @@ use roaring::RoaringBitmap;
|
|||||||
use tempfile::tempfile;
|
use tempfile::tempfile;
|
||||||
|
|
||||||
use super::channel::*;
|
use super::channel::*;
|
||||||
use super::{Deletion, DocumentChange, Insertion, KvReaderDelAdd, KvReaderFieldId, Update};
|
|
||||||
use super::extract::FacetKind;
|
use super::extract::FacetKind;
|
||||||
|
use super::{Deletion, DocumentChange, Insertion, KvReaderDelAdd, KvReaderFieldId, Update};
|
||||||
use crate::update::del_add::DelAdd;
|
use crate::update::del_add::DelAdd;
|
||||||
use crate::update::new::channel::MergerOperation;
|
use crate::update::new::channel::MergerOperation;
|
||||||
use crate::update::MergeDeladdCboRoaringBitmaps;
|
use crate::update::MergeDeladdCboRoaringBitmaps;
|
||||||
@ -29,7 +29,7 @@ pub fn merge_grenad_entries(
|
|||||||
index: &Index,
|
index: &Index,
|
||||||
mut global_fields_ids_map: GlobalFieldsIdsMap<'_>,
|
mut global_fields_ids_map: GlobalFieldsIdsMap<'_>,
|
||||||
) -> Result<()> {
|
) -> Result<()> {
|
||||||
let mut buffer = Vec::new();
|
let mut buffer: Vec<u8> = Vec::new();
|
||||||
let mut documents_ids = index.documents_ids(rtxn)?;
|
let mut documents_ids = index.documents_ids(rtxn)?;
|
||||||
let mut geo_extractor = GeoExtractor::new(rtxn, index)?;
|
let mut geo_extractor = GeoExtractor::new(rtxn, index)?;
|
||||||
|
|
||||||
@ -46,8 +46,7 @@ pub fn merge_grenad_entries(
|
|||||||
rtxn,
|
rtxn,
|
||||||
&mut buffer,
|
&mut buffer,
|
||||||
sender.docids::<ExactWordDocids>(),
|
sender.docids::<ExactWordDocids>(),
|
||||||
|_key| Ok(()),
|
|_, _key| Ok(()),
|
||||||
|_key| Ok(()),
|
|
||||||
)?;
|
)?;
|
||||||
}
|
}
|
||||||
MergerOperation::FidWordCountDocidsMerger(merger) => {
|
MergerOperation::FidWordCountDocidsMerger(merger) => {
|
||||||
@ -59,13 +58,12 @@ pub fn merge_grenad_entries(
|
|||||||
rtxn,
|
rtxn,
|
||||||
&mut buffer,
|
&mut buffer,
|
||||||
sender.docids::<FidWordCountDocids>(),
|
sender.docids::<FidWordCountDocids>(),
|
||||||
|_key| Ok(()),
|
|_, _key| Ok(()),
|
||||||
|_key| Ok(()),
|
|
||||||
)?;
|
)?;
|
||||||
}
|
}
|
||||||
MergerOperation::WordDocidsMerger(merger) => {
|
MergerOperation::WordDocidsMerger(merger) => {
|
||||||
let mut add_words_fst = SetBuilder::new(BufWriter::new(tempfile()?))?;
|
let words_fst = index.words_fst(rtxn)?;
|
||||||
let mut del_words_fst = SetBuilder::new(BufWriter::new(tempfile()?))?;
|
let mut word_fst_builder = WordFstBuilder::new(&words_fst, 4)?;
|
||||||
{
|
{
|
||||||
let span =
|
let span =
|
||||||
tracing::trace_span!(target: "indexing::documents::merge", "word_docids");
|
tracing::trace_span!(target: "indexing::documents::merge", "word_docids");
|
||||||
@ -77,8 +75,7 @@ pub fn merge_grenad_entries(
|
|||||||
rtxn,
|
rtxn,
|
||||||
&mut buffer,
|
&mut buffer,
|
||||||
sender.docids::<WordDocids>(),
|
sender.docids::<WordDocids>(),
|
||||||
|key| add_words_fst.insert(key),
|
|deladd, key| word_fst_builder.register_word(deladd, key),
|
||||||
|key| del_words_fst.insert(key),
|
|
||||||
)?;
|
)?;
|
||||||
}
|
}
|
||||||
|
|
||||||
@ -86,9 +83,8 @@ pub fn merge_grenad_entries(
|
|||||||
let span =
|
let span =
|
||||||
tracing::trace_span!(target: "indexing::documents::merge", "words_fst");
|
tracing::trace_span!(target: "indexing::documents::merge", "words_fst");
|
||||||
let _entered = span.enter();
|
let _entered = span.enter();
|
||||||
// Move that into a dedicated function
|
|
||||||
let words_fst = index.words_fst(rtxn)?;
|
let mmap = word_fst_builder.build()?;
|
||||||
let mmap = compute_new_words_fst(add_words_fst, del_words_fst, words_fst)?;
|
|
||||||
sender.main().write_words_fst(mmap).unwrap();
|
sender.main().write_words_fst(mmap).unwrap();
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@ -102,8 +98,7 @@ pub fn merge_grenad_entries(
|
|||||||
rtxn,
|
rtxn,
|
||||||
&mut buffer,
|
&mut buffer,
|
||||||
sender.docids::<WordFidDocids>(),
|
sender.docids::<WordFidDocids>(),
|
||||||
|_key| Ok(()),
|
|_, _key| Ok(()),
|
||||||
|_key| Ok(()),
|
|
||||||
)?;
|
)?;
|
||||||
}
|
}
|
||||||
MergerOperation::WordPairProximityDocidsMerger(merger) => {
|
MergerOperation::WordPairProximityDocidsMerger(merger) => {
|
||||||
@ -115,8 +110,7 @@ pub fn merge_grenad_entries(
|
|||||||
rtxn,
|
rtxn,
|
||||||
&mut buffer,
|
&mut buffer,
|
||||||
sender.docids::<WordPairProximityDocids>(),
|
sender.docids::<WordPairProximityDocids>(),
|
||||||
|_key| Ok(()),
|
|_, _key| Ok(()),
|
||||||
|_key| Ok(()),
|
|
||||||
)?;
|
)?;
|
||||||
}
|
}
|
||||||
MergerOperation::WordPositionDocidsMerger(merger) => {
|
MergerOperation::WordPositionDocidsMerger(merger) => {
|
||||||
@ -128,8 +122,7 @@ pub fn merge_grenad_entries(
|
|||||||
rtxn,
|
rtxn,
|
||||||
&mut buffer,
|
&mut buffer,
|
||||||
sender.docids::<WordPositionDocids>(),
|
sender.docids::<WordPositionDocids>(),
|
||||||
|_key| Ok(()),
|
|_, _key| Ok(()),
|
||||||
|_key| Ok(()),
|
|
||||||
)?;
|
)?;
|
||||||
}
|
}
|
||||||
MergerOperation::InsertDocument { docid, document } => {
|
MergerOperation::InsertDocument { docid, document } => {
|
||||||
@ -199,6 +192,142 @@ pub fn merge_grenad_entries(
|
|||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
|
struct WordFstBuilder<'a> {
|
||||||
|
stream: fst::set::Stream<'a>,
|
||||||
|
word_fst_builder: SetBuilder<BufWriter<File>>,
|
||||||
|
prefix_fst_builders: Vec<SetBuilder<BufWriter<File>>>,
|
||||||
|
max_prefix_length: usize,
|
||||||
|
last_word: Vec<u8>,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl<'a> WordFstBuilder<'a> {
|
||||||
|
pub fn new(
|
||||||
|
words_fst: &'a Set<std::borrow::Cow<'a, [u8]>>,
|
||||||
|
max_prefix_length: usize,
|
||||||
|
) -> Result<Self> {
|
||||||
|
let mut prefix_fst_builders = Vec::new();
|
||||||
|
for _ in 0..max_prefix_length {
|
||||||
|
prefix_fst_builders.push(SetBuilder::new(BufWriter::new(tempfile()?))?);
|
||||||
|
}
|
||||||
|
|
||||||
|
Ok(Self {
|
||||||
|
stream: words_fst.stream(),
|
||||||
|
word_fst_builder: SetBuilder::new(BufWriter::new(tempfile()?))?,
|
||||||
|
prefix_fst_builders,
|
||||||
|
max_prefix_length,
|
||||||
|
last_word: Vec::new(),
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
pub fn register_word(&mut self, deladd: DelAdd, key: &[u8]) -> Result<()> {
|
||||||
|
match deladd {
|
||||||
|
DelAdd::Addition => self.add_word(key),
|
||||||
|
DelAdd::Deletion => self.del_word(key),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
pub fn add_word(&mut self, word: &[u8]) -> Result<()> {
|
||||||
|
if !self.last_word.is_empty() {
|
||||||
|
let next = self.last_word.as_slice();
|
||||||
|
match next.cmp(word) {
|
||||||
|
std::cmp::Ordering::Less => {
|
||||||
|
// We need to insert the last word from the current fst
|
||||||
|
self.word_fst_builder.insert(next)?;
|
||||||
|
self.last_word.clear();
|
||||||
|
}
|
||||||
|
std::cmp::Ordering::Equal => {
|
||||||
|
// We insert the word and drop the last word
|
||||||
|
self.word_fst_builder.insert(next)?;
|
||||||
|
self.last_word.clear();
|
||||||
|
return Ok(());
|
||||||
|
}
|
||||||
|
std::cmp::Ordering::Greater => {
|
||||||
|
// We insert the word and keep the last word
|
||||||
|
self.word_fst_builder.insert(word)?;
|
||||||
|
|
||||||
|
return Ok(());
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
while let Some(next) = self.stream.next() {
|
||||||
|
match next.cmp(word) {
|
||||||
|
std::cmp::Ordering::Less => {
|
||||||
|
// We need to insert the last word from the current fst
|
||||||
|
self.word_fst_builder.insert(next)?;
|
||||||
|
}
|
||||||
|
std::cmp::Ordering::Equal => {
|
||||||
|
// We insert the word
|
||||||
|
self.word_fst_builder.insert(next)?;
|
||||||
|
|
||||||
|
return Ok(());
|
||||||
|
}
|
||||||
|
std::cmp::Ordering::Greater => {
|
||||||
|
// We insert the word and keep the last word
|
||||||
|
self.word_fst_builder.insert(word)?;
|
||||||
|
self.last_word.clear();
|
||||||
|
self.last_word.extend_from_slice(next);
|
||||||
|
|
||||||
|
return Ok(());
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
||||||
|
pub fn del_word(&mut self, word: &[u8]) -> Result<()> {
|
||||||
|
if !self.last_word.is_empty() {
|
||||||
|
let next = self.last_word.as_slice();
|
||||||
|
match next.cmp(word) {
|
||||||
|
std::cmp::Ordering::Less => {
|
||||||
|
// We insert the word from the current fst because the next word to delete is greater
|
||||||
|
self.word_fst_builder.insert(next)?;
|
||||||
|
self.last_word.clear();
|
||||||
|
}
|
||||||
|
std::cmp::Ordering::Equal => {
|
||||||
|
// We delete the word by not inserting it in the new fst and drop the last word
|
||||||
|
self.last_word.clear();
|
||||||
|
return Ok(());
|
||||||
|
}
|
||||||
|
std::cmp::Ordering::Greater => {
|
||||||
|
// keep the current word until the next word to delete is greater or equal
|
||||||
|
return Ok(());
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
while let Some(next) = self.stream.next() {
|
||||||
|
match next.cmp(word) {
|
||||||
|
std::cmp::Ordering::Less => {
|
||||||
|
// We insert the word from the current fst because the next word to delete is greater
|
||||||
|
self.word_fst_builder.insert(next)?;
|
||||||
|
}
|
||||||
|
std::cmp::Ordering::Equal => {
|
||||||
|
// We delete the word by not inserting it in the new fst and drop the last word
|
||||||
|
return Ok(());
|
||||||
|
}
|
||||||
|
std::cmp::Ordering::Greater => {
|
||||||
|
// keep the current word until the next word to delete is greater or equal
|
||||||
|
self.last_word.clear();
|
||||||
|
self.last_word.extend_from_slice(next);
|
||||||
|
|
||||||
|
return Ok(());
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
||||||
|
pub fn build(mut self) -> Result<Mmap> {
|
||||||
|
let words_fst_file = self.word_fst_builder.into_inner()?.into_inner().unwrap();
|
||||||
|
let words_fst_mmap = unsafe { Mmap::map(&words_fst_file)? };
|
||||||
|
|
||||||
|
Ok(words_fst_mmap)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
pub struct GeoExtractor {
|
pub struct GeoExtractor {
|
||||||
rtree: Option<rstar::RTree<GeoPoint>>,
|
rtree: Option<rstar::RTree<GeoPoint>>,
|
||||||
}
|
}
|
||||||
@ -247,30 +376,6 @@ impl GeoExtractor {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
fn compute_new_words_fst(
|
|
||||||
add_words_fst: SetBuilder<BufWriter<File>>,
|
|
||||||
del_words_fst: SetBuilder<BufWriter<File>>,
|
|
||||||
words_fst: Set<std::borrow::Cow<'_, [u8]>>,
|
|
||||||
) -> Result<Mmap> {
|
|
||||||
let add_words_fst_file = add_words_fst.into_inner()?;
|
|
||||||
let add_words_fst_mmap = unsafe { Mmap::map(&add_words_fst_file.into_inner().unwrap())? };
|
|
||||||
let add_words_fst = Set::new(&add_words_fst_mmap)?;
|
|
||||||
|
|
||||||
let del_words_fst_file = del_words_fst.into_inner()?;
|
|
||||||
let del_words_fst_mmap = unsafe { Mmap::map(&del_words_fst_file.into_inner().unwrap())? };
|
|
||||||
let del_words_fst = Set::new(&del_words_fst_mmap)?;
|
|
||||||
|
|
||||||
let diff = words_fst.op().add(&del_words_fst).difference();
|
|
||||||
let stream = add_words_fst.op().add(diff).union();
|
|
||||||
|
|
||||||
let mut words_fst = SetBuilder::new(tempfile()?)?;
|
|
||||||
words_fst.extend_stream(stream)?;
|
|
||||||
let words_fst_file = words_fst.into_inner()?;
|
|
||||||
let words_fst_mmap = unsafe { Mmap::map(&words_fst_file)? };
|
|
||||||
|
|
||||||
Ok(words_fst_mmap)
|
|
||||||
}
|
|
||||||
|
|
||||||
#[tracing::instrument(level = "trace", skip_all, target = "indexing::merge")]
|
#[tracing::instrument(level = "trace", skip_all, target = "indexing::merge")]
|
||||||
fn merge_and_send_docids(
|
fn merge_and_send_docids(
|
||||||
merger: Merger<File, MergeDeladdCboRoaringBitmaps>,
|
merger: Merger<File, MergeDeladdCboRoaringBitmaps>,
|
||||||
@ -278,8 +383,7 @@ fn merge_and_send_docids(
|
|||||||
rtxn: &RoTxn<'_>,
|
rtxn: &RoTxn<'_>,
|
||||||
buffer: &mut Vec<u8>,
|
buffer: &mut Vec<u8>,
|
||||||
docids_sender: impl DocidsSender,
|
docids_sender: impl DocidsSender,
|
||||||
mut add_key: impl FnMut(&[u8]) -> fst::Result<()>,
|
mut register_key: impl FnMut(DelAdd, &[u8]) -> Result<()>,
|
||||||
mut del_key: impl FnMut(&[u8]) -> fst::Result<()>,
|
|
||||||
) -> Result<()> {
|
) -> Result<()> {
|
||||||
let mut merger_iter = merger.into_stream_merger_iter().unwrap();
|
let mut merger_iter = merger.into_stream_merger_iter().unwrap();
|
||||||
while let Some((key, deladd)) = merger_iter.next().unwrap() {
|
while let Some((key, deladd)) = merger_iter.next().unwrap() {
|
||||||
@ -292,11 +396,11 @@ fn merge_and_send_docids(
|
|||||||
Operation::Write(bitmap) => {
|
Operation::Write(bitmap) => {
|
||||||
let value = cbo_bitmap_serialize_into_vec(&bitmap, buffer);
|
let value = cbo_bitmap_serialize_into_vec(&bitmap, buffer);
|
||||||
docids_sender.write(key, value).unwrap();
|
docids_sender.write(key, value).unwrap();
|
||||||
add_key(key)?;
|
register_key(DelAdd::Addition, key)?;
|
||||||
}
|
}
|
||||||
Operation::Delete => {
|
Operation::Delete => {
|
||||||
docids_sender.delete(key).unwrap();
|
docids_sender.delete(key).unwrap();
|
||||||
del_key(key)?;
|
register_key(DelAdd::Deletion, key)?;
|
||||||
}
|
}
|
||||||
Operation::Ignore => (),
|
Operation::Ignore => (),
|
||||||
}
|
}
|
||||||
|
Loading…
Reference in New Issue
Block a user