aboutsummaryrefslogtreecommitdiffstats
path: root/log_db/src
diff options
context:
space:
mode:
authorJan Tuomi <jan@jantuomi.fi>2024-10-10 23:41:10 +0300
committerJan Tuomi <jan@jantuomi.fi>2024-10-10 23:41:10 +0300
commitee0a81cf0277a7ff8683f94cc95a23f29996d980 (patch)
tree263ae30f5057c402bf9d73ae6e26e08ad0ec22e8 /log_db/src
parent336afff93bdd09cfaad75bde25d6d6627ef7c24f (diff)
Add invariant that memtables should always have the same set of primary key references
Diffstat (limited to 'log_db/src')
-rw-r--r--log_db/src/lib.rs106
-rw-r--r--log_db/src/primary_memtable.rs31
-rw-r--r--log_db/src/secondary_memtable.rs36
3 files changed, 139 insertions, 34 deletions
diff --git a/log_db/src/lib.rs b/log_db/src/lib.rs
index d704c7d..bcdaefe 100644
--- a/log_db/src/lib.rs
+++ b/log_db/src/lib.rs
@@ -234,7 +234,7 @@ impl<Field: Eq + Clone + Debug> DB<Field> {
let secondary_memtables = config
.secondary_keys
.iter()
- .map(|key| SecondaryMemtable::new(key))
+ .map(|key| SecondaryMemtable::new(&config.fields, key, primary_key_index))
.collect();
let mut db = DB::<Field> {
@@ -251,8 +251,7 @@ impl<Field: Eq + Clone + Debug> DB<Field> {
let forward_log_reader = ForwardLogReader::new(&mut file);
for record in forward_log_reader {
- db.update_primary_index(&record);
- db.update_secondary_indexes(&record);
+ db.insert_to_memtables(&record);
}
info!("Database ready.");
@@ -351,11 +350,8 @@ impl<Field: Eq + Clone + Debug> DB<Field> {
debug!("Record appended to log file, lock released");
- debug!("Updating primary memtable");
- self.update_primary_index(record);
-
- debug!("Updating secondary memtables");
- self.update_secondary_indexes(record);
+ debug!("Updating memtables");
+ self.insert_to_memtables(record);
Ok(())
}
@@ -456,11 +452,8 @@ impl<Field: Eq + Clone + Debug> DB<Field> {
debug!("Found matching record in log file.");
- debug!("Updating primary memtable");
- self.update_primary_index(&result_value);
-
- debug!("Updating secondary memtables");
- self.update_secondary_indexes(&result_value);
+ debug!("Updating memtables");
+ self.insert_to_memtables(&result_value);
Ok(result)
}
@@ -593,14 +586,30 @@ impl<Field: Eq + Clone + Debug> DB<Field> {
}
}
- fn update_primary_index(&mut self, record: &Record) {
+ fn insert_to_memtables(&mut self, record: &Record) {
let key = record.values[self.primary_key_index]
.as_indexable()
.expect("A non-indexable value was stored at key index");
+
+ if self.primary_memtable.capacity == 0 {
+ return;
+ }
+
+ debug!(
+ "Inserting/updating record in primary memtable with key {:?} = {:?}",
+ &key, &record,
+ );
+
+ if let Some(evicted) = self.primary_memtable.evict_if_necessary() {
+ self.secondary_memtables
+ .iter_mut()
+ .for_each(|secondary_memtable| {
+ secondary_memtable.remove(&evicted);
+ });
+ }
+
self.primary_memtable.set(&key, record);
- }
- fn update_secondary_indexes(&mut self, record: &Record) {
self.secondary_memtables
.iter_mut()
.for_each(|secondary_memtable| {
@@ -840,3 +849,68 @@ impl<Field: Eq + Clone + Debug> DB<Field> {
Ok(())
}
}
+
+#[cfg(test)]
+mod tests {
+ use super::*;
+ use rand::distributions::Alphanumeric;
+ use rand::Rng;
+ use std::collections::HashSet;
+ use tempfile::tempdir;
+
+ #[derive(Eq, PartialEq, Clone, Debug)]
+ enum Field {
+ Id,
+ Name,
+ Data,
+ }
+
+ fn tmp_dir() -> String {
+ let dir = tempdir()
+ .expect("Failed to create temporary directory")
+ .path()
+ .to_str()
+ .expect("Failed to convert temporary directory path to string")
+ .to_string();
+ fs::create_dir_all(&dir).expect("Failed to create temporary directory");
+ dir
+ }
+
+ #[test]
+ fn memtables_always_have_the_same_primary_keys() {
+ let data_dir = tmp_dir();
+
+ let mut db = DB::configure()
+ .data_dir(&data_dir)
+ .fields(vec![
+ (Field::Id, RecordField::int()),
+ (Field::Name, RecordField::string()),
+ ])
+ .primary_key(Field::Id)
+ .secondary_keys(vec![Field::Name])
+ .initialize()
+ .expect("Failed to initialize DB instance");
+
+ let mut rng = rand::thread_rng();
+ for _ in 0..100 {
+ let id = rng.gen_range(0..100);
+ let name = (0..5).map(|_| rng.sample(Alphanumeric) as char).collect();
+
+ let record = Record {
+ values: vec![RecordValue::Int(id), RecordValue::String(name)],
+ };
+ db.upsert(&record).expect("Failed to upsert record");
+
+ let p_set: HashSet<&IndexableValue> = db.primary_memtable.records.keys().collect();
+ let mut s_set: HashSet<&IndexableValue> = HashSet::new();
+
+ for table in db.secondary_memtables.iter() {
+ table.records.values().for_each(|r| {
+ s_set.extend(r);
+ });
+ }
+
+ assert_eq!(p_set, s_set);
+ }
+ }
+}
diff --git a/log_db/src/primary_memtable.rs b/log_db/src/primary_memtable.rs
index 141238f..180ae72 100644
--- a/log_db/src/primary_memtable.rs
+++ b/log_db/src/primary_memtable.rs
@@ -6,7 +6,7 @@ pub struct PrimaryMemtable {
/// Maximum number of records that can be stored in the memtable
/// before evicting the oldest records. The oldest record is
/// determined by the `evict_policy`.
- capacity: usize,
+ pub capacity: usize,
/// 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
@@ -14,7 +14,7 @@ pub struct PrimaryMemtable {
///
/// Note: it must be invariant that all memtables (primary and secondary)
/// contain the same keys.
- records: BTreeMap<IndexableValue, Record>,
+ pub records: BTreeMap<IndexableValue, Record>,
/// A max heap priority queue of keys. The record with least priority is evicted
/// from the primary memtable and any secondary memtables that reference it, when
/// the memtable reaches capacity.
@@ -40,20 +40,6 @@ impl PrimaryMemtable {
}
pub fn set(&mut self, key: &IndexableValue, value: &Record) {
- if self.capacity == 0 {
- return;
- }
-
- debug!(
- "Inserting/updating record in primary memtable with key {:?} = {:?}",
- &key, &value,
- );
-
- if self.records.len() >= self.capacity {
- let (evict_key, _prio) = self.evict_queue.pop().expect("Evict queue was empty");
- self.records.remove(&evict_key);
- }
-
self.records.insert(key.clone(), value.clone());
if self.evict_policy == MemtableEvictPolicy::LeastWritten
@@ -94,4 +80,17 @@ impl PrimaryMemtable {
self.n_operations += 1;
ret
}
+
+ pub fn evict_if_necessary(&mut self) -> Option<Record> {
+ if self.records.len() >= self.capacity {
+ let (evict_key, _prio) = self.evict_queue.pop().expect("Evict queue was empty");
+ let removed = self
+ .records
+ .remove(&evict_key)
+ .expect("Key was not found in records");
+ Some(removed)
+ } else {
+ None
+ }
+ }
}
diff --git a/log_db/src/secondary_memtable.rs b/log_db/src/secondary_memtable.rs
index 194d4a0..994ee0d 100644
--- a/log_db/src/secondary_memtable.rs
+++ b/log_db/src/secondary_memtable.rs
@@ -5,17 +5,30 @@ use std::fmt::Debug;
pub struct SecondaryMemtable<Field: Eq + Clone + Debug> {
pub field: Field,
+ field_index: usize,
+ primary_key_index: usize,
/// 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>>,
+ pub records: BTreeMap<IndexableValue, HashSet<IndexableValue>>,
}
impl<Field: Eq + Clone + Debug> SecondaryMemtable<Field> {
- pub fn new(field: &Field) -> SecondaryMemtable<Field> {
+ pub fn new(
+ field_schema: &Vec<(Field, RecordField)>,
+ field: &Field,
+ primary_key_index: usize,
+ ) -> SecondaryMemtable<Field> {
+ let field_index = field_schema
+ .iter()
+ .position(|(f, _)| f == field)
+ .expect("Field not found in schema");
+
SecondaryMemtable {
field: field.clone(),
+ field_index,
+ primary_key_index,
records: BTreeMap::new(),
}
}
@@ -76,4 +89,23 @@ impl<Field: Eq + Clone + Debug> SecondaryMemtable<Field> {
.collect(),
}
}
+
+ pub fn remove(&mut self, record: &Record) {
+ let key = record.values[self.field_index]
+ .as_indexable()
+ .expect("Field is not indexable");
+
+ let primary_key = record.values[self.primary_key_index]
+ .as_indexable()
+ .expect("Primary key is not indexable");
+
+ match self.records.get_mut(&key) {
+ Some(set) => {
+ set.remove(&primary_key);
+ }
+ None => {
+ panic!("Record not found in secondary memtable");
+ }
+ }
+ }
}