aboutsummaryrefslogtreecommitdiffstats
path: root/log_db
diff options
context:
space:
mode:
authorJan Tuomi <jan@jantuomi.fi>2025-01-08 10:13:00 +0200
committerJan Tuomi <jan@jantuomi.fi>2025-01-10 19:48:50 +0200
commit3ffa6702cdb39c21a26ddc434f46d17d0104cb7b (patch)
treeb0d8e4dcbb967d90b90a338eb0cb50fc0f10c518 /log_db
parentf19c7e6a41768e73a0f727deb22d65f1646a0777 (diff)
Fix compaction issue where get of valid logkey produced unset metadata row (len = 0)
Diffstat (limited to 'log_db')
-rw-r--r--log_db/benches/benchmark.rs170
-rw-r--r--log_db/benches/utils.rs22
-rw-r--r--log_db/src/lib.rs119
-rw-r--r--log_db/src/record.rs24
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);
+ }
+}