diff options
| -rw-r--r-- | log_db/benches/benchmark.rs | 170 | ||||
| -rw-r--r-- | log_db/benches/utils.rs | 22 | ||||
| -rw-r--r-- | log_db/src/lib.rs | 119 | ||||
| -rw-r--r-- | log_db/src/record.rs | 24 |
4 files changed, 190 insertions, 145 deletions
diff --git a/log_db/benches/benchmark.rs b/log_db/benches/benchmark.rs index 38f1eb3..0f72aed 100644 --- a/log_db/benches/benchmark.rs +++ b/log_db/benches/benchmark.rs @@ -5,67 +5,57 @@ use log_db::*; use tempfile; use utils::*; -pub fn upsert_various_initial_sizes(c: &mut Criterion) { - let mut group = c.benchmark_group("upsert_various_initial_sizes"); +pub fn upsert_compacted(c: &mut Criterion) { + let mut group = c.benchmark_group("upsert_compacted"); + let data_dir_obj = tempfile::tempdir().expect("Failed to get tmpdir"); + let data_dir = &data_dir_obj + .path() + .to_str() + .expect("Failed to convert tmpdir path to str"); + let mut db = DB::<Inst>::configure() + .data_dir(&data_dir) + .initialize() + .expect("Failed to initialize DB"); - for size in [0, 1000, 10_000, 100_000, 1_000_000, 10_000_000] { - let data_dir_obj = tempfile::tempdir().expect("Failed to get tmpdir"); - let data_dir = &data_dir_obj - .path() - .to_str() - .expect("Failed to convert tmpdir path to str"); - let mut db = DB::<Inst>::configure() - .data_dir(&data_dir) - .initialize() - .expect("Failed to initialize DB"); - - prefill_db(&mut db, size, false).expect("Failed to prefill DB"); + let mut insts = Vec::new(); + for size in [100_000, 1_000_000, 10_000_000] { + println!("Prefilling DB to {} entries", size); + prefill_db(&mut db, &mut insts, size, true).expect("Failed to prefill DB"); group.bench_with_input(BenchmarkId::from_parameter(size), &size, |b, &_size| { b.iter(|| { let inst = random_inst(0, size as i64 + 1); - let _ = db.upsert(black_box(inst)); + let result = db.upsert(black_box(inst)).unwrap(); + assert!(result == ()); }); }); } } -pub fn upsert_various_initial_sizes_compacted(c: &mut Criterion) { - let mut group = c.benchmark_group("upsert_various_initial_sizes_compacted"); - - for size in [0, 1000, 10_000, 100_000, 1_000_000, 10_000_000] { - let data_dir_obj = tempfile::tempdir().expect("Failed to get tmpdir"); - let data_dir = &data_dir_obj - .path() - .to_str() - .expect("Failed to convert tmpdir path to str"); - - let record_length = 1 + // tombstone tag - 1 + 8 + // int tag + int value - 1 + 8 + 5 + // string tag + string length + string value - 1 + 8 + 10; // bytes tag + bytes length + bytes value - - let mut db = DB::<Inst>::configure() - .data_dir(&data_dir) - .segment_size(1000 * record_length) - .initialize() - .expect("Failed to initialize DB"); +pub fn delete_existing_compacted(c: &mut Criterion) { + let mut group = c.benchmark_group("delete_existing_compacted"); + let data_dir_obj = tempfile::tempdir().expect("Failed to get tmpdir"); + let data_dir = &data_dir_obj + .path() + .to_str() + .expect("Failed to convert tmpdir path to str"); + let mut db = DB::<Inst>::configure() + .data_dir(&data_dir) + .initialize() + .expect("Failed to initialize DB"); - prefill_db(&mut db, size, true).expect("Failed to prefill DB"); - db.do_maintenance_tasks() - .expect("Failed to do maintenance tasks"); + let mut insts = Vec::new(); + for size in [100_000, 1_000_000, 10_000_000] { + println!("Prefilling DB to {} entries", size); + prefill_db(&mut db, &mut insts, size, true).expect("Failed to prefill DB"); + let mut inst_it = insts.iter(); group.bench_with_input(BenchmarkId::from_parameter(size), &size, |b, &_size| { - let mut i = 0; b.iter(|| { - let inst = random_inst(0, size as i64 + 1); - let _ = db.upsert(black_box(inst)); - - if i % 100 == 0 { - db.do_maintenance_tasks() - .expect("Failed to do maintenance tasks"); - } - i += 1; + let result = db + .delete(black_box(&Value::Int(inst_it.next().unwrap().id))) + .unwrap(); + assert!(result.is_some()) }); }); } @@ -89,56 +79,72 @@ pub fn upsert_write_durability(c: &mut Criterion) { b.iter(|| { let inst = random_inst(0, 1000); - let _ = db.upsert(black_box(inst)); + let result = db.upsert(black_box(inst)).unwrap(); + assert!(result == ()); }); }); } } -pub fn get_from_disk_various_initial_sizes(c: &mut Criterion) { - let mut group = c.benchmark_group("get_from_disk_various_initial_sizes"); +pub fn get_existing_compacted(c: &mut Criterion) { + let mut group = c.benchmark_group("get_existing_compacted"); + + let data_dir_obj = tempfile::tempdir().expect("Failed to get tmpdir"); + let data_dir = &data_dir_obj + .path() + .to_str() + .expect("Failed to convert tmpdir path to str"); + let mut db = DB::<Inst>::configure() + .data_dir(&data_dir) + .initialize() + .expect("Failed to initialize DB"); - for size in [0, 1000, 10_000, 100_000, 1_000_000, 10_000_000] { - let data_dir_obj = tempfile::tempdir().expect("Failed to get tmpdir"); - let data_dir = &data_dir_obj - .path() - .to_str() - .expect("Failed to convert tmpdir path to str"); - let mut db = DB::<Inst>::configure() - .data_dir(&data_dir) - .initialize() - .expect("Failed to initialize DB"); - prefill_db(&mut db, size, false).expect("Failed to prefill DB"); + let mut insts = Vec::new(); + for size in [100_000, 1_000_000, 10_000_000] { + println!("Prefilling DB to {} entries", size); + prefill_db(&mut db, &mut insts, size, true).expect("Failed to prefill DB"); + let mut inst_it = insts.iter(); group.bench_with_input(BenchmarkId::from_parameter(size), &size, |b, &_size| { b.iter(|| { - let id = random_int(0, size as i64 + 1); - let _ = db.get(black_box(&Value::Int(id))); + let result = db + .get(black_box(&Value::Int(inst_it.next().unwrap().id))) + .unwrap(); + assert!(result.is_some()) }); }); } } -pub fn get_from_disk_various_initial_sizes_compacted(c: &mut Criterion) { - let mut group = c.benchmark_group("get_from_disk_various_initial_sizes_compacted"); +pub fn find_by_existing_compacted(c: &mut Criterion) { + let mut group = c.benchmark_group("find_by_existing_compacted"); - for size in [0, 1000, 10_000, 100_000, 1_000_000, 10_000_000] { - let data_dir_obj = tempfile::tempdir().expect("Failed to get tmpdir"); - let data_dir = &data_dir_obj - .path() - .to_str() - .expect("Failed to convert tmpdir path to str"); - let mut db = DB::<Inst>::configure() - .data_dir(&data_dir) - .initialize() - .expect("Failed to initialize DB"); + let data_dir_path = tempfile::tempdir() + .expect("Failed to get tmpdir") + .into_path(); + let data_dir = data_dir_path + .to_str() + .expect("Failed to convert tmpdir path to str"); + let mut db = DB::<Inst>::configure() + .data_dir(&data_dir) + .initialize() + .expect("Failed to initialize DB"); - prefill_db(&mut db, size, true).expect("Failed to prefill DB"); + let mut insts = Vec::new(); + let mut inst_index = 0; + for size in [100_000, 1_000_000, 10_000_000] { + println!("Prefilling DB to {} entries", size); + prefill_db(&mut db, &mut insts, size, true).expect("Failed to prefill DB"); + println!("DB prefilling done, insts.len = {}", insts.len()); group.bench_with_input(BenchmarkId::from_parameter(size), &size, |b, &_size| { b.iter(|| { - let id = random_int(0, size as i64 + 1); - let _ = db.get(black_box(&Value::Int(id))); + let name = insts[inst_index].name.clone(); + let result = db + .find_by(black_box(&Field::Name), black_box(&Value::String(name))) + .unwrap(); + assert!(result.len() > 0); + inst_index = (inst_index + 1) % insts.len(); }); }); } @@ -147,10 +153,10 @@ pub fn get_from_disk_various_initial_sizes_compacted(c: &mut Criterion) { // Register the benchmark group criterion_group!( benches, - upsert_various_initial_sizes, - upsert_various_initial_sizes_compacted, + upsert_compacted, + delete_existing_compacted, upsert_write_durability, - get_from_disk_various_initial_sizes, - get_from_disk_various_initial_sizes_compacted, + get_existing_compacted, + find_by_existing_compacted, ); criterion_main!(benches); diff --git a/log_db/benches/utils.rs b/log_db/benches/utils.rs index 05ea30f..0f9c620 100644 --- a/log_db/benches/utils.rs +++ b/log_db/benches/utils.rs @@ -12,9 +12,9 @@ pub enum Field { #[derive(PartialEq, Eq, Debug, Clone)] pub struct Inst { - id: i64, - name: String, - data: Vec<u8>, + pub id: i64, + pub name: String, + pub data: Vec<u8>, } impl Recordable for Inst { type Field = Field; @@ -28,6 +28,10 @@ impl Recordable for Inst { fn primary_key() -> Self::Field { Field::Id } + fn secondary_keys() -> Vec<Self::Field> { + vec![Field::Name] + } + fn into_record(self) -> Vec<Value> { vec![ Value::Int(self.id), @@ -81,11 +85,17 @@ pub fn random_inst(from_id: i64, to_id: i64) -> Inst { } } -pub fn prefill_db(db: &mut DB<Inst>, n_records: usize, compact: bool) -> Result<(), DBError> { - for _ in 0..n_records { +pub fn prefill_db( + db: &mut DB<Inst>, + insts: &mut Vec<Inst>, + n_records: usize, + compact: bool, +) -> Result<(), DBError> { + for i in 0..(n_records - insts.len()) { let inst = random_inst(0, n_records as i64); + insts.push(inst.clone()); db.upsert(inst)?; - if compact { + if i % 1000 == 0 && compact { db.do_maintenance_tasks()?; } } diff --git a/log_db/src/lib.rs b/log_db/src/lib.rs index eb9cfe9..918f5e1 100644 --- a/log_db/src/lib.rs +++ b/log_db/src/lib.rs @@ -428,9 +428,7 @@ impl<R: Recordable> DB<R> { self.insert_record_to_memtables(&log_key, &record); - // These post-condition asserts are commented out since they seemed - // to sometimes report false positives. - // + // 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); @@ -810,50 +808,78 @@ impl<R: Recordable> DB<R> { self.active_data_file.lock_shared()?; let original_data_len = self.active_data_file.seek(SeekFrom::End(0))?; - let metadata_size = self.active_metadata_file.seek(SeekFrom::End(0))?; + + let active_target = fs::read_link(&self.data_dir.join(ACTIVE_SYMLINK_FILENAME))?; + let active_num = parse_segment_number(&active_target)?; debug!("Reading segment data into a BTreeMap"); - let mut map = BTreeMap::new(); - let forward_log_reader = ForwardLogReader::new( + let mut pk_to_item_map = BTreeMap::new(); + let forward_read_items: Vec<(IndexableValue, Record)> = ForwardLogReader::new( self.active_metadata_file.try_clone()?, self.active_data_file.try_clone()?, - ); + ) + .map(|item| { + ( + item.record + .at(self.primary_key_index) + .as_indexable() + .expect("Primary key was not indexable"), + item.record, + ) + }) + .collect(); self.active_data_file.unlock()?; - let mut read_n = 0; - for (original_index, item) in forward_log_reader.enumerate() { - let pk = item - .record - .at(self.primary_key_index) - .as_indexable() - .expect("Primary key was not indexable"); - + for (pk, record) in forward_read_items.iter() { // If the record is a tombstone, remove the PK from the map - if item.record.tombstone { - map.remove(&pk); + if record.tombstone { + pk_to_item_map.remove(pk); } else { - map.insert(pk, (original_index, item.record)); + pk_to_item_map.insert(pk, record); } - - read_n += 1; } debug!( "Read {} records, out of which {} were unique", - read_n, - map.len() + forward_read_items.len(), + pk_to_item_map.len() ); - debug!("Opening temporary files for writing compacted data"); - // Create a new log data file + + // Create a new log data file and write it + debug!("Opening new data file and writing compacted data"); let (new_data_uuid, new_data_path) = create_segment_data_file(&self.data_dir)?; let mut new_data_file = APPEND_MODE.open(&new_data_path)?; + let mut pk_to_data_map = BTreeMap::new(); + let mut offset = 0u64; + for (pk, record) in pk_to_item_map.into_iter() { + let serialized = record.serialize(); + let len = serialized.len() as u64; + new_data_file.write_all(&serialized)?; + + pk_to_data_map.insert(pk, (offset, len)); + offset += len; + } + + // Sync the data file to disk. + // This is fine to do without consulting WriteDurability because this is a one-off + // operation that is not part of the normal write path. + new_data_file.flush()?; + new_data_file.sync_all()?; + + let final_data_len = new_data_file.seek(io::SeekFrom::End(0))?; + debug!( + "Wrote compacted data, reduced data size: {} -> {}", + original_data_len, final_data_len + ); + + // Create a new log metadata file and write it + debug!("Opening temp metadata file and writing pointers to compacted data file"); let temp_metadata_file = tempfile::NamedTempFile::new()?; let temp_metadata_path = temp_metadata_file.as_ref(); let mut temp_metadata_file = WRITE_MODE.open(temp_metadata_path)?; - debug!("Writing compacted data to temporary files"); let metadata_header = MetadataHeader { version: 1, uuid: new_data_uuid, @@ -861,37 +887,19 @@ impl<R: Recordable> DB<R> { temp_metadata_file.write_all(&metadata_header.serialize())?; - let mut metadata_rows_buf = vec![0; metadata_size as usize - METADATA_FILE_HEADER_SIZE]; - - let mut offset = 0u64; - for (original_index, record) in map.values() { - let serialized = record.serialize(); - let len = serialized.len() as u64; - new_data_file.write_all(&serialized)?; - - let metadata_offset = original_index * 16; + for (pk, _) in forward_read_items.iter() { + let (offset, len) = pk_to_data_map.get(&pk).unwrap(); - for (i, byte) in offset.to_be_bytes().into_iter().enumerate() { - metadata_rows_buf[metadata_offset + i] = byte; - } - for (i, byte) in len.to_be_bytes().into_iter().enumerate() { - metadata_rows_buf[metadata_offset + 8 + i] = byte; - } + let mut metadata_buf = vec![]; + metadata_buf.extend(offset.to_be_bytes().into_iter()); + metadata_buf.extend(len.to_be_bytes().into_iter()); - offset += len; + temp_metadata_file.write_all(&metadata_buf)?; } - temp_metadata_file.write_all(&metadata_rows_buf)?; - - // Sync the temporary files to disk - // This is fine to do without consulting WriteDurability because this is a one-off - // operation that is not part of the normal write path. + // Sync the metadata file to disk, see comment above about sync. temp_metadata_file.flush()?; - new_data_file.flush()?; - - let final_len = new_data_file.seek(io::SeekFrom::End(0))?; - - let active_num = greatest_segment_number(&self.data_dir)?; + temp_metadata_file.sync_all()?; debug!("Moving temporary files to their final locations"); let new_data_path = &self.data_dir.join(new_data_uuid.to_string()); @@ -899,10 +907,7 @@ impl<R: Recordable> DB<R> { fs::rename(&temp_metadata_path, &active_metadata_path)?; - debug!( - "Compaction complete, reduced data size: {} -> {}", - original_data_len, final_len - ); + debug!("Compaction complete, creating new segment"); let new_segment_num = active_num + 1; let new_metadata_path = self.data_dir.join(metadata_filename(new_segment_num)); @@ -925,8 +930,8 @@ impl<R: Recordable> DB<R> { // The new active log file is not locked by this client so it cannot be touched. debug!( - "Active log file rotated and compacted, new segment: {}", - new_segment_num + "Active log file {} rotated and compacted, new segment: {}", + active_num, new_segment_num ); Ok(()) diff --git a/log_db/src/record.rs b/log_db/src/record.rs index 61939e1..4a30902 100644 --- a/log_db/src/record.rs +++ b/log_db/src/record.rs @@ -120,3 +120,27 @@ pub trait Recordable { /// Convert a vector of database values into the data structure implementing the `Recordable` trait. fn from_record(record: Vec<Value>) -> Self; } + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn test_record_serialize_deserialize() { + let record = Record { + values: vec![ + Value::Int(1), + Value::String("hello".to_string()), + Value::Bytes(vec![0, 1, 2, 3]), + ], + tombstone: true, + }; + + let serialized = record.serialize(); + let deserialized = Record::deserialize(&serialized); + let reserialized = deserialized.serialize(); + + assert_eq!(serialized.len(), reserialized.len()); + assert_eq!(record.values, deserialized.values); + } +} |
