Skip to content

Commit 4b29aae

Browse files
committed
feat(metrics): compute compaction live-bytes from latest visible versions
- estimate compaction live bytes by scanning SSTables and selecting the latest version per user key - treat tombstones as non-live entries while preserving undecodable rows as live fallback bytes - update space metric recomputation to return errors and run after flush/compaction state transitions - avoid unnecessary space-metric recomputation on put/delete hot paths - extend compaction integration assertions to validate live/total byte invariants
1 parent e1594e5 commit 4b29aae

2 files changed

Lines changed: 129 additions & 9 deletions

File tree

src/storage/engine.rs

Lines changed: 117 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,4 @@
1-
use std::collections::HashSet;
1+
use std::collections::{HashMap, HashSet};
22
use std::path::{Path, PathBuf};
33
use std::sync::Arc;
44
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
@@ -357,7 +357,7 @@ impl StorageEngine {
357357

358358
{
359359
let mut state = engine.state.lock();
360-
engine.refresh_compaction_space_metrics_locked(&state);
360+
engine.refresh_compaction_space_metrics_locked(&state)?;
361361
engine.schedule_pending_flushes_locked(&mut state)?;
362362
engine.schedule_pending_compactions_locked(&mut state)?;
363363
}
@@ -389,7 +389,6 @@ impl StorageEngine {
389389
state.memtables.put(key, sequence, value);
390390
self.schedule_pending_flushes_locked(&mut state)?;
391391
self.schedule_pending_compactions_locked(&mut state)?;
392-
self.refresh_compaction_space_metrics_locked(&state);
393392
self.metrics.puts.fetch_add(1, Ordering::Relaxed);
394393
self.metrics
395394
.compaction
@@ -416,7 +415,6 @@ impl StorageEngine {
416415
state.memtables.delete(key, sequence);
417416
self.schedule_pending_flushes_locked(&mut state)?;
418417
self.schedule_pending_compactions_locked(&mut state)?;
419-
self.refresh_compaction_space_metrics_locked(&state);
420418
self.metrics.deletes.fetch_add(1, Ordering::Relaxed);
421419
self.metrics.compaction.lock().record_user_write(key.len() as u64);
422420

@@ -847,7 +845,7 @@ impl StorageEngine {
847845
state.memtables.remove_immutable(&completed.memtable);
848846
state.sstables.insert(0, SSTableRuntime { metadata, reader });
849847
self.metrics.flushes_completed.fetch_add(1, Ordering::Relaxed);
850-
self.refresh_compaction_space_metrics_locked(&state);
848+
self.refresh_compaction_space_metrics_locked(&state)?;
851849
self.schedule_pending_compactions_locked(&mut state)?;
852850
}
853851
FlushResponse::Empty { memtable } => {
@@ -933,7 +931,7 @@ impl StorageEngine {
933931
completed.task_id,
934932
&completed.input_tables,
935933
);
936-
self.refresh_compaction_space_metrics_locked(&state);
934+
self.refresh_compaction_space_metrics_locked(&state)?;
937935
self.schedule_pending_compactions_locked(&mut state)?;
938936
removed_file_names
939937
};
@@ -982,7 +980,7 @@ impl StorageEngine {
982980
.sstables
983981
.retain(|runtime| !input_ids.contains(&runtime.metadata.table_id));
984982
self.release_compaction_task_locked(&mut state, task_id, &input_tables);
985-
self.refresh_compaction_space_metrics_locked(&state);
983+
self.refresh_compaction_space_metrics_locked(&state)?;
986984
self.schedule_pending_compactions_locked(&mut state)?;
987985
}
988986

@@ -1008,10 +1006,15 @@ impl StorageEngine {
10081006
Ok(())
10091007
}
10101008

1011-
fn refresh_compaction_space_metrics_locked(&self, state: &EngineState) {
1009+
fn refresh_compaction_space_metrics_locked(
1010+
&self,
1011+
state: &EngineState,
1012+
) -> Result<(), EngineError> {
10121013
let total_bytes =
10131014
state.sstables.iter().map(|runtime| runtime.metadata.file_size_bytes).sum::<u64>();
1014-
self.metrics.compaction.lock().set_space_bytes(total_bytes, total_bytes);
1015+
let live_bytes = estimate_live_user_bytes_for_tables(&state.sstables)?;
1016+
self.metrics.compaction.lock().set_space_bytes(live_bytes, total_bytes);
1017+
Ok(())
10151018
}
10161019

10171020
fn release_compaction_task(&self, task_id: u64, input_tables: &[SSTableMetadata]) {
@@ -1290,6 +1293,60 @@ fn merge_compaction_rows(
12901293
collapsed
12911294
}
12921295

1296+
fn estimate_live_user_bytes_for_tables(
1297+
tables: &[SSTableRuntime],
1298+
) -> Result<u64, SSTableReadError> {
1299+
let mut latest = HashMap::<Vec<u8>, (u64, ValueType, u64)>::new();
1300+
let mut undecodable_live_bytes = 0_u64;
1301+
1302+
for table in tables {
1303+
let rows = table.reader.scan_range(None, None)?;
1304+
for (internal_key, value) in rows {
1305+
accumulate_live_user_bytes(
1306+
&mut latest,
1307+
&internal_key,
1308+
value.len(),
1309+
&mut undecodable_live_bytes,
1310+
);
1311+
}
1312+
}
1313+
1314+
let decoded_live_bytes = latest
1315+
.values()
1316+
.map(|(_, value_type, logical_bytes)| {
1317+
if *value_type == ValueType::Put {
1318+
*logical_bytes
1319+
} else {
1320+
0
1321+
}
1322+
})
1323+
.sum::<u64>();
1324+
1325+
Ok(decoded_live_bytes.saturating_add(undecodable_live_bytes))
1326+
}
1327+
1328+
fn accumulate_live_user_bytes(
1329+
latest: &mut HashMap<Vec<u8>, (u64, ValueType, u64)>,
1330+
internal_key: &[u8],
1331+
value_len: usize,
1332+
undecodable_live_bytes: &mut u64,
1333+
) {
1334+
let Some(decoded) = decode_internal_key(internal_key) else {
1335+
*undecodable_live_bytes = undecodable_live_bytes
1336+
.saturating_add((internal_key.len().saturating_add(value_len)) as u64);
1337+
return;
1338+
};
1339+
1340+
let logical_bytes = (decoded.user_key.len().saturating_add(value_len)) as u64;
1341+
let should_replace = latest
1342+
.get(decoded.user_key)
1343+
.map(|(sequence, _, _)| decoded.sequence > *sequence)
1344+
.unwrap_or(true);
1345+
if should_replace {
1346+
latest.insert(decoded.user_key.to_vec(), (decoded.sequence, decoded.value_type, logical_bytes));
1347+
}
1348+
}
1349+
12931350
fn recover_from_wal(
12941351
wal_dir: &Path,
12951352
memtables: &mut MemTableManager,
@@ -1462,4 +1519,55 @@ mod tests {
14621519
.expect("duplicate key row");
14631520
assert_eq!(dup_row.1, b"new".to_vec());
14641521
}
1522+
1523+
#[test]
1524+
fn live_user_bytes_prefers_latest_sequence_and_ignores_tombstone_payload() {
1525+
let mut latest = HashMap::<Vec<u8>, (u64, ValueType, u64)>::new();
1526+
let mut undecodable = 0_u64;
1527+
1528+
accumulate_live_user_bytes(
1529+
&mut latest,
1530+
&encode_internal_key(b"user:1", 1, ValueType::Put),
1531+
5,
1532+
&mut undecodable,
1533+
);
1534+
accumulate_live_user_bytes(
1535+
&mut latest,
1536+
&encode_internal_key(b"user:1", 2, ValueType::Delete),
1537+
0,
1538+
&mut undecodable,
1539+
);
1540+
accumulate_live_user_bytes(
1541+
&mut latest,
1542+
&encode_internal_key(b"user:2", 3, ValueType::Put),
1543+
4,
1544+
&mut undecodable,
1545+
);
1546+
1547+
let live = latest
1548+
.values()
1549+
.map(|(_, value_type, logical_bytes)| {
1550+
if *value_type == ValueType::Put {
1551+
*logical_bytes
1552+
} else {
1553+
0
1554+
}
1555+
})
1556+
.sum::<u64>()
1557+
.saturating_add(undecodable);
1558+
1559+
// user:1 latest is delete => 0 live bytes for that key.
1560+
// user:2 latest is put with logical bytes = len("user:2") + value_len(4) = 10.
1561+
assert_eq!(live, 10);
1562+
}
1563+
1564+
#[test]
1565+
fn live_user_bytes_counts_undecodable_keys_as_live() {
1566+
let mut latest = HashMap::<Vec<u8>, (u64, ValueType, u64)>::new();
1567+
let mut undecodable = 0_u64;
1568+
1569+
accumulate_live_user_bytes(&mut latest, b"\x01\x02\x03", 7, &mut undecodable);
1570+
assert_eq!(undecodable, 10);
1571+
assert!(latest.is_empty());
1572+
}
14651573
}

tests/integration/engine_compaction.rs

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -59,6 +59,11 @@ fn engine_runs_background_compaction_and_preserves_reads() {
5959
compaction_metrics.completed_compactions >= 1,
6060
"expected at least one completed compaction, got {compaction_metrics:?}"
6161
);
62+
assert!(compaction_metrics.total_bytes_on_disk > 0);
63+
assert!(
64+
compaction_metrics.live_bytes <= compaction_metrics.total_bytes_on_disk,
65+
"live bytes should not exceed total bytes: {compaction_metrics:?}"
66+
);
6267

6368
for sample in [0_u32, 37, 83, 119] {
6469
let key = format!("k{sample:04}");
@@ -126,6 +131,13 @@ fn compaction_collapses_old_versions_for_same_user_key() {
126131
hot_versions, 1,
127132
"expected one compacted version for hot-key, found {hot_versions}"
128133
);
134+
135+
let metrics = engine.compaction_metrics();
136+
assert!(metrics.total_bytes_on_disk > 0);
137+
assert!(
138+
metrics.live_bytes < metrics.total_bytes_on_disk,
139+
"expected compaction to reduce live/total ratio below 1.0, got {metrics:?}"
140+
);
129141
}
130142

131143
fs::remove_dir_all(dir).expect("cleanup temp dir");

0 commit comments

Comments
 (0)