aboutsummaryrefslogtreecommitdiffstats
path: root/log_db/src
diff options
context:
space:
mode:
Diffstat (limited to 'log_db/src')
-rw-r--r--log_db/src/common.rs4
-rw-r--r--log_db/src/lib.rs139
-rw-r--r--log_db/src/log_reader_reverse.rs118
-rw-r--r--log_db/src/memtable_primary.rs8
-rw-r--r--log_db/src/memtable_secondary.rs32
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
}