aboutsummaryrefslogtreecommitdiffstats
diff options
context:
space:
mode:
authorJan Tuomi <jan@jantuomi.fi>2024-11-21 18:49:12 +0200
committerJan Tuomi <jan@jantuomi.fi>2024-11-21 18:49:12 +0200
commit34730bbe58b77c7ee3d5e6344846a8a8e6595f17 (patch)
tree5f1ee575dde374cffac59a4677fdf921f88b34b3
parent1ee30efe415695e1b7d9b6390778a50e8168df04 (diff)
Refactor compact
-rw-r--r--log_db/src/common.rs6
-rw-r--r--log_db/src/lib.rs228
-rw-r--r--log_db/src/log_reader_forward.rs4
3 files changed, 132 insertions, 106 deletions
diff --git a/log_db/src/common.rs b/log_db/src/common.rs
index 8f55ab2..74bf97f 100644
--- a/log_db/src/common.rs
+++ b/log_db/src/common.rs
@@ -542,11 +542,7 @@ pub fn create_segment_metadata_file(
let metadata_filename = format!("metadata.{}", new_num);
let metadata_path = data_dir_path.join(metadata_filename);
- let mut metadata_file = fs::OpenOptions::new()
- .create(true)
- .write(true)
- .append(true)
- .open(&metadata_path)?;
+ let mut metadata_file = APPEND_MODE.clone().create(true).open(&metadata_path)?;
let metadata_header = MetadataHeader {
version: 1,
diff --git a/log_db/src/lib.rs b/log_db/src/lib.rs
index 8b342dd..f049295 100644
--- a/log_db/src/lib.rs
+++ b/log_db/src/lib.rs
@@ -612,130 +612,128 @@ impl<Field: Eq + Clone + Debug> DB<Field> {
/// You may call this function in a separate thread or process to avoid blocking the main thread.
/// However, the database will be exclusively locked, so all writes will be blocked during the tasks.
pub fn do_maintenance_tasks(&mut self) -> Result<(), io::Error> {
- let active_log_path = &self.data_dir.join(ACTIVE_SYMLINK_FILENAME);
- let active_target = fs::read_link(&active_log_path)?;
- let active_metadata_path = &self.data_dir.join(active_target);
- let active_log_md = fs::metadata(&active_log_path)?;
+ request_exclusive_lock(&self.data_dir, &mut self.active_metadata_file)?;
- let mut already_locked = false;
- match is_metadata_file_valid(&mut self.active_metadata_file)? {
- IsMetadatafileValidResult::Ok => {}
- _ => {
- debug!("Active metadata is invalid, acquiring exclusive lock...");
- request_exclusive_lock(&self.data_dir, &mut self.active_metadata_file)?;
- already_locked = true;
+ ensure_active_metadata_is_valid(&self.data_dir, &mut self.active_metadata_file)?;
- ensure_active_metadata_is_valid(&self.data_dir, &mut self.active_metadata_file)?;
- }
- };
+ let metadata_size = self.active_metadata_file.seek(SeekFrom::End(0))?;
+ if metadata_size >= self.config.segment_size as u64 {
+ debug!("Active log size exceeds threshold, starting rotation and compaction...");
- if active_log_md.size() >= self.config.segment_size as u64 {
- // Rotate the active log file
+ self.active_data_file.lock_shared()?;
+ let original_data_len = self.active_data_file.seek(SeekFrom::End(0))?;
- debug!("Starting rotation");
- if !already_locked {
- debug!("Requesting exclusive lock on active log file...");
- request_exclusive_lock(&self.data_dir, &mut self.active_metadata_file)?;
+ debug!("Reading segment data into a BTreeMap");
+ let mut map = BTreeMap::new();
+ let forward_log_reader = ForwardLogReader::new(
+ self.active_metadata_file.try_clone()?,
+ self.active_data_file.try_clone()?,
+ );
+
+ self.active_data_file.unlock()?;
+
+ let mut read_n = 0;
+ for (original_index, record) in forward_log_reader.enumerate() {
+ let primary_key = record
+ .at(self.primary_key_index)
+ .as_indexable()
+ .expect("Primary key was not indexable");
+ map.insert(primary_key, (original_index, record));
+ read_n += 1;
}
- debug!("Exclusive lock acquired, rotating active log file...");
+ debug!(
+ "Read {} records, out of which {} were unique",
+ read_n,
+ map.len()
+ );
+ debug!("Opening temporary files for writing compacted data");
+ // Create a new log data file
+ 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)?;
- // Create a new active log segment
- let (data_uuid, _) = create_segment_data_file(&self.data_dir)?;
- let (new_segment_num, _) = create_segment_metadata_file(&self.data_dir, &data_uuid)?;
- set_active_segment(&self.data_dir, new_segment_num)?;
+ 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)?;
- // The new active log file is not locked by this client so it cannot be touched.
- debug!("Active log file rotated, new segment: {}", new_segment_num);
+ debug!("Writing compacted data to temporary files");
+ let metadata_header = MetadataHeader {
+ version: 1,
+ uuid: new_data_uuid,
+ };
- let new_metadata_path = &self.data_dir.join(format!("metadata.{}", new_segment_num));
- self.active_metadata_file = APPEND_MODE.open(&new_metadata_path)?;
- let new_data_path = &self.data_dir.join(data_uuid.to_string());
- self.active_data_file = APPEND_MODE.open(&new_data_path)?;
+ temp_metadata_file.write_all(&metadata_header.serialize())?;
- // Compact the rotated segment without a lock.
- // Since the rotated segment and the compacted segment based on it will be
- // a) read-only, and b) identical in effective content, there is no need to lock it.
- self.compact_segment(&active_metadata_path)?;
+ let mut metadata_rows_buf = vec![0; metadata_size as usize - METADATA_FILE_HEADER_SIZE];
- // The new active log file is not locked by this client so it cannot be touched.
- debug!("Active log file rotated, new segment: {}", new_segment_num);
- debug!("Segment compacted");
- }
+ 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)?;
- self.active_metadata_file.unlock()?;
- self.active_data_file.unlock()?;
+ let metadata_offset = original_index * 16;
- Ok(())
- }
+ for (i, byte) in offset.to_be_bytes().iter().enumerate() {
+ metadata_rows_buf[metadata_offset + i] = *byte;
+ }
+ for (i, byte) in len.to_be_bytes().iter().enumerate() {
+ metadata_rows_buf[metadata_offset + 8 + i] = *byte;
+ }
- fn compact_segment(&self, metadata_path: &Path) -> Result<(), io::Error> {
- debug!("Opening segment file {:?} for compaction", metadata_path);
- let mut metadata_file = READ_MODE.open(metadata_path)?;
- let metadata_header = read_metadata_header(&mut metadata_file)?;
- validate_metadata_header(&metadata_header)?;
+ offset += len;
+ }
- let data_file_path = &self.data_dir.join(metadata_header.uuid.to_string());
- let data_file = READ_MODE.open(&data_file_path)?;
+ temp_metadata_file.write_all(&metadata_rows_buf)?;
- debug!("Reading segment data into a BTreeMap");
- let mut map = BTreeMap::new();
- let forward_log_reader = ForwardLogReader::new(metadata_file, data_file);
- for entry in forward_log_reader {
- let primary_key = entry
- .at(self.primary_key_index)
- .as_indexable()
- .expect("Primary key was not indexable");
- map.insert(primary_key, entry);
- }
+ // 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.
+ temp_metadata_file.flush()?;
+ new_data_file.flush()?;
- debug!("Opening temporary files for writing compacted data");
- let temp_data_file = tempfile::NamedTempFile::new()?;
- let temp_data_path = temp_data_file.as_ref();
- let mut temp_data_file = APPEND_MODE.open(temp_data_path)?;
+ let final_len = new_data_file.seek(io::SeekFrom::End(0))?;
- let temp_metadata_file = tempfile::NamedTempFile::new()?;
- let temp_metadata_path = temp_metadata_file.as_ref();
- let mut temp_metadata_file = APPEND_MODE.open(temp_metadata_path)?;
+ let active_num = greatest_segment_number(&self.data_dir)?;
- let new_data_uuid = Uuid::new_v4();
+ debug!("Moving temporary files to their final locations");
+ let new_data_path = &self.data_dir.join(new_data_uuid.to_string());
+ let active_metadata_path = &self.data_dir.join(format!("metadata.{}", active_num)); // overwrite active
- debug!("Writing compacted data to temporary files");
- let metadata_header = MetadataHeader {
- version: 1,
- uuid: new_data_uuid,
- };
+ fs::rename(&temp_metadata_path, &active_metadata_path)?;
- temp_metadata_file.write_all(&metadata_header.serialize())?;
+ debug!(
+ "Compaction complete, reduced data size: {} -> {}",
+ original_data_len, final_len
+ );
- let mut offset = 0u64;
- for entry in map.values() {
- let serialized = entry.serialize();
- let len = serialized.len() as u64;
- temp_data_file.write_all(&serialized)?;
+ let new_segment_num = active_num + 1;
+ let new_metadata_path = self.data_dir.join(format!("metadata.{}", new_segment_num));
+ let mut new_metadata_file =
+ APPEND_MODE.clone().create(true).open(&new_metadata_path)?;
- let mut metadata_buf = vec![];
- metadata_buf.extend(&offset.to_be_bytes());
- metadata_buf.extend(&len.to_be_bytes());
- temp_metadata_file.write_all(&metadata_buf)?;
+ let new_metadata_header = MetadataHeader {
+ version: 1,
+ uuid: new_data_uuid,
+ };
- offset += len;
- }
+ new_metadata_file.write_all(&new_metadata_header.serialize())?;
- // 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.
- temp_metadata_file.flush()?;
- temp_data_file.flush()?;
+ set_active_segment(&self.data_dir, new_segment_num)?;
+ self.active_metadata_file.unlock()?;
- let final_len = temp_metadata_file.seek(io::SeekFrom::End(0))?;
+ self.active_metadata_file = APPEND_MODE.open(&new_metadata_path)?;
+ self.active_data_file = APPEND_MODE.open(&new_data_path)?;
- debug!("Moving temporary files to their final locations");
- let target_data_file_path = &self.data_dir.join(new_data_uuid.to_string());
- fs::rename(&temp_data_path, &target_data_file_path)?;
- fs::rename(&temp_metadata_path, metadata_path)?;
+ // 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
+ );
+ } else {
+ self.active_metadata_file.unlock()?;
+ }
- debug!("Compaction complete, resulting size: {}", final_len);
Ok(())
}
}
@@ -771,6 +769,15 @@ mod tests {
db.upsert(&record).expect("Failed to insert record");
}
+ let mut segment1_file = READ_MODE.open(data_dir.join("metadata.1")).unwrap();
+ let segment1_metadata_size_original = segment1_file.seek(io::SeekFrom::End(0)).unwrap();
+
+ let segment1_header = read_metadata_header(&mut segment1_file).unwrap();
+ let mut segment1_data_file = READ_MODE
+ .open(data_dir.join(segment1_header.uuid.to_string()))
+ .unwrap();
+ let segment1_data_size_original = segment1_data_file.seek(io::SeekFrom::End(0)).unwrap();
+
// Rotate and compact
db.do_maintenance_tasks()
.expect("Failed to do maintenance tasks");
@@ -785,9 +792,32 @@ mod tests {
// Note negation here
assert!(!fs::exists(data_dir.join("metadata.3")).unwrap());
- // Check that the first segment was compacted
- let segment1_size = fs::metadata(data_dir.join("metadata.1")).unwrap().len();
- assert_eq!(segment1_size, 2 * 8 + METADATA_FILE_HEADER_SIZE as u64);
+ // Check that the compacted metadata file has the same size
+ let mut segment1_metadata_file_compacted =
+ READ_MODE.open(data_dir.join("metadata.1")).unwrap();
+ let segment1_metadata_size_compacted = segment1_metadata_file_compacted
+ .seek(io::SeekFrom::End(0))
+ .unwrap();
+ assert_eq!(
+ segment1_metadata_size_compacted,
+ segment1_metadata_size_original
+ );
+
+ // Check that the compacted data file is smaller
+ let segment1_header_compacted =
+ read_metadata_header(&mut segment1_metadata_file_compacted).unwrap();
+ let mut segment1_data_file_compacted = READ_MODE
+ .open(data_dir.join(segment1_header_compacted.uuid.to_string()))
+ .unwrap();
+ let segment1_data_size_compacted = segment1_data_file_compacted
+ .seek(io::SeekFrom::End(0))
+ .unwrap();
+ assert!(
+ segment1_data_size_compacted < segment1_data_size_original,
+ "Original: {}, Compacted: {}",
+ segment1_data_size_original,
+ segment1_data_size_compacted
+ );
// Check that the records can be read
let rec0 = db
diff --git a/log_db/src/log_reader_forward.rs b/log_db/src/log_reader_forward.rs
index 343d7ff..0e681b9 100644
--- a/log_db/src/log_reader_forward.rs
+++ b/log_db/src/log_reader_forward.rs
@@ -37,8 +37,8 @@ impl<'a> ForwardLogReader {
// Use .seek_relative instead of .seek to avoid dropping the BufReader internal buffer when
// the seek distance is small
- let seek_distance = entry_offset - self.data_reader.stream_position()?;
- self.data_reader.seek_relative(seek_distance as i64)?;
+ let seek_distance = entry_offset as i64 - self.data_reader.stream_position()? as i64;
+ self.data_reader.seek_relative(seek_distance)?;
let mut result_buf = vec![0; entry_length as usize];
self.data_reader.read_exact(&mut result_buf)?;