diff options
| author | Jan Tuomi <jan@jantuomi.fi> | 2024-11-21 18:49:12 +0200 |
|---|---|---|
| committer | Jan Tuomi <jan@jantuomi.fi> | 2024-11-21 18:49:12 +0200 |
| commit | 34730bbe58b77c7ee3d5e6344846a8a8e6595f17 (patch) | |
| tree | 5f1ee575dde374cffac59a4677fdf921f88b34b3 | |
| parent | 1ee30efe415695e1b7d9b6390778a50e8168df04 (diff) | |
Refactor compact
| -rw-r--r-- | log_db/src/common.rs | 6 | ||||
| -rw-r--r-- | log_db/src/lib.rs | 228 | ||||
| -rw-r--r-- | log_db/src/log_reader_forward.rs | 4 |
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)?; |
