From fba44b78b627241c70cacc40ab7fcd5fea5891a4 Mon Sep 17 00:00:00 2001 From: Paul Masurel Date: Thu, 19 Jan 2017 23:28:57 +0900 Subject: [PATCH] issue/43 Added delete doc file --- src/core/index_meta.rs | 2 + src/core/segment.rs | 17 ++++-- src/core/segment_component.rs | 33 +++++------- src/core/segment_id.rs | 2 +- src/directory/directory.rs | 6 ++- src/directory/mmap_directory.rs | 24 +++++++-- src/directory/ram_directory.rs | 19 +++++++ src/fastfield/delete.rs | 91 +++++++++++++++++++++++++++++++++ src/fastfield/mod.rs | 1 + src/indexer/index_writer.rs | 36 ++++++++----- src/indexer/segment_updater.rs | 1 + src/indexer/segment_writer.rs | 55 +++++++++++--------- 12 files changed, 218 insertions(+), 69 deletions(-) create mode 100644 src/fastfield/delete.rs diff --git a/src/core/index_meta.rs b/src/core/index_meta.rs index a2623f9d0..d82f865f2 100644 --- a/src/core/index_meta.rs +++ b/src/core/index_meta.rs @@ -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, } } } \ No newline at end of file diff --git a/src/core/segment.rs b/src/core/segment.rs index 3e8bc9a42..354ddb34a 100644 --- a/src/core/segment.rs +++ b/src/core/segment.rs @@ -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. diff --git a/src/core/segment_component.rs b/src/core/segment_component.rs index a55ea19dc..62610af3f 100644 --- a/src/core/segment_component.rs +++ b/src/core/segment_component.rs @@ -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 { - 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) } } } diff --git a/src/core/segment_id.rs b/src/core/segment_id.rs index 3d77668b3..a9916cb83 100644 --- a/src/core/segment_id.rs +++ b/src/core/segment_id.rs @@ -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) } } diff --git a/src/directory/directory.rs b/src/directory/directory.rs index 6171ed606..1b29592b3 100644 --- a/src/directory/directory.rs +++ b/src/directory/directory.rs @@ -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; + + /// Returns the list of files starting by a given + /// prefix. + fn ls_starting_with(&self, prefix: &str) -> io::Result>; } diff --git a/src/directory/mmap_directory.rs b/src/directory/mmap_directory.rs index e4e595fec..b905bd7f5 100644 --- a/src/directory/mmap_directory.rs +++ b/src/directory/mmap_directory.rs @@ -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 { 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> { + 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() + }) + + } } diff --git a/src/directory/ram_directory.rs b/src/directory/ram_directory.rs index 6ced0ae7e..cb3c2757a 100644 --- a/src/directory/ram_directory.rs +++ b/src/directory/ram_directory.rs @@ -130,6 +130,20 @@ impl InnerDirectory { .contains_key(path) } + fn ls_starting_with(&self, prefix: &str) -> Vec { + 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> { + Ok(self.fs.ls_starting_with(prefix)) + } + } diff --git a/src/fastfield/delete.rs b/src/fastfield/delete.rs new file mode 100644 index 000000000..e700bc3e7 --- /dev/null +++ b/src/fastfield/delete.rs @@ -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); + } + } +} \ No newline at end of file diff --git a/src/fastfield/mod.rs b/src/fastfield/mod.rs index b51a5d15b..00de208b9 100644 --- a/src/fastfield/mod.rs +++ b/src/fastfield/mod.rs @@ -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}; diff --git a/src/indexer/index_writer.rs b/src/indexer/index_writer.rs index bc59efcb9..07aa2bd82 100644 --- a/src/indexer/index_writer.rs +++ b/src/indexer/index_writer.rs @@ -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, 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 = 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); diff --git a/src/indexer/segment_updater.rs b/src/indexer/segment_updater.rs index 2a1b3e577..f7141d46d 100644 --- a/src/indexer/segment_updater.rs +++ b/src/indexer/segment_updater.rs @@ -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()); diff --git a/src/indexer/segment_writer.rs b/src/indexer/segment_writer.rs index 8a7930a86..074f6e7d1 100644 --- a/src/indexer/segment_writer.rs +++ b/src/indexer/segment_writer.rs @@ -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> { - 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::>() - } + // pub fn compute_doc_mapping_after_delete(&self, mut delete_queue_cursor: DeleteQueueCursor) -> Vec> { + // 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::>() + // } - 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 { + 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) }