From bfa8fab553eab1b0daa062010fb2ccf5d365d21f Mon Sep 17 00:00:00 2001 From: Jan Tuomi Date: Sat, 9 Nov 2024 22:30:51 +0200 Subject: Refresh indexes on read or at maintenance --- log_db/src/common.rs | 179 +++++++++++++++++++++-- log_db/src/lib.rs | 296 +++++++++++++++++++++++++++------------ log_db/src/log_reader_forward.rs | 53 +++++-- log_db/src/memtable_primary.rs | 19 +-- log_db/src/memtable_secondary.rs | 46 +++--- log_db/tests/integration.rs | 5 - 6 files changed, 437 insertions(+), 161 deletions(-) diff --git a/log_db/src/common.rs b/log_db/src/common.rs index fb3124c..d86e6c6 100644 --- a/log_db/src/common.rs +++ b/log_db/src/common.rs @@ -1,10 +1,11 @@ use fs2::{lock_contended_error, FileExt}; use once_cell::sync::Lazy; use std::cmp::Ordering; -use std::collections::HashSet; +use std::collections::{BTreeMap, HashSet}; use std::fmt::Display; use std::fs::{self, metadata, File}; use std::io::{self, Read, Seek, SeekFrom, Write}; +use std::num::ParseIntError; use std::path::{Path, PathBuf}; use std::thread; use uuid::Uuid; @@ -43,6 +44,10 @@ impl LogKey { LogKey((segment_num as u64) << 48 | index) } + pub fn from(u: u64) -> Self { + LogKey(u) + } + pub fn segment_num(&self) -> u16 { (self.0 >> 48) as u16 } @@ -196,6 +201,22 @@ pub enum WriteDurability { FlushSync, } +#[derive(Debug, Clone, Eq, PartialEq)] +pub enum ReadConsistency { + /// Reads are **not** guaranteed to see the latest writes. + /// Indexes are only updated when running maintenance tasks. + /// This is the fastest option, but can produce stale reads. + Eventual, + /// Reads are guaranteed to see the latest writes from the same client, but + /// not from other clients. Written values are indexed after writing to file. + /// Indexes are updated when running maintenance tasks. + ReadMyWrites, + /// Reads are guaranteed to see the latest writes. + /// Indexes are updated synchronously before reads. + /// This is the slowest option, and can cause long waits for reads. + Strong, +} + impl Display for WriteDurability { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> Result<(), std::fmt::Error> { write!(f, "{:?}", self)?; @@ -268,6 +289,20 @@ pub enum Value { Bytes(Vec), } +impl PartialEq for Value { + fn eq(&self, other: &Self) -> bool { + match (self, other) { + (Value::Int(a), Value::Int(b)) => a == b, + (Value::Float(a), Value::Float(b)) => a == b, + (Value::String(a), Value::String(b)) => a == b, + (Value::Bytes(a), Value::Bytes(b)) => a == b, + (Value::Null, Value::Null) => true, + _ => false, + } + } +} +impl Eq for Value {} + impl Value { pub fn serialize(&self) -> Vec { match self { @@ -448,6 +483,30 @@ pub fn get_secondary_memtable_index_by_field( sks.iter().position(|schema_field| schema_field == field) } +pub struct GetSecondaryIndexPositionsResult { + pub schema_index: usize, + pub secondary_index: usize, +} + +pub fn get_secondary_index_positions( + sks: &Vec, + schema: &Vec<(Field, RecordField)>, +) -> Vec { + let mut ret = vec![]; + for (i, (schema_field, _)) in schema.iter().enumerate() { + for (j, sk) in sks.iter().enumerate() { + if schema_field == sk { + ret.push(GetSecondaryIndexPositionsResult { + schema_index: i, + secondary_index: j, + }); + break; + } + } + } + ret +} + /// A path to a log segment file along with its type pub enum SegmentPath { /// A symbolic link to the active log file @@ -531,7 +590,8 @@ pub fn create_segment_metadata_file( uuid: *data_file_uuid, }; - metadata_file.write_all(&metadata_header.serialize())?; + let header_serialized = metadata_header.serialize(); + metadata_file.write_all(&header_serialized)?; metadata_file.flush()?; let len = metadata_file.seek(io::SeekFrom::End(0))?; @@ -541,6 +601,19 @@ pub fn create_segment_metadata_file( Ok((new_num, metadata_path)) } +/// Parse number from format "metadata.1" +pub fn parse_segment_num(segment_filename: &Path) -> Result { + segment_filename + .file_name() + .expect("Not a valid path") + .to_str() + .expect("Not a valid UTF-8 string") + .split(".") + .last() + .expect("No extension") + .parse::() +} + /// Get the number of the segment with the greatest ordinal. /// This is the newest segment, i.e. the one that is pointed to by the `active` symlink. /// If there are no segments yet, returns 0. @@ -552,20 +625,8 @@ pub fn greatest_segment_number(data_dir_path: &Path) -> Result { } let segment_metadata_path = fs::read_link(&active_symlink)?; - let filename = segment_metadata_path - .file_name() - .expect("No filename in symlink") - .to_str() - .expect("Filename was not valid UTF-8"); - // parse number from format "metadata.1" - let segment_number = filename - .split('.') - .last() - .expect("Filename did not have a number") - .parse::(); - - segment_number.map_err(|_| { + parse_segment_num(&segment_metadata_path).map_err(|_| { io::Error::new( io::ErrorKind::InvalidData, "Failed to parse segment number from filename", @@ -775,3 +836,91 @@ pub fn request_exclusive_lock(data_dir: &Path, file: &mut fs::File) -> Result<() Ok(()) } + +pub fn get_record_by_log_key(data_dir: &Path, log_key: &LogKey) -> Result { + let segment_num = log_key.segment_num(); + let segment_index = log_key.index(); + + let metadata_path = data_dir.join(format!("metadata.{}", segment_num)); + let mut metadata_file = READ_MODE.open(metadata_path)?; + + request_shared_lock(data_dir, &mut metadata_file)?; + + let metadata_header = read_metadata_header(&mut metadata_file)?; + validate_metadata_header(&metadata_header)?; + + let data_file_path = data_dir.join(metadata_header.uuid.to_string()); + let mut data_file = READ_MODE.open(data_file_path)?; + + request_shared_lock(data_dir, &mut data_file)?; + + let metadata_offset = METADATA_FILE_HEADER_SIZE as u64 + segment_index * 16; + let mut metadata_buf = vec![0; 16]; + metadata_file.seek(SeekFrom::Start(metadata_offset))?; + metadata_file.read_exact(&mut metadata_buf)?; + + let data_offset = u64::from_be_bytes(metadata_buf[0..8].try_into().unwrap()); + let data_len = u64::from_be_bytes(metadata_buf[8..16].try_into().unwrap()); + + let mut data_buf = vec![0; data_len as usize]; + data_file.seek(SeekFrom::Start(data_offset))?; + data_file.read_exact(&mut data_buf)?; + + let record = Record::deserialize(&data_buf); + + Ok(record) +} + +pub fn get_records_by_log_keys( + data_dir: &Path, + log_keys: &[LogKey], +) -> Result, io::Error> { + // Partition by segment number + let mut segments: BTreeMap> = BTreeMap::new(); + for log_key in log_keys { + segments + .entry(log_key.segment_num()) + .or_default() + .push(log_key.index()); + } + + let mut records = Vec::new(); + + // Process newest (largest segment num) first + let mut segments_sorted = segments.iter().collect::>(); + segments_sorted.sort_by_key(|(&segment_num, _)| -(segment_num as i32)); + + for (segment_num, indices) in segments_sorted { + let metadata_path = data_dir.join(format!("metadata.{}", segment_num)); + let mut metadata_file = READ_MODE.open(metadata_path)?; + + request_shared_lock(data_dir, &mut metadata_file)?; + + let metadata_header = read_metadata_header(&mut metadata_file)?; + validate_metadata_header(&metadata_header)?; + + let data_file_path = data_dir.join(metadata_header.uuid.to_string()); + let mut data_file = READ_MODE.open(data_file_path)?; + + request_shared_lock(data_dir, &mut data_file)?; + + for index in indices { + let metadata_offset = METADATA_FILE_HEADER_SIZE as u64 + index * 16; + let mut metadata_buf = vec![0; 16]; + metadata_file.seek(SeekFrom::Start(metadata_offset))?; + metadata_file.read_exact(&mut metadata_buf)?; + + let data_offset = u64::from_be_bytes(metadata_buf[0..8].try_into().unwrap()); + let data_len = u64::from_be_bytes(metadata_buf[8..16].try_into().unwrap()); + + let mut data_buf = vec![0; data_len as usize]; + data_file.seek(SeekFrom::Start(data_offset))?; + data_file.read_exact(&mut data_buf)?; + + let record = Record::deserialize(&data_buf); + records.push(record); + } + } + + Ok(records) +} diff --git a/log_db/src/lib.rs b/log_db/src/lib.rs index 68db7fd..318334c 100644 --- a/log_db/src/lib.rs +++ b/log_db/src/lib.rs @@ -31,6 +31,7 @@ pub struct ConfigBuilder { primary_key: Option, secondary_keys: Option>, write_durability: Option, + read_consistency: Option, } impl<'a, Field: Eq + Clone + Debug> ConfigBuilder { @@ -43,6 +44,7 @@ impl<'a, Field: Eq + Clone + Debug> ConfigBuilder { primary_key: None, secondary_keys: None, write_durability: None, + read_consistency: None, } } @@ -97,6 +99,13 @@ impl<'a, Field: Eq + Clone + Debug> ConfigBuilder { self } + /// The read consistency policy for the database. + /// A less strict policy will result in faster reads, but may return stale data. + pub fn read_consistency(&mut self, read_consistency: ReadConsistency) -> &mut Self { + self.read_consistency = Some(read_consistency); + self + } + pub fn initialize(&self) -> Result, io::Error> { let config = Config:: { data_dir: self.data_dir.clone().unwrap_or("db_data".to_string()), @@ -119,6 +128,10 @@ impl<'a, Field: Eq + Clone + Debug> ConfigBuilder { .write_durability .clone() .unwrap_or(WriteDurability::Flush), + read_consistency: self + .read_consistency + .clone() + .unwrap_or(ReadConsistency::Strong), }; DB::initialize(&config) @@ -134,16 +147,21 @@ struct Config { pub primary_key: Field, pub secondary_keys: Vec, pub write_durability: WriteDurability, + pub read_consistency: ReadConsistency, } pub struct DB { config: Config, data_dir: PathBuf, + active_segment_num: u16, active_metadata_file: fs::File, active_data_file: fs::File, primary_key_index: usize, primary_memtable: PrimaryMemtable, secondary_memtables: Vec, + + /// The position of the latest log entry that has been read into memtable indexes, plus one. + next_index_refresh_index: LogKey, } impl DB { @@ -240,6 +258,8 @@ impl DB { let active_symlink = Path::new(&config.data_dir).join(ACTIVE_SYMLINK_FILENAME); let active_target = fs::read_link(&active_symlink)?; + let active_segment_num = parse_segment_num(&active_target) + .map_err(|e| io::Error::new(io::ErrorKind::InvalidData, e))?; let active_metadata_path = Path::new(&config.data_dir).join(active_target); let mut active_metadata_file = APPEND_MODE.open(&active_metadata_path)?; @@ -253,11 +273,13 @@ impl DB { let db = DB:: { config: config.clone(), data_dir: data_dir_path.to_path_buf(), + active_segment_num, active_metadata_file, active_data_file, primary_key_index, primary_memtable, secondary_memtables, + next_index_refresh_index: LogKey::new(1, 0), }; // info!("Rebuilding memtable indexes..."); @@ -314,6 +336,7 @@ impl DB { metadata_buf.extend(&record_length.to_be_bytes()); assert_eq!(metadata_buf.len(), 16); + let metadata_offset = self.active_metadata_file.seek(SeekFrom::End(0))?; self.active_metadata_file.write_all(&metadata_buf)?; // Flush and sync metadata to disk @@ -331,13 +354,37 @@ impl DB { debug!("Record appended to log file, lock released"); + if self.config.read_consistency == ReadConsistency::ReadMyWrites { + debug!("Updating memtable indexes with new write..."); + + let sk_positions = + get_secondary_index_positions(&self.config.secondary_keys, &self.config.fields); + let index = (metadata_offset - METADATA_FILE_HEADER_SIZE as u64) / 16; + let log_key = LogKey::new(self.active_segment_num, index); + + let pk = record + .at(self.primary_key_index) + .as_indexable() + .expect("Primary key must be an IndexableValue"); + self.primary_memtable.set(&pk, &log_key); + + for sk_pos in sk_positions { + let sk = record + .at(sk_pos.schema_index) + .as_indexable() + .expect("Secondary key must be an IndexableValue"); + self.secondary_memtables[sk_pos.secondary_index].set(&sk, &log_key); + } + + debug!("Memtable indexes updated"); + } + 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); - // Depending on write durability, the data might not be written to disk yet - //let data_file_len = self.active_data_file.seek(SeekFrom::End(0))?; - //assert_eq!(data_file_len, record_offset + record_length); + let data_file_len = self.active_data_file.seek(SeekFrom::End(0))?; + assert_eq!(data_file_len, record_offset + record_length); Ok(()) } @@ -371,68 +418,20 @@ impl DB { } }; - debug!("Looking up key {:?} in primary memtable", query_key); - let found = self.primary_memtable.get(&query_key); - if let Some(record) = found { - debug!("Found record in primary memtable: {:?}", record); - return Ok(Some(record.clone())); + if self.config.read_consistency == ReadConsistency::Strong { + debug!("Refreshing indexes before reading"); + self.refresh_indexes()?; } - debug!( - "No memtable entry found, looking up key {:?} in log file", - query_key - ); - - debug!( - "Matching records based on value at primary key index ({})", - &self.primary_key_index - ); - - let greatest = greatest_segment_number(&self.data_dir)?; - debug!("Searching segments {} through 1", greatest); - - let mut found_record: Option = None; - for segment_num in (1..=greatest).rev() { - let segment_path = &self.data_dir.join(format!("metadata.{}", segment_num)); - - debug!( - "Opening segment {} in read mode and acquiring shared lock...", - segment_num - ); - - let mut metadata_file = READ_MODE.open(&segment_path)?; - - request_shared_lock(&self.data_dir, &mut metadata_file)?; - - let metadata_header = read_metadata_header(&mut metadata_file)?; - - validate_metadata_header(&metadata_header)?; - - let data_path = &self.data_dir.join(metadata_header.uuid.to_string()); - let data_file = READ_MODE.open(&data_path)?; - - // We should not "request_shared_lock()" here because we do not want - // to give way to writers at this point. That would possibly lead to a deadlock. - data_file.lock_shared()?; - - let mut reader = ReverseLogReader::new(metadata_file, data_file)?; - - if let Some(found) = reader.find(|record| { - let record_key = record - .at(self.primary_key_index) - .as_indexable() - .expect("Primary key must be indexable"); - record_key == query_key - }) { - found_record = Some(found); - break; - } + debug!("Looking up key {:?} in primary memtable", query_key); + if let Some(log_key) = self.primary_memtable.get(&query_key) { + debug!("Found record log key in primary memtable, reading from file"); + let record = get_record_by_log_key(&self.data_dir, log_key)?; + return Ok(Some(record)); } - debug!("Record search complete"); - debug!("Found matching record in log file."); - - Ok(found_record) + debug!("Key not found in primary memtable, returning None"); + Ok(None) } /// Get a collection of records based on a field value. @@ -473,25 +472,29 @@ impl DB { } }; + if self.config.read_consistency == ReadConsistency::Strong { + debug!("Refreshing indexes before reading"); + self.refresh_indexes()?; + } + // Try to find a memtable with the queried key let found_memtable_index = get_secondary_memtable_index_by_field(&self.config.secondary_keys, field); + // If the requested key is indexed, we can just look it up in the memtable + // and return the matching records. if let Some(memtable_index) = found_memtable_index { debug!( "Found suitable secondary index. Looking up key {:?} in the memtable", query_key ); - let records = self.secondary_memtables[memtable_index] - .find_all(&self.primary_memtable, &query_key); - debug!("Found matching key"); - return Ok(records.iter().map(|record| record.clone()).collect()); + let log_keys = self.secondary_memtables[memtable_index].find_all(&query_key); + let log_keys = log_keys.iter().cloned().collect::>(); + let records = get_records_by_log_keys(&self.data_dir, &log_keys)?; + return Ok(records); } - debug!( - "No memtable entry found, looking up key {:?} in log file", - query_key - ); + debug!("Key is not indexed, looking up matching values in log file"); // Get the index of the requested field let key_index = self @@ -548,26 +551,11 @@ impl DB { } } - debug!("Record search complete"); - debug!( - "Number of matching records found in log file: {}", + "Record search complete, number of matching records found in log file: {}", found_records.len() ); - if let Some(memtable_index) = found_memtable_index { - debug!("Inserting result set into secondary index"); - let primary_values: Vec = found_records - .iter() - .map(|r| { - r.at(self.primary_key_index) - .as_indexable() - .expect("A non-indexable value was stored at primary key index") - }) - .collect(); - self.secondary_memtables[memtable_index].set_all(&query_key, &primary_values); - } - Ok(found_records) } @@ -593,6 +581,8 @@ impl DB { self.active_metadata_file = metadata_file; self.active_data_file = APPEND_MODE.open(&data_file_path)?; + self.active_segment_num = parse_segment_num(&active_metadata_path) + .map_err(|e| io::Error::new(io::ErrorKind::InvalidData, e))?; return Ok(false); } else { @@ -648,6 +638,8 @@ impl DB { self.active_metadata_file = APPEND_MODE.open(&new_metadata_path)?; let new_data_path = &self.data_dir.join(data_uuid.to_string()); self.active_data_file = APPEND_MODE.open(&new_data_path)?; + self.active_segment_num = parse_segment_num(&new_metadata_path) + .map_err(|e| io::Error::new(io::ErrorKind::InvalidData, e))?; // Compact the rotated segment without a lock. // Since the rotated segment and the compacted segment based on it will be @@ -662,6 +654,11 @@ impl DB { self.active_metadata_file.unlock()?; self.active_data_file.unlock()?; + debug!("Refreshing indexes"); + self.refresh_indexes()?; + + debug!("Maintenance tasks complete"); + Ok(()) } @@ -677,12 +674,13 @@ impl DB { debug!("Reading segment data into a BTreeMap"); let mut map = BTreeMap::new(); let forward_log_reader = ForwardLogReader::new(metadata_file, data_file); - for entry in forward_log_reader { - let primary_key = entry + for item in forward_log_reader { + let primary_key = item + .record .at(self.primary_key_index) .as_indexable() .expect("Primary key was not indexable"); - map.insert(primary_key, entry); + map.insert(primary_key, item.record); } debug!("Opening temporary files for writing compacted data"); @@ -734,6 +732,77 @@ impl DB { debug!("Compaction complete, resulting size: {}", final_len); Ok(()) } + + fn refresh_indexes(&mut self) -> Result<(), io::Error> { + let sk_positions = + get_secondary_index_positions(&self.config.secondary_keys, &self.config.fields); + + // Read segments in order, starting from `self.next_index_refresh_index` and going up + let greatest = greatest_segment_number(&self.data_dir)?; + + let next_segment_num = self.next_index_refresh_index.segment_num(); + let mut next_segment_index = self.next_index_refresh_index.index(); + debug!( + "Refreshing indexes from segment {} index {} onwards", + next_segment_num, next_segment_index + ); + + for segment_num in next_segment_num..=greatest { + let metadata_path = self.data_dir.join(format!("metadata.{}", segment_num)); + let mut metadata_file = READ_MODE.open(&metadata_path)?; + + request_shared_lock(&self.data_dir, &mut metadata_file)?; + + let metadata_header = read_metadata_header(&mut metadata_file)?; + validate_metadata_header(&metadata_header)?; + + let data_file_path = self.data_dir.join(metadata_header.uuid.to_string()); + let data_file = READ_MODE.open(&data_file_path)?; + + data_file.lock_shared()?; + + let forward_log_reader = + ForwardLogReader::new_with_index(metadata_file, data_file, next_segment_index); + + for item in forward_log_reader { + let pk = item + .record + .at(self.primary_key_index) + .as_indexable() + .expect("Primary key was not indexable"); + let log_key = LogKey::new(segment_num, item.metadata_index); + + debug!( + "Inserting primary key {:?} -> segment {} index {}", + pk, segment_num, item.metadata_index + ); + + self.primary_memtable.set(&pk, &log_key); + for sk_pos in &sk_positions { + let sk = item + .record + .at(sk_pos.schema_index) + .as_indexable() + .expect("Secondary key was not indexable"); + self.secondary_memtables[sk_pos.secondary_index].set(&sk, &log_key); + } + + next_segment_index = item.metadata_index + 1; + } + + if segment_num < greatest { + next_segment_index = 0; + } + } + + self.next_index_refresh_index = LogKey::new(greatest, next_segment_index); + debug!( + "Finished refreshing indexes, stored checkpoint to segment {} index {}", + greatest, next_segment_index + ); + + Ok(()) + } } #[cfg(test)] @@ -851,4 +920,55 @@ mod tests { let len = file.seek(SeekFrom::End(0)).expect("Failed to seek"); assert_eq!(len, METADATA_FILE_HEADER_SIZE as u64 + n_recs * 16); } + + #[test] + fn test_get_one_by_log_key() { + let _ = env_logger::builder().is_test(true).try_init(); + let temp_dir = tempfile::tempdir().unwrap(); + let data_dir = temp_dir.path(); + + // Create segment + let (data_uuid1, _) = create_segment_data_file(&data_dir).unwrap(); + let (segment_num1, _) = create_segment_metadata_file(&data_dir, &data_uuid1).unwrap(); + set_active_segment(&data_dir, segment_num1).unwrap(); + + // Open segment1 and write record 3 times + let mut metadata_file = APPEND_MODE + .open(&data_dir.join(format!("metadata.{}", segment_num1))) + .unwrap(); + let mut data_file = APPEND_MODE + .open(&data_dir.join(data_uuid1.to_string())) + .unwrap(); + + let mut data_offset: u64 = 0; + for i in 0..3 { + let record = Record::from(&[Value::Int(i)]); + let serialized = record.serialize(); + let data_len = serialized.len() as u64; + + data_file.write_all(&serialized).unwrap(); + + metadata_file.write_all(&data_offset.to_be_bytes()).unwrap(); + metadata_file.write_all(&data_len.to_be_bytes()).unwrap(); + + metadata_file.flush().unwrap(); + data_file.flush().unwrap(); + + data_offset += data_len; + } + + // Check that `get_record_by_log_key` works + let rec = get_record_by_log_key(&data_dir, &LogKey::new(1, 0)).unwrap(); + assert_eq!(rec.at(0), &Value::Int(0)); + + // Check that `get_records_by_log_keys` works + let recs = get_records_by_log_keys( + &data_dir, + &[LogKey::new(1, 0), LogKey::new(1, 1), LogKey::new(1, 2)], + ) + .unwrap(); + assert_eq!(recs[0].at(0), &Value::Int(0)); + assert_eq!(recs[1].at(0), &Value::Int(1)); + assert_eq!(recs[2].at(0), &Value::Int(2)); + } } diff --git a/log_db/src/log_reader_forward.rs b/log_db/src/log_reader_forward.rs index 343d7ff..bfd329b 100644 --- a/log_db/src/log_reader_forward.rs +++ b/log_db/src/log_reader_forward.rs @@ -21,7 +21,29 @@ impl<'a> ForwardLogReader { ret } - fn read_record(&mut self) -> Result, io::Error> { + pub fn new_with_index( + metadata_file: fs::File, + data_file: fs::File, + index: u64, + ) -> ForwardLogReader { + let mut ret = ForwardLogReader { + metadata_reader: io::BufReader::new(metadata_file), + data_reader: io::BufReader::new(data_file), + }; + + ret.metadata_reader + .seek(io::SeekFrom::Start( + METADATA_FILE_HEADER_SIZE as u64 + 16 * index, + )) + .expect("Seek failed"); + + ret + } + + fn read_record(&mut self) -> Result, io::Error> { + let metadata_index = + (self.metadata_reader.stream_position()? - METADATA_FILE_HEADER_SIZE as u64) / 16; + let mut metadata_entry_buf = vec![0; 16]; // 2x u64 if let Err(e) = self.metadata_reader.read_exact(&mut metadata_entry_buf) { if e.kind() == io::ErrorKind::UnexpectedEof { @@ -32,24 +54,37 @@ impl<'a> ForwardLogReader { } // 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(metadata_entry_buf[0..8].try_into().unwrap()); - let entry_length = u64::from_be_bytes(metadata_entry_buf[8..16].try_into().unwrap()); + let data_offset = u64::from_be_bytes(metadata_entry_buf[0..8].try_into().unwrap()); + let data_length = u64::from_be_bytes(metadata_entry_buf[8..16].try_into().unwrap()); // Use .seek_relative instead of .seek to avoid dropping the BufReader internal buffer when // the seek distance is small - let seek_distance = entry_offset - self.data_reader.stream_position()?; + let seek_distance = data_offset - self.data_reader.stream_position()?; self.data_reader.seek_relative(seek_distance as i64)?; - let mut result_buf = vec![0; entry_length as usize]; + let mut result_buf = vec![0; data_length as usize]; self.data_reader.read_exact(&mut result_buf)?; let record = Record::deserialize(&result_buf); - Ok(Some(record)) + let item = ForwardLogReaderItem { + metadata_index, + data_offset, + data_length, + record, + }; + Ok(Some(item)) } } +pub struct ForwardLogReaderItem { + pub metadata_index: u64, + pub data_offset: u64, + pub data_length: u64, + pub record: Record, +} + impl Iterator for ForwardLogReader { - type Item = Record; + type Item = ForwardLogReaderItem; fn next(&mut self) -> Option { match self.read_record() { @@ -83,10 +118,10 @@ mod tests { // There are two records in the log with "schema" with one field: Bytes - let first_record = forward_log_reader + let first_item = forward_log_reader .next() .expect("Failed to read the first record"); - assert!(match first_record.values() { + assert!(match first_item.record.values() { [Value::Bytes(bytes)] => bytes.len() == 256, _ => false, }); diff --git a/log_db/src/memtable_primary.rs b/log_db/src/memtable_primary.rs index 9d7e1e8..8340cdd 100644 --- a/log_db/src/memtable_primary.rs +++ b/log_db/src/memtable_primary.rs @@ -2,14 +2,9 @@ use super::common::*; use std::collections::BTreeMap; pub struct PrimaryMemtable { - /// Map of records indexed by key. Used as a shared heap of records - /// for all secondary memtables also. Secondary memtables store an - /// IndexableValue as their record value, which is used to get - /// the actual record from the primary memtable `records` map. - /// - /// Note: it must be invariant that all memtables (primary and secondary) - /// contain the same keys. - records: BTreeMap, + /// Map of record locations indexed by key. Use the `LogKey` values to look up + /// the actual records in the log files. + records: BTreeMap, } impl PrimaryMemtable { @@ -19,15 +14,11 @@ impl PrimaryMemtable { } } - pub fn set(&mut self, key: &IndexableValue, value: &Record) { + pub fn set(&mut self, key: &IndexableValue, value: &LogKey) { self.records.insert(key.clone(), value.clone()); } - pub fn get(&mut self, key: &IndexableValue) -> Option<&Record> { - self.records.get(key) - } - - pub fn get_without_update(&self, key: &IndexableValue) -> Option<&Record> { + pub fn get(&self, key: &IndexableValue) -> Option<&LogKey> { self.records.get(key) } } diff --git a/log_db/src/memtable_secondary.rs b/log_db/src/memtable_secondary.rs index 3f9863b..8300ada 100644 --- a/log_db/src/memtable_secondary.rs +++ b/log_db/src/memtable_secondary.rs @@ -1,14 +1,15 @@ use super::*; +use once_cell::sync::Lazy; use std::collections::BTreeMap; use std::collections::HashSet; pub struct SecondaryMemtable { - /// Map of records indexed by key. The value is the set of primary key values of records - /// that have the secondary key value. The actual `Record` objects are stored in the - /// primary memtable, which acts as the shared heap. - records: BTreeMap>, + /// Map of records indexed by key. Values are non-empty sets of `LogKey` values. + records: BTreeMap, } +static EMPTY_SET: Lazy> = Lazy::new(|| HashSet::new()); + impl SecondaryMemtable { pub fn new() -> SecondaryMemtable { SecondaryMemtable { @@ -16,10 +17,12 @@ impl SecondaryMemtable { } } - pub fn set(&mut self, key: &IndexableValue, value: &IndexableValue) { + pub fn set(&mut self, key: &IndexableValue, value: &LogKey) { debug!( - "Inserting/updating record in secondary memtable with key {:?} = {:?}", - &key, &value, + "Inserting/updating record in secondary memtable with key {:?} = segment {} index {}", + &key, + &value.segment_num(), + &value.index() ); match self.records.get_mut(key) { @@ -32,44 +35,27 @@ impl SecondaryMemtable { } None => { debug!("No existing entry found, creating one."); - let mut set = HashSet::with_capacity(1); - set.insert(value.clone()); + let set = LogKeySet::new_with_initial(&value); self.records.insert(key.clone(), set); } } } - pub fn set_all(&mut self, key: &IndexableValue, values: &[IndexableValue]) { + pub fn set_all(&mut self, key: &IndexableValue, values: &[LogKey]) { debug!( "Replacing set of records in secondary memtable with key {:?} ({} values)", &key, &values.len(), ); - let mut set = HashSet::with_capacity(values.len()); - values.iter().for_each(|value| { - set.insert(value.clone()); - }); - + let set = LogKeySet::from_slice(values); self.records.insert(key.clone(), set); } - pub fn find_all( - &mut self, - primary_memtable: &PrimaryMemtable, - key: &IndexableValue, - ) -> Vec { + pub fn find_all(&mut self, key: &IndexableValue) -> &HashSet { match self.records.get(key) { - None => vec![], - Some(set) => set - .iter() - .map(|key| { - primary_memtable - .get_without_update(key) - .expect("Record not found") - .clone() - }) - .collect(), + Some(set) => set.log_keys(), + None => &EMPTY_SET, } } } diff --git a/log_db/tests/integration.rs b/log_db/tests/integration.rs index fa81542..f178738 100644 --- a/log_db/tests/integration.rs +++ b/log_db/tests/integration.rs @@ -207,7 +207,6 @@ fn test_upsert_fails_on_invalid_value_type() { } #[test] -#[ignore] fn test_upsert_and_get_from_secondary_memtable() { let data_dir = tmp_dir(); let mut db = DB::configure() @@ -244,10 +243,6 @@ fn test_upsert_and_get_from_secondary_memtable() { ]); db.upsert(&record2).unwrap(); - // Delete the DB so that any results must come from a memtable - fs::remove_file(Path::new(&data_dir).join(ACTIVE_SYMLINK_FILENAME)) - .expect("Failed to delete the DB log file"); - // There should be 2 Johns let johns = db .find_all(&Field::Name, &Value::String("John".to_string())) -- cgit v1.3