diff options
Diffstat (limited to 'log_db/src')
| -rw-r--r-- | log_db/src/common.rs | 4 | ||||
| -rw-r--r-- | log_db/src/lib.rs | 139 | ||||
| -rw-r--r-- | log_db/src/log_reader_reverse.rs | 118 | ||||
| -rw-r--r-- | log_db/src/memtable_primary.rs | 8 | ||||
| -rw-r--r-- | log_db/src/memtable_secondary.rs | 32 |
5 files changed, 91 insertions, 210 deletions
diff --git a/log_db/src/common.rs b/log_db/src/common.rs index 48c85ea..2adca5d 100644 --- a/log_db/src/common.rs +++ b/log_db/src/common.rs @@ -99,9 +99,9 @@ impl PartialOrd for LogKeySet { impl LogKeySet { /// Create a new LogKeySet with an initial LogKey. /// The initial LogKey is required since LogKeySet must be non-empty. - pub fn new_with_initial(key: &LogKey) -> Self { + pub fn new_with_initial(key: LogKey) -> Self { let mut set = HashSet::new(); - set.insert(key.clone()); + set.insert(key); LogKeySet { set } } diff --git a/log_db/src/lib.rs b/log_db/src/lib.rs index 3c185b4..50e5a65 100644 --- a/log_db/src/lib.rs +++ b/log_db/src/lib.rs @@ -4,7 +4,6 @@ extern crate log; #[macro_use] mod common; mod log_reader_forward; -mod log_reader_reverse; mod memtable_primary; mod memtable_secondary; mod record; @@ -13,7 +12,6 @@ pub use common::*; use fs2::FileExt; pub use log_reader_forward::ForwardLogReader; use log_reader_forward::ForwardLogReaderItem; -pub use log_reader_reverse::ReverseLogReader; use memtable_primary::PrimaryMemtable; use memtable_secondary::SecondaryMemtable; pub use record::Recordable; @@ -275,7 +273,7 @@ impl<R: Recordable> DB<R> { if record.tombstone { self.remove_record_from_memtables(&record); } else { - self.insert_record_to_memtables(&log_key, &record); + self.insert_record_to_memtables(log_key, record); } // Update from_index in case this is the last iteration: we need to know the next @@ -295,10 +293,7 @@ impl<R: Recordable> DB<R> { Ok(()) } - fn insert_record_to_memtables(&mut self, log_key: &LogKey, record: &Record) { - let pk = record.at(self.primary_key_index).as_indexable().unwrap(); - self.primary_memtable.set(&pk, &log_key); - + fn insert_record_to_memtables(&mut self, log_key: LogKey, record: Record) { for (sk_index, sk_field) in self.config.secondary_keys.iter().enumerate() { let secondary_memtable = &mut self.secondary_memtables[sk_index]; let sk_field_index = self @@ -309,8 +304,12 @@ impl<R: Recordable> DB<R> { .unwrap(); let sk = record.at(sk_field_index).as_indexable().unwrap(); - secondary_memtable.set(&sk, &log_key); + secondary_memtable.set(sk, log_key.clone()); } + + // Doing this last because this moves log_key + let pk = record.at(self.primary_key_index).as_indexable().unwrap(); + self.primary_memtable.set(pk, log_key); } fn remove_record_from_memtables(&mut self, record: &Record) { @@ -341,10 +340,30 @@ impl<R: Recordable> DB<R> { record.validate(&self.config.fields)?; debug!("Record is valid"); - self.upsert_record(&record) + self.batch_upsert_records(std::iter::once(record)) + } + + /// Insert a batch of records into the database. If the primary key value for a record already exists, + /// the existing record will be replaced by the supplied one. Records are inserted in the order they are given. + pub fn batch_upsert(&mut self, recordables: Vec<R>) -> Result<(), DBError> { + let records = recordables + .into_iter() + .map(|r| Record::from(&r.into_record())) + .collect::<Vec<Record>>(); + debug!("Batch upserting {} records", records.len()); + + for record in &records { + record.validate(&self.config.fields)?; + } + debug!("Records are valid"); + + self.batch_upsert_records(records.into_iter()) } - fn upsert_record(&mut self, record: &Record) -> Result<(), DBError> { + fn batch_upsert_records( + &mut self, + records: impl Iterator<Item = Record>, + ) -> Result<(), DBError> { debug!("Opening file in append mode and acquiring exclusive lock..."); // Acquire an exclusive lock for writing @@ -355,7 +374,7 @@ impl<R: Recordable> DB<R> { { // The log file has been rotated, so we must try again self.active_metadata_file.unlock()?; - return self.upsert_record(record); + return self.batch_upsert_records(records); } self.active_data_file.lock_exclusive()?; @@ -366,61 +385,59 @@ impl<R: Recordable> DB<R> { debug!("Exclusive lock acquired, appending to log file"); - // Write the record to the log - let serialized = &record.serialize(); - let record_offset = self.active_data_file.seek(SeekFrom::End(0))?; - let record_length = serialized.len() as u64; - assert!(record_length > 0); + let mut serialized_data: Vec<u8> = vec![]; + let mut serialized_metadata: Vec<u8> = vec![]; + let mut pending_memtable_insertions: Vec<(LogKey, Record)> = vec![]; + for record in records { + // Write the record to the log + let serialized = &record.serialize(); + let record_offset = self.active_data_file.seek(SeekFrom::End(0))?; + let record_length = serialized.len() as u64; + assert!(record_length > 0); - self.active_data_file.write_all(serialized)?; + serialized_data.extend(serialized); - // Flush and sync data to disk - if self.config.write_durability == WriteDurability::Flush { - self.active_data_file.flush()?; - } - if self.config.write_durability == WriteDurability::FlushSync { - self.active_data_file.flush()?; - self.active_data_file.sync_all()?; - } + let metadata_pos = self.active_metadata_file.seek(SeekFrom::End(0))?; + let metadata_index = + (metadata_pos - METADATA_FILE_HEADER_SIZE as u64) / METADATA_ROW_LENGTH as u64; + + // Write the record metadata to the metadata file + let mut metadata_buf = vec![]; + metadata_buf.extend(record_offset.to_be_bytes().into_iter()); + metadata_buf.extend(record_length.to_be_bytes().into_iter()); + + assert_eq!(metadata_buf.len(), 16); - let metadata_pos = self.active_metadata_file.seek(SeekFrom::End(0))?; - let metadata_index = - (metadata_pos - METADATA_FILE_HEADER_SIZE as u64) / METADATA_ROW_LENGTH as u64; + serialized_metadata.extend(metadata_buf); - // Write the record metadata to the metadata file - let mut metadata_buf = vec![]; - metadata_buf.extend(record_offset.to_be_bytes().into_iter()); - metadata_buf.extend(record_length.to_be_bytes().into_iter()); + let log_key = LogKey::new(segment_num, metadata_index); - assert_eq!(metadata_buf.len(), 16); - self.active_metadata_file.write_all(&metadata_buf)?; + pending_memtable_insertions.push((log_key, record)); + } + + self.active_data_file.write_all(&serialized_data)?; + self.active_metadata_file.write_all(&serialized_metadata)?; - // Flush and sync metadata to disk + // Flush and sync data and metadata to disk if self.config.write_durability == WriteDurability::Flush { + self.active_data_file.flush()?; self.active_metadata_file.flush()?; - } - if self.config.write_durability == WriteDurability::FlushSync { + } else if self.config.write_durability == WriteDurability::FlushSync { + self.active_data_file.flush()?; + self.active_data_file.sync_all()?; self.active_metadata_file.flush()?; self.active_metadata_file.sync_all()?; } - debug!("Record appended to log file, releasing locks"); + debug!("Records appended to log file, releasing locks"); // Manually release the locks because the file handles are left open self.active_data_file.unlock()?; self.active_metadata_file.unlock()?; - debug!("Update memtables with newly written data"); - let log_key = LogKey::new(segment_num, metadata_index); - - self.insert_record_to_memtables(&log_key, &record); - - // These post-condition asserts feel like they sometimes report false positives. - let len = self.active_metadata_file.seek(SeekFrom::End(0))?; - assert!(len >= METADATA_FILE_HEADER_SIZE as u64); - assert_eq!((len - METADATA_FILE_HEADER_SIZE as u64) % 16, 0); - let data_file_len = self.active_data_file.seek(SeekFrom::End(0))?; - assert_eq!(data_file_len, record_offset + record_length); + for (log_key, record) in pending_memtable_insertions { + self.insert_record_to_memtables(log_key, record); + } Ok(()) } @@ -505,7 +522,7 @@ impl<R: Recordable> DB<R> { if field == &self.config.primary_key { let opt = self.primary_memtable.get(&query_key); let log_keys = match opt { - Some(log_key) => vec![log_key.clone()], + Some(log_key) => vec![log_key], None => vec![], }; Ok(log_keys) @@ -525,21 +542,23 @@ impl<R: Recordable> DB<R> { let log_keys = self.secondary_memtables[smemtable_index] .find_by(&query_key) .into_iter() - .cloned() .collect(); Ok(log_keys) } }) - .collect::<Result<Vec<Vec<LogKey>>, DBError>>()?; + .collect::<Result<Vec<Vec<&LogKey>>, DBError>>()?; debug!("Found log keys in memtable: {:?}", log_key_batches); - let tagged = log_key_batches - .into_iter() - .enumerate() - .flat_map(|(tag, log_keys)| log_keys.into_iter().map(move |log_key| (tag, log_key))); + let mut tagged = vec![]; + let mut tag: usize = 0; + for batch in log_key_batches { + let mapped = batch.into_iter().map(|log_key| (tag, log_key)); + tagged.extend(mapped); + tag += 1; + } - let tagged_records = self.read_tagged_log_keys(tagged)?; + let tagged_records = self.read_tagged_log_keys(tagged.into_iter())?; debug!("Read {} records", tagged_records.len()); @@ -548,9 +567,9 @@ impl<R: Recordable> DB<R> { /// Read records from segment files based on log keys. /// The log keys are accompanied by an integer tag that can be used to identify and group them later. - fn read_tagged_log_keys( - &mut self, - log_keys: impl Iterator<Item = (usize, LogKey)>, + fn read_tagged_log_keys<'a>( + &self, + log_keys: impl Iterator<Item = (usize, &'a LogKey)>, ) -> Result<Vec<(usize, Record)>, DBError> { let mut records = vec![]; let mut log_keys_map = BTreeMap::new(); diff --git a/log_db/src/log_reader_reverse.rs b/log_db/src/log_reader_reverse.rs deleted file mode 100644 index ef2b0dc..0000000 --- a/log_db/src/log_reader_reverse.rs +++ /dev/null @@ -1,118 +0,0 @@ -use super::common::*; -use super::record::*; -use std::fs::{self}; -use std::io::{self, Read, Seek}; - -pub struct ReverseLogReader { - /// The contents of the metadata file are read into this buffer in one go. - metadata_buf: Vec<u8>, - - /// The data log file that is read based on offset + length information from the metadata file. - data_reader: io::BufReader<fs::File>, - - /// The current position in metadata_buf. - metadata_pos: usize, -} - -impl ReverseLogReader { - pub fn new( - mut metadata_file: fs::File, - data_file: fs::File, - ) -> Result<ReverseLogReader, io::Error> { - let mut metadata_buf = vec![]; - metadata_file.seek(io::SeekFrom::Start(0))?; - metadata_file.read_to_end(&mut metadata_buf)?; - let len = metadata_buf.len(); - - let ret = ReverseLogReader { - metadata_buf, - data_reader: io::BufReader::new(data_file), - metadata_pos: len, - }; - Ok(ret) - } - - fn read_record(&mut self) -> Result<Option<Record>, io::Error> { - loop { - assert!( - (self.metadata_pos - METADATA_FILE_HEADER_SIZE) % 16 == 0, - "metadata_pos is not aligned" - ); - // Return None if we have read the entire metadata file and - // reached the end of the header - if self.metadata_pos == METADATA_FILE_HEADER_SIZE { - return Ok(None); - } - - self.metadata_pos -= 16; - let i = self.metadata_pos; // shorter alias - - // First u64 is the offset of the record in the data file, second is the length of the record - let entry_offset = u64::from_be_bytes(self.metadata_buf[i..i + 8].try_into().unwrap()); - let entry_length = - u64::from_be_bytes(self.metadata_buf[i + 8..i + 16].try_into().unwrap()); - - if entry_offset == 0 && entry_length == 0 { - // This is an unused entry in the metadata file, skip - continue; - } - - self.data_reader.seek(io::SeekFrom::Start(entry_offset))?; - let mut result_buf = vec![0; entry_length as usize]; - self.data_reader.read_exact(&mut result_buf)?; - - let record = Record::deserialize(&result_buf); - - assert!( - (self.metadata_pos - METADATA_FILE_HEADER_SIZE) % 16 == 0, - "metadata_pos is not aligned" - ); - - return Ok(Some(record)); - } - } -} - -#[cfg(test)] -mod tests { - use super::*; - use std::path::Path; - - #[test] - fn test_reverse_log_reader_fixture_db1() { - let _ = env_logger::builder().is_test(true).try_init(); - let metadata_path = Path::new(TEST_RESOURCES_DIR).join("test_metadata_1"); - let metadata_file = fs::OpenOptions::new() - .read(true) - .open(&metadata_path) - .expect("Failed to open file"); - let data_file = fs::OpenOptions::new() - .read(true) - .open(Path::new(TEST_RESOURCES_DIR).join("test_data_1")) - .expect("Failed to open file"); - - let mut reverse_log_reader = ReverseLogReader::new(metadata_file, data_file).unwrap(); - - // There are two records in the log with "schema": Int - - let last_record = reverse_log_reader - .next() - .expect("Failed to read the last record"); - assert!(match &last_record.values[..] { - [Value::Bytes(bytes)] => bytes.len() == 256, - _ => false, - }); - - assert!(reverse_log_reader.next().is_none()); - } -} - -impl Iterator for ReverseLogReader { - type Item = Record; - - fn next(&mut self) -> Option<Self::Item> { - self.read_record().unwrap_or_else(|err| { - panic!("Error reading record: {:?}", err); - }) - } -} diff --git a/log_db/src/memtable_primary.rs b/log_db/src/memtable_primary.rs index e2921a6..573592e 100644 --- a/log_db/src/memtable_primary.rs +++ b/log_db/src/memtable_primary.rs @@ -19,8 +19,8 @@ impl PrimaryMemtable { } } - pub fn set(&mut self, key: &IndexableValue, value: &LogKey) { - self.records.insert(key.clone(), value.clone()); + pub fn set(&mut self, key: IndexableValue, value: LogKey) { + self.records.insert(key, value); } pub fn get(&self, key: &IndexableValue) -> Option<&LogKey> { @@ -31,10 +31,10 @@ impl PrimaryMemtable { self.records.remove(key) } - pub fn range<B: RangeBounds<IndexableValue>>(&self, range: B) -> Vec<LogKey> { + pub fn range<B: RangeBounds<IndexableValue>>(&self, range: B) -> Vec<&LogKey> { self.records .range(range) - .map(|(_, log_key)| log_key.clone()) + .map(|(_, log_key)| log_key) .collect() } } diff --git a/log_db/src/memtable_secondary.rs b/log_db/src/memtable_secondary.rs index dc161c8..cd6c252 100644 --- a/log_db/src/memtable_secondary.rs +++ b/log_db/src/memtable_secondary.rs @@ -19,14 +19,13 @@ impl SecondaryMemtable { } } - pub fn set(&mut self, key: &IndexableValue, value: &LogKey) { - match self.records.get_mut(key) { + pub fn set(&mut self, key: IndexableValue, value: LogKey) { + match self.records.get_mut(&key) { Some(set) => { - set.insert(value.clone()); + set.insert(value); } None => { - self.records - .insert(key.clone(), LogKeySet::new_with_initial(&value)); + self.records.insert(key, LogKeySet::new_with_initial(value)); } }; } @@ -38,11 +37,6 @@ impl SecondaryMemtable { } } - // Remove all log keys associated with the given key - pub fn remove_all(&mut self, key: &IndexableValue) -> Option<LogKeySet> { - self.records.remove(key) - } - // Remove a single log key associated with the given key. Returns `true` // if the log key existed and was removed, `false` otherwise. pub fn remove(&mut self, key: &IndexableValue, log_key: &LogKey) -> bool { @@ -62,24 +56,10 @@ impl SecondaryMemtable { } } - // Remove all log keys associated with the given log key - // Note: This is a linear time operation, prefer using the `remove` method - // if you know the secondary key associated with the log key. - pub fn scan_remove(&mut self, log_key: &LogKey) -> u64 { - let mut removed = 0; - self.records.iter_mut().for_each(|(_, set)| { - if let Ok(_) = set.remove(&log_key) { - removed += 1; - } - }); - - removed - } - - pub fn range<B: RangeBounds<IndexableValue>>(&self, range: B) -> Vec<LogKey> { + pub fn range<B: RangeBounds<IndexableValue>>(&self, range: B) -> Vec<&LogKey> { let mut keys = Vec::new(); for (_, set) in self.records.range(range) { - keys.extend(set.log_keys().iter().cloned()); + keys.extend(set.log_keys().iter()); } keys } |
