mirror of
https://github.com/quickwit-oss/tantivy.git
synced 2026-08-18 12:08:22 +00:00
issue/43 Added delete doc file
This commit is contained in:
@@ -34,6 +34,7 @@ impl IndexMeta {
|
||||
pub struct SegmentMeta {
|
||||
pub segment_id: SegmentId,
|
||||
pub num_docs: u32,
|
||||
pub num_deleted_docs: u32,
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
@@ -42,6 +43,7 @@ impl SegmentMeta {
|
||||
SegmentMeta {
|
||||
segment_id: segment_id,
|
||||
num_docs: num_docs,
|
||||
num_deleted_docs: 0,
|
||||
}
|
||||
}
|
||||
}
|
||||
+12
-5
@@ -65,11 +65,18 @@ impl Segment {
|
||||
/// # Disclaimer
|
||||
/// If deletion of a file fails (e.g. a file
|
||||
/// was read-only.), the method does not
|
||||
/// fail and just logs an error
|
||||
pub fn delete(&self,) {
|
||||
for component in SegmentComponent::values() {
|
||||
let rel_path = self.relative_path(component);
|
||||
if let Err(err) = self.index.directory().delete(&rel_path) {
|
||||
/// fail and just logs an error when it fails.
|
||||
pub fn delete(&self) {
|
||||
info!("Deleting segment {:?}", self.segment_id);
|
||||
let segment_filepaths_res = self.index.directory().ls_starting_with(
|
||||
&*self.segment_id.uuid_string()
|
||||
);
|
||||
if segment_filepaths_res.is_err() {
|
||||
error!("Failed to list files of segment {:?} for deletion.", self.segment_id.uuid_string());
|
||||
return;
|
||||
}
|
||||
for segment_filepath in &segment_filepaths_res.unwrap() {
|
||||
if let Err(err) = self.index.directory().delete(&segment_filepath) {
|
||||
match err {
|
||||
FileError::FileDoesNotExist(_) => {
|
||||
// this is normal behavior.
|
||||
|
||||
@@ -1,5 +1,3 @@
|
||||
use std::vec::IntoIter;
|
||||
|
||||
#[derive(Copy, Clone)]
|
||||
pub enum SegmentComponent {
|
||||
INFO,
|
||||
@@ -9,30 +7,23 @@ pub enum SegmentComponent {
|
||||
FIELDNORMS,
|
||||
TERMS,
|
||||
STORE,
|
||||
DELETE(u64), //< The argument here is an opstamp.
|
||||
// All of the deletes with an opstamp smaller or equal
|
||||
// to this opstamp have been taken in account.
|
||||
}
|
||||
|
||||
impl SegmentComponent {
|
||||
pub fn values() -> IntoIter<SegmentComponent> {
|
||||
vec!(
|
||||
SegmentComponent::INFO,
|
||||
SegmentComponent::POSTINGS,
|
||||
SegmentComponent::POSITIONS,
|
||||
SegmentComponent::FASTFIELDS,
|
||||
SegmentComponent::FIELDNORMS,
|
||||
SegmentComponent::TERMS,
|
||||
SegmentComponent::STORE,
|
||||
).into_iter()
|
||||
}
|
||||
|
||||
pub fn path_suffix(&self)-> &'static str {
|
||||
pub fn path_suffix(&self)-> String {
|
||||
match *self {
|
||||
SegmentComponent::POSITIONS => ".pos",
|
||||
SegmentComponent::INFO => ".info",
|
||||
SegmentComponent::POSTINGS => ".idx",
|
||||
SegmentComponent::TERMS => ".term",
|
||||
SegmentComponent::STORE => ".store",
|
||||
SegmentComponent::FASTFIELDS => ".fast",
|
||||
SegmentComponent::FIELDNORMS => ".fieldnorm",
|
||||
SegmentComponent::POSITIONS => ".pos".to_string(),
|
||||
SegmentComponent::INFO => ".info".to_string(),
|
||||
SegmentComponent::POSTINGS => ".idx".to_string(),
|
||||
SegmentComponent::TERMS => ".term".to_string(),
|
||||
SegmentComponent::STORE => ".store".to_string(),
|
||||
SegmentComponent::FASTFIELDS => ".fast".to_string(),
|
||||
SegmentComponent::FIELDNORMS => ".fieldnorm".to_string(),
|
||||
SegmentComponent::DELETE(opstamp) => format!("{}.del", opstamp)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -50,7 +50,7 @@ impl SegmentId {
|
||||
}
|
||||
|
||||
pub fn relative_path(&self, component: SegmentComponent) -> PathBuf {
|
||||
let filename = self.uuid_string() + component.path_suffix();
|
||||
let filename = self.uuid_string() + &*component.path_suffix();
|
||||
PathBuf::from(filename)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
use std::marker::Send;
|
||||
use std::fmt;
|
||||
use std::path::Path;
|
||||
use std::path::{Path, PathBuf};
|
||||
use directory::error::{FileError, OpenWriteError};
|
||||
use directory::{ReadOnlySource, WritePtr};
|
||||
use std::result;
|
||||
@@ -70,6 +70,10 @@ pub trait Directory: fmt::Debug + Send + Sync + 'static {
|
||||
|
||||
/// Clones the directory and boxes the clone
|
||||
fn box_clone(&self) -> Box<Directory>;
|
||||
|
||||
/// Returns the list of files starting by a given
|
||||
/// prefix.
|
||||
fn ls_starting_with(&self, prefix: &str) -> io::Result<Vec<PathBuf>>;
|
||||
}
|
||||
|
||||
|
||||
|
||||
@@ -4,6 +4,7 @@ use std::collections::HashMap;
|
||||
use std::collections::hash_map::Entry as HashMapEntry;
|
||||
use fst::raw::MmapReadOnly;
|
||||
use std::fs::File;
|
||||
use std::fs::ReadDir;
|
||||
use atomicwrites;
|
||||
use std::sync::RwLock;
|
||||
use std::fmt;
|
||||
@@ -128,8 +129,6 @@ impl Seek for SafeFileWriter {
|
||||
|
||||
impl Directory for MmapDirectory {
|
||||
|
||||
|
||||
|
||||
fn open_read(&self, path: &Path) -> result::Result<ReadOnlySource, FileError> {
|
||||
debug!("Open Read {:?}", path);
|
||||
let full_path = self.resolve_path(path);
|
||||
@@ -203,7 +202,7 @@ impl Directory for MmapDirectory {
|
||||
}
|
||||
|
||||
fn delete(&self, path: &Path) -> result::Result<(), FileError> {
|
||||
debug!("Delete {:?}", path);
|
||||
debug!("Deleting file {:?}", path);
|
||||
let full_path = self.resolve_path(path);
|
||||
let mut mmap_cache = try!(self.mmap_cache
|
||||
.write()
|
||||
@@ -239,4 +238,23 @@ impl Directory for MmapDirectory {
|
||||
Box::new(self.clone())
|
||||
}
|
||||
|
||||
fn ls_starting_with(&self, prefix: &str) -> io::Result<Vec<PathBuf>> {
|
||||
fs::read_dir(&self.root_path)
|
||||
.map(|paths: ReadDir| {
|
||||
paths
|
||||
.filter_map(|dir_entry_res|
|
||||
dir_entry_res
|
||||
.ok()
|
||||
.map(|dir_entry| dir_entry.path())
|
||||
)
|
||||
.filter(|path|
|
||||
path.to_str()
|
||||
.map(|filepath| filepath.starts_with(prefix))
|
||||
.unwrap_or(false)
|
||||
)
|
||||
.map(PathBuf::from)
|
||||
.collect()
|
||||
})
|
||||
|
||||
}
|
||||
}
|
||||
|
||||
@@ -130,6 +130,20 @@ impl InnerDirectory {
|
||||
.contains_key(path)
|
||||
}
|
||||
|
||||
fn ls_starting_with(&self, prefix: &str) -> Vec<PathBuf> {
|
||||
self.0
|
||||
.read()
|
||||
.expect("Failed to get read lock directory.")
|
||||
.keys()
|
||||
.filter(|path: &&PathBuf|
|
||||
path.to_str()
|
||||
.map(|p: &str| p.starts_with(prefix))
|
||||
.unwrap_or(false)
|
||||
)
|
||||
.cloned()
|
||||
.collect()
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
impl fmt::Debug for RAMDirectory {
|
||||
@@ -198,4 +212,9 @@ impl Directory for RAMDirectory {
|
||||
Box::new(self.clone())
|
||||
}
|
||||
|
||||
|
||||
fn ls_starting_with(&self, prefix: &str) -> io::Result<Vec<PathBuf>> {
|
||||
Ok(self.fs.ls_starting_with(prefix))
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -0,0 +1,91 @@
|
||||
use bit_set::BitSet;
|
||||
use directory::WritePtr;
|
||||
use std::io::Write;
|
||||
use std::io;
|
||||
use directory::ReadOnlySource;
|
||||
use DocId;
|
||||
|
||||
pub fn write_delete_bitset(delete_bitset: &BitSet, writer: &mut WritePtr) -> io::Result<()> {
|
||||
let max_doc = delete_bitset.capacity();
|
||||
let mut byte = 0u8;
|
||||
let mut shift = 0u8;
|
||||
for doc in 0..max_doc {
|
||||
if delete_bitset.contains(doc) {
|
||||
byte |= 1 << shift;
|
||||
}
|
||||
if shift == 7 {
|
||||
writer.write(&[byte])?;
|
||||
shift = 0;
|
||||
byte = 0;
|
||||
}
|
||||
else {
|
||||
shift += 1;
|
||||
}
|
||||
}
|
||||
if max_doc % 8 > 0 {
|
||||
writer.write(&[byte])?;
|
||||
}
|
||||
writer.flush()
|
||||
}
|
||||
|
||||
pub struct DeleteBitSet(ReadOnlySource);
|
||||
|
||||
impl DeleteBitSet {
|
||||
|
||||
pub fn open(data: ReadOnlySource) -> DeleteBitSet {
|
||||
DeleteBitSet(data)
|
||||
}
|
||||
|
||||
pub fn is_deleted(&self, doc: DocId) -> bool {
|
||||
let byte_offset = doc / 8u32;
|
||||
let b: u8 = (*self.0)[byte_offset as usize];
|
||||
let shift = (doc & 7u32) as u8;
|
||||
b & (1u8 << shift) != 0
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use std::path::PathBuf;
|
||||
use bit_set::BitSet;
|
||||
use directory::*;
|
||||
use super::*;
|
||||
|
||||
fn test_delete_bitset_helper(bitset: &BitSet) {
|
||||
let test_path = PathBuf::from("test");
|
||||
let mut directory = RAMDirectory::create();
|
||||
{
|
||||
let mut writer = directory.open_write(&*test_path).unwrap();
|
||||
write_delete_bitset(bitset, &mut writer).unwrap();
|
||||
}
|
||||
{
|
||||
let source = directory.open_read(&test_path).unwrap();
|
||||
let delete_bitset = DeleteBitSet::open(source);
|
||||
let n = bitset.capacity();
|
||||
for doc in 0..n {
|
||||
assert_eq!(bitset.contains(doc), delete_bitset.is_deleted(doc as DocId));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_delete_bitset() {
|
||||
{
|
||||
let mut bitset = BitSet::with_capacity(10);
|
||||
bitset.insert(1);
|
||||
bitset.insert(9);
|
||||
test_delete_bitset_helper(&bitset);
|
||||
}
|
||||
{
|
||||
let mut bitset = BitSet::with_capacity(8);
|
||||
bitset.insert(1);
|
||||
bitset.insert(2);
|
||||
bitset.insert(3);
|
||||
bitset.insert(5);
|
||||
bitset.insert(7);
|
||||
test_delete_bitset_helper(&bitset);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -13,6 +13,7 @@
|
||||
mod reader;
|
||||
mod writer;
|
||||
mod serializer;
|
||||
pub mod delete;
|
||||
|
||||
pub use self::writer::{U32FastFieldsWriter, U32FastFieldWriter};
|
||||
pub use self::reader::{U32FastFieldsReader, U32FastFieldReader};
|
||||
|
||||
+23
-13
@@ -9,9 +9,11 @@ use schema::Term;
|
||||
use std::thread::JoinHandle;
|
||||
use indexer::{MergePolicy, DefaultMergePolicy};
|
||||
use indexer::SegmentWriter;
|
||||
use core::SegmentComponent;
|
||||
use super::directory_lock::DirectoryLock;
|
||||
use std::clone::Clone;
|
||||
use std::io;
|
||||
use fastfield::delete;
|
||||
use std::thread;
|
||||
use std::mem;
|
||||
use indexer::merger::IndexMerger;
|
||||
@@ -87,7 +89,7 @@ impl !Sync for IndexWriter {}
|
||||
|
||||
|
||||
fn index_documents(heap: &mut Heap,
|
||||
segment: Segment,
|
||||
mut segment: Segment,
|
||||
schema: &Schema,
|
||||
document_iterator: &mut Iterator<Item=AddOperation>,
|
||||
segment_update_sender: &mut SegmentUpdateSender,
|
||||
@@ -95,7 +97,7 @@ fn index_documents(heap: &mut Heap,
|
||||
-> Result<()> {
|
||||
heap.clear();
|
||||
let segment_id = segment.id();
|
||||
let mut segment_writer = try!(SegmentWriter::for_segment(heap, segment, &schema));
|
||||
let mut segment_writer = try!(SegmentWriter::for_segment(heap, segment.clone(), &schema));
|
||||
for doc in document_iterator {
|
||||
try!(segment_writer.add_document(&doc, &schema));
|
||||
if segment_writer.is_buffer_full() {
|
||||
@@ -105,23 +107,29 @@ fn index_documents(heap: &mut Heap,
|
||||
}
|
||||
}
|
||||
let num_docs = segment_writer.max_doc();
|
||||
assert!(num_docs > 0);
|
||||
let first_opstamp: u64 = segment_writer.first_opstamp();
|
||||
let last_opstamp: u64 = segment_writer.last_opstamp();
|
||||
|
||||
delete_cursor.skip_to(first_opstamp);
|
||||
|
||||
let delete_cursor_clone = delete_cursor.clone();
|
||||
assert!(num_docs > 0);
|
||||
|
||||
let deleted_docset_opt = segment_writer.compute_deleted_bitset(delete_cursor);
|
||||
|
||||
let last_opstamp = segment_writer.last_opstamp();
|
||||
|
||||
let num_deleted_docs;
|
||||
|
||||
if let Some(deleted_docset) = deleted_docset_opt {
|
||||
let mut delete_write = segment.open_write(SegmentComponent::DELETE(last_opstamp))?;
|
||||
delete::write_delete_bitset(&deleted_docset, &mut delete_write)?;
|
||||
num_deleted_docs = deleted_docset.len();
|
||||
}
|
||||
else {
|
||||
num_deleted_docs = 0;
|
||||
}
|
||||
|
||||
let doc_mapping = segment_writer.compute_doc_mapping_after_delete(delete_cursor_clone);
|
||||
|
||||
let segment_meta = SegmentMeta {
|
||||
segment_id: segment_id,
|
||||
num_docs: num_docs,
|
||||
num_deleted_docs: num_deleted_docs as u32,
|
||||
};
|
||||
|
||||
delete_cursor.skip_to(last_opstamp);
|
||||
|
||||
try!(segment_writer.finalize());
|
||||
segment_update_sender.send(SegmentUpdate::AddSegment(segment_meta));
|
||||
Ok(())
|
||||
@@ -330,9 +338,11 @@ impl IndexWriter {
|
||||
let num_docs = try!(merger.write(segment_serializer));
|
||||
let merged_segment_ids: Vec<SegmentId> =
|
||||
segments.iter().map(|segment| segment.id()).collect();
|
||||
|
||||
let segment_meta = SegmentMeta {
|
||||
segment_id: merged_segment.id(),
|
||||
num_docs: num_docs,
|
||||
num_deleted_docs: 0,
|
||||
};
|
||||
|
||||
segment_manager.end_merge(&merged_segment_ids, &segment_meta);
|
||||
|
||||
@@ -221,6 +221,7 @@ impl SegmentUpdater {
|
||||
let segment_meta = SegmentMeta {
|
||||
segment_id: merged_segment.id(),
|
||||
num_docs: num_docs,
|
||||
num_deleted_docs: 0u32,
|
||||
};
|
||||
let segment_update = SegmentUpdate::EndMerge(merging_thread_id, segment_ids.clone(), segment_meta.clone());
|
||||
segment_update_sender_clone.send(segment_update.clone());
|
||||
|
||||
@@ -156,25 +156,25 @@ impl<'a> SegmentWriter<'a> {
|
||||
doc_id as DocId
|
||||
}
|
||||
|
||||
pub fn compute_doc_mapping_after_delete(&self, mut delete_queue_cursor: DeleteQueueCursor) -> Vec<Option<DocId>> {
|
||||
let delete_docs = self.compute_delete_mask(&mut delete_queue_cursor);
|
||||
let max_doc: usize = self.max_doc as usize;
|
||||
let mut doc_autoinc = 0u32;
|
||||
(0..max_doc)
|
||||
.map(|doc| {
|
||||
if delete_docs.contains(doc) {
|
||||
None
|
||||
}
|
||||
else {
|
||||
let new_doc = doc_autoinc;
|
||||
doc_autoinc += 1;
|
||||
Some(new_doc)
|
||||
}
|
||||
})
|
||||
.collect::<Vec<_>>()
|
||||
}
|
||||
// pub fn compute_doc_mapping_after_delete(&self, mut delete_queue_cursor: DeleteQueueCursor) -> Vec<Option<DocId>> {
|
||||
// let delete_docs = self.compute_delete_mask(&mut delete_queue_cursor);
|
||||
// let max_doc: usize = self.max_doc as usize;
|
||||
// let mut doc_autoinc = 0u32;
|
||||
// (0..max_doc)
|
||||
// .map(|doc| {
|
||||
// if delete_docs.contains(doc) {
|
||||
// None
|
||||
// }
|
||||
// else {
|
||||
// let new_doc = doc_autoinc;
|
||||
// doc_autoinc += 1;
|
||||
// Some(new_doc)
|
||||
// }
|
||||
// })
|
||||
// .collect::<Vec<_>>()
|
||||
// }
|
||||
|
||||
pub fn first_opstamp(&self) -> u64 {
|
||||
fn first_opstamp(&self) -> u64 {
|
||||
*(self.doc_opstamps
|
||||
.first()
|
||||
.expect("Last doc opstamp called on an empty segment writer"))
|
||||
@@ -186,17 +186,21 @@ impl<'a> SegmentWriter<'a> {
|
||||
.expect("Last doc opstamp called on an empty segment writer"))
|
||||
}
|
||||
|
||||
fn compute_delete_mask(&self, delete_queue_cursor: &mut DeleteQueueCursor) -> BitSet {
|
||||
if let Some(min_opstamp) = self.doc_opstamps.first() {
|
||||
if !delete_queue_cursor.skip_to(*min_opstamp) {
|
||||
return BitSet::new();
|
||||
pub fn compute_deleted_bitset(&self, delete_queue_cursor: &mut DeleteQueueCursor) -> Option<BitSet> {
|
||||
if let Some(first_opstamp) = self.doc_opstamps.first() {
|
||||
if !delete_queue_cursor.skip_to(*first_opstamp) {
|
||||
return None;
|
||||
}
|
||||
}
|
||||
else {
|
||||
return BitSet::new();
|
||||
return None;
|
||||
}
|
||||
let last_opstamp = *self.doc_opstamps.last().unwrap();
|
||||
let mut deleted_docs = BitSet::with_capacity(self.max_doc as usize);
|
||||
while let Some(delete_operation) = delete_queue_cursor.consume() {
|
||||
while let Some(delete_operation) = delete_queue_cursor.peek() {
|
||||
if delete_operation.opstamp > last_opstamp {
|
||||
break;
|
||||
}
|
||||
// We can skip computing delete operations that
|
||||
// are older than our oldest document.
|
||||
//
|
||||
@@ -210,8 +214,9 @@ impl<'a> SegmentWriter<'a> {
|
||||
deleted_docs: &mut deleted_docs
|
||||
};
|
||||
postings_writer.push_documents(delete_term.value(), &mut document_deleter);
|
||||
delete_queue_cursor.consume();
|
||||
}
|
||||
deleted_docs
|
||||
Some(deleted_docs)
|
||||
}
|
||||
|
||||
|
||||
|
||||
Reference in New Issue
Block a user