aboutsummaryrefslogtreecommitdiffstats
path: root/log_db
diff options
context:
space:
mode:
authorJan Tuomi <jan@jantuomi.fi>2024-11-21 14:47:20 +0200
committerJan Tuomi <jan@jantuomi.fi>2024-11-21 14:47:20 +0200
commitf759208ae41a5bbeb163fda4b66f43b8c0e84b2b (patch)
tree10ac7d07e39b2437ca4ccb6b0f5123f8c64d3338 /log_db
parenta01b696134e56d0e0a57cb5ff6c6dcb5914d05bb (diff)
Revert "Refresh indexes on read or at maintenance"
This reverts commit bfa8fab553eab1b0daa062010fb2ccf5d365d21f.
Diffstat (limited to 'log_db')
-rw-r--r--log_db/src/common.rs179
-rw-r--r--log_db/src/lib.rs296
-rw-r--r--log_db/src/log_reader_forward.rs53
-rw-r--r--log_db/src/memtable_primary.rs19
-rw-r--r--log_db/src/memtable_secondary.rs46
-rw-r--r--log_db/tests/integration.rs5
6 files changed, 161 insertions, 437 deletions
diff --git a/log_db/src/common.rs b/log_db/src/common.rs
index 3590cf9..5fd9d51 100644
--- a/log_db/src/common.rs
+++ b/log_db/src/common.rs
@@ -1,11 +1,10 @@
use fs2::{lock_contended_error, FileExt};
use once_cell::sync::Lazy;
use std::cmp::Ordering;
-use std::collections::{BTreeMap, HashSet};
+use std::collections::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;
@@ -44,10 +43,6 @@ 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
}
@@ -201,22 +196,6 @@ 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)?;
@@ -289,20 +268,6 @@ pub enum Value {
Bytes(Vec<u8>),
}
-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<u8> {
match self {
@@ -491,30 +456,6 @@ pub fn get_secondary_memtable_index_by_field<Field: Eq>(
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<Field: Eq>(
- sks: &Vec<Field>,
- schema: &Vec<(Field, RecordField)>,
-) -> Vec<GetSecondaryIndexPositionsResult> {
- 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
@@ -598,8 +539,7 @@ pub fn create_segment_metadata_file(
uuid: *data_file_uuid,
};
- let header_serialized = metadata_header.serialize();
- metadata_file.write_all(&header_serialized)?;
+ metadata_file.write_all(&metadata_header.serialize())?;
metadata_file.flush()?;
let len = metadata_file.seek(io::SeekFrom::End(0))?;
@@ -609,19 +549,6 @@ 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<u16, ParseIntError> {
- segment_filename
- .file_name()
- .expect("Not a valid path")
- .to_str()
- .expect("Not a valid UTF-8 string")
- .split(".")
- .last()
- .expect("No extension")
- .parse::<u16>()
-}
-
/// 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.
@@ -633,8 +560,20 @@ pub fn greatest_segment_number(data_dir_path: &Path) -> Result<u16, io::Error> {
}
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_segment_num(&segment_metadata_path).map_err(|_| {
+ // parse number from format "metadata.1"
+ let segment_number = filename
+ .split('.')
+ .last()
+ .expect("Filename did not have a number")
+ .parse::<u16>();
+
+ segment_number.map_err(|_| {
io::Error::new(
io::ErrorKind::InvalidData,
"Failed to parse segment number from filename",
@@ -844,91 +783,3 @@ 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<Record, io::Error> {
- 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<Vec<Record>, io::Error> {
- // Partition by segment number
- let mut segments: BTreeMap<u16, Vec<u64>> = 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::<Vec<_>>();
- 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 318334c..68db7fd 100644
--- a/log_db/src/lib.rs
+++ b/log_db/src/lib.rs
@@ -31,7 +31,6 @@ pub struct ConfigBuilder<Field: Eq + Clone + Debug> {
primary_key: Option<Field>,
secondary_keys: Option<Vec<Field>>,
write_durability: Option<WriteDurability>,
- read_consistency: Option<ReadConsistency>,
}
impl<'a, Field: Eq + Clone + Debug> ConfigBuilder<Field> {
@@ -44,7 +43,6 @@ impl<'a, Field: Eq + Clone + Debug> ConfigBuilder<Field> {
primary_key: None,
secondary_keys: None,
write_durability: None,
- read_consistency: None,
}
}
@@ -99,13 +97,6 @@ impl<'a, Field: Eq + Clone + Debug> ConfigBuilder<Field> {
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<DB<Field>, io::Error> {
let config = Config::<Field> {
data_dir: self.data_dir.clone().unwrap_or("db_data".to_string()),
@@ -128,10 +119,6 @@ impl<'a, Field: Eq + Clone + Debug> ConfigBuilder<Field> {
.write_durability
.clone()
.unwrap_or(WriteDurability::Flush),
- read_consistency: self
- .read_consistency
- .clone()
- .unwrap_or(ReadConsistency::Strong),
};
DB::initialize(&config)
@@ -147,21 +134,16 @@ struct Config<Field: Eq + Clone> {
pub primary_key: Field,
pub secondary_keys: Vec<Field>,
pub write_durability: WriteDurability,
- pub read_consistency: ReadConsistency,
}
pub struct DB<Field: Eq + Clone + Debug> {
config: Config<Field>,
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<SecondaryMemtable>,
-
- /// The position of the latest log entry that has been read into memtable indexes, plus one.
- next_index_refresh_index: LogKey,
}
impl<Field: Eq + Clone + Debug> DB<Field> {
@@ -258,8 +240,6 @@ impl<Field: Eq + Clone + Debug> DB<Field> {
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)?;
@@ -273,13 +253,11 @@ impl<Field: Eq + Clone + Debug> DB<Field> {
let db = DB::<Field> {
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...");
@@ -336,7 +314,6 @@ impl<Field: Eq + Clone + Debug> DB<Field> {
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
@@ -354,37 +331,13 @@ impl<Field: Eq + Clone + Debug> DB<Field> {
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);
- let data_file_len = self.active_data_file.seek(SeekFrom::End(0))?;
- assert_eq!(data_file_len, record_offset + record_length);
+ // 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);
Ok(())
}
@@ -418,20 +371,68 @@ impl<Field: Eq + Clone + Debug> DB<Field> {
}
};
- if self.config.read_consistency == ReadConsistency::Strong {
- debug!("Refreshing indexes before reading");
- self.refresh_indexes()?;
+ 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()));
}
- 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!(
+ "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<Record> = 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!("Key not found in primary memtable, returning None");
- Ok(None)
+ debug!("Record search complete");
+ debug!("Found matching record in log file.");
+
+ Ok(found_record)
}
/// Get a collection of records based on a field value.
@@ -472,29 +473,25 @@ impl<Field: Eq + Clone + Debug> DB<Field> {
}
};
- 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 log_keys = self.secondary_memtables[memtable_index].find_all(&query_key);
- let log_keys = log_keys.iter().cloned().collect::<Vec<_>>();
- let records = get_records_by_log_keys(&self.data_dir, &log_keys)?;
- return Ok(records);
+ 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());
}
- debug!("Key is not indexed, looking up matching values in log file");
+ debug!(
+ "No memtable entry found, looking up key {:?} in log file",
+ query_key
+ );
// Get the index of the requested field
let key_index = self
@@ -551,11 +548,26 @@ impl<Field: Eq + Clone + Debug> DB<Field> {
}
}
+ debug!("Record search complete");
+
debug!(
- "Record search complete, number of matching records found in log file: {}",
+ "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<IndexableValue> = 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)
}
@@ -581,8 +593,6 @@ impl<Field: Eq + Clone + Debug> DB<Field> {
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 {
@@ -638,8 +648,6 @@ impl<Field: Eq + Clone + Debug> DB<Field> {
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
@@ -654,11 +662,6 @@ impl<Field: Eq + Clone + Debug> DB<Field> {
self.active_metadata_file.unlock()?;
self.active_data_file.unlock()?;
- debug!("Refreshing indexes");
- self.refresh_indexes()?;
-
- debug!("Maintenance tasks complete");
-
Ok(())
}
@@ -674,13 +677,12 @@ impl<Field: Eq + Clone + Debug> DB<Field> {
debug!("Reading segment data into a BTreeMap");
let mut map = BTreeMap::new();
let forward_log_reader = ForwardLogReader::new(metadata_file, data_file);
- for item in forward_log_reader {
- let primary_key = item
- .record
+ for entry in forward_log_reader {
+ let primary_key = entry
.at(self.primary_key_index)
.as_indexable()
.expect("Primary key was not indexable");
- map.insert(primary_key, item.record);
+ map.insert(primary_key, entry);
}
debug!("Opening temporary files for writing compacted data");
@@ -732,77 +734,6 @@ impl<Field: Eq + Clone + Debug> DB<Field> {
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)]
@@ -920,55 +851,4 @@ 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 bfd329b..343d7ff 100644
--- a/log_db/src/log_reader_forward.rs
+++ b/log_db/src/log_reader_forward.rs
@@ -21,29 +21,7 @@ impl<'a> ForwardLogReader {
ret
}
- 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<Option<ForwardLogReaderItem>, io::Error> {
- let metadata_index =
- (self.metadata_reader.stream_position()? - METADATA_FILE_HEADER_SIZE as u64) / 16;
-
+ fn read_record(&mut self) -> Result<Option<Record>, io::Error> {
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 {
@@ -54,37 +32,24 @@ impl<'a> ForwardLogReader {
}
// First u64 is the offset of the record in the data file, second is the length of the record
- 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());
+ 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());
// Use .seek_relative instead of .seek to avoid dropping the BufReader internal buffer when
// the seek distance is small
- let seek_distance = data_offset - self.data_reader.stream_position()?;
+ let seek_distance = entry_offset - self.data_reader.stream_position()?;
self.data_reader.seek_relative(seek_distance as i64)?;
- let mut result_buf = vec![0; data_length as usize];
+ let mut result_buf = vec![0; entry_length as usize];
self.data_reader.read_exact(&mut result_buf)?;
let record = Record::deserialize(&result_buf);
- let item = ForwardLogReaderItem {
- metadata_index,
- data_offset,
- data_length,
- record,
- };
- Ok(Some(item))
+ Ok(Some(record))
}
}
-pub struct ForwardLogReaderItem {
- pub metadata_index: u64,
- pub data_offset: u64,
- pub data_length: u64,
- pub record: Record,
-}
-
impl Iterator for ForwardLogReader {
- type Item = ForwardLogReaderItem;
+ type Item = Record;
fn next(&mut self) -> Option<Self::Item> {
match self.read_record() {
@@ -118,10 +83,10 @@ mod tests {
// There are two records in the log with "schema" with one field: Bytes
- let first_item = forward_log_reader
+ let first_record = forward_log_reader
.next()
.expect("Failed to read the first record");
- assert!(match first_item.record.values() {
+ assert!(match first_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 8340cdd..9d7e1e8 100644
--- a/log_db/src/memtable_primary.rs
+++ b/log_db/src/memtable_primary.rs
@@ -2,9 +2,14 @@ use super::common::*;
use std::collections::BTreeMap;
pub struct PrimaryMemtable {
- /// Map of record locations indexed by key. Use the `LogKey` values to look up
- /// the actual records in the log files.
- records: BTreeMap<IndexableValue, LogKey>,
+ /// 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<IndexableValue, Record>,
}
impl PrimaryMemtable {
@@ -14,11 +19,15 @@ impl PrimaryMemtable {
}
}
- pub fn set(&mut self, key: &IndexableValue, value: &LogKey) {
+ pub fn set(&mut self, key: &IndexableValue, value: &Record) {
self.records.insert(key.clone(), value.clone());
}
- pub fn get(&self, key: &IndexableValue) -> Option<&LogKey> {
+ pub fn get(&mut self, key: &IndexableValue) -> Option<&Record> {
+ self.records.get(key)
+ }
+
+ pub fn get_without_update(&self, key: &IndexableValue) -> Option<&Record> {
self.records.get(key)
}
}
diff --git a/log_db/src/memtable_secondary.rs b/log_db/src/memtable_secondary.rs
index 8300ada..3f9863b 100644
--- a/log_db/src/memtable_secondary.rs
+++ b/log_db/src/memtable_secondary.rs
@@ -1,15 +1,14 @@
use super::*;
-use once_cell::sync::Lazy;
use std::collections::BTreeMap;
use std::collections::HashSet;
pub struct SecondaryMemtable {
- /// Map of records indexed by key. Values are non-empty sets of `LogKey` values.
- records: BTreeMap<IndexableValue, LogKeySet>,
+ /// 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<IndexableValue, HashSet<IndexableValue>>,
}
-static EMPTY_SET: Lazy<HashSet<LogKey>> = Lazy::new(|| HashSet::new());
-
impl SecondaryMemtable {
pub fn new() -> SecondaryMemtable {
SecondaryMemtable {
@@ -17,12 +16,10 @@ impl SecondaryMemtable {
}
}
- pub fn set(&mut self, key: &IndexableValue, value: &LogKey) {
+ pub fn set(&mut self, key: &IndexableValue, value: &IndexableValue) {
debug!(
- "Inserting/updating record in secondary memtable with key {:?} = segment {} index {}",
- &key,
- &value.segment_num(),
- &value.index()
+ "Inserting/updating record in secondary memtable with key {:?} = {:?}",
+ &key, &value,
);
match self.records.get_mut(key) {
@@ -35,27 +32,44 @@ impl SecondaryMemtable {
}
None => {
debug!("No existing entry found, creating one.");
- let set = LogKeySet::new_with_initial(&value);
+ let mut set = HashSet::with_capacity(1);
+ set.insert(value.clone());
self.records.insert(key.clone(), set);
}
}
}
- pub fn set_all(&mut self, key: &IndexableValue, values: &[LogKey]) {
+ pub fn set_all(&mut self, key: &IndexableValue, values: &[IndexableValue]) {
debug!(
"Replacing set of records in secondary memtable with key {:?} ({} values)",
&key,
&values.len(),
);
- let set = LogKeySet::from_slice(values);
+ let mut set = HashSet::with_capacity(values.len());
+ values.iter().for_each(|value| {
+ set.insert(value.clone());
+ });
+
self.records.insert(key.clone(), set);
}
- pub fn find_all(&mut self, key: &IndexableValue) -> &HashSet<LogKey> {
+ pub fn find_all(
+ &mut self,
+ primary_memtable: &PrimaryMemtable,
+ key: &IndexableValue,
+ ) -> Vec<Record> {
match self.records.get(key) {
- Some(set) => set.log_keys(),
- None => &EMPTY_SET,
+ None => vec![],
+ Some(set) => set
+ .iter()
+ .map(|key| {
+ primary_memtable
+ .get_without_update(key)
+ .expect("Record not found")
+ .clone()
+ })
+ .collect(),
}
}
}
diff --git a/log_db/tests/integration.rs b/log_db/tests/integration.rs
index f178738..fa81542 100644
--- a/log_db/tests/integration.rs
+++ b/log_db/tests/integration.rs
@@ -207,6 +207,7 @@ 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()
@@ -243,6 +244,10 @@ 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()))