diff options
Diffstat (limited to 'log_db/src')
| -rw-r--r-- | log_db/src/lib.rs | 106 | ||||
| -rw-r--r-- | log_db/src/primary_memtable.rs | 31 | ||||
| -rw-r--r-- | log_db/src/secondary_memtable.rs | 36 |
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"); + } + } + } } |
