1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
use std::sync::Arc;
use rocksdb::{BlockBasedIndexType, BlockBasedOptions, Cache, Options, WriteBufferManager};
use super::TARGET;
use crate::kvs::Result;
use crate::kvs::rocksdb::RocksDbConfig;
use crate::mem::{MemoryReporter, cleanup_memory_reporters, register_memory_reporter};
pub(super) struct MemoryManager {
/// The RocksDB block cache, shared with the row and blob caches.
///
/// Held so [`MemoryReporter::memory_allocated`] can read its usage, not to
/// keep it alive: `Options` clones both this cache and the write buffer
/// manager into its own `OptionsMustOutliveDB`, and an opened database
/// stores a clone of that, so both outlive the database regardless of what
/// this struct holds.
cache: Cache,
}
impl MemoryReporter for MemoryManager {
/// Reports the RocksDB memory that the global allocator tracker cannot see.
///
/// RocksDB allocates in C++, so none of its memory passes through Rust's
/// `GlobalAlloc` and none of it is counted by the tracking allocator. This
/// reporter is the only channel by which it enters the tracked total.
///
/// The [`WriteBufferManager`] built by [`Self::configure`] is constructed
/// against the same [`Cache`], so memtable arena bytes are charged into it
/// as reservation entries and `Cache::get_usage` covers memtables too.
/// Adding the manager's own usage on top would count those bytes twice.
///
/// Not covered: per-SST table-reader memory (top-level index and filter
/// blocks, table properties), compaction and flush buffers, iterator
/// readahead, and the pending writes held by an open transaction's
/// `WriteBatchWithIndex`. Callers sizing a memory threshold against a
/// container limit must leave headroom for those.
fn memory_allocated(&self) -> usize {
self.cache.get_usage()
}
}
impl MemoryManager {
/// Pre-configure the disk space manager
pub(super) fn configure(opts: &mut Options, config: &RocksDbConfig) -> Result<Self> {
// Get the configuration options
let block_cache_size = config.block_cache_size;
let write_buffer_size = config.write_buffer_size;
// Get the total write buffer size
let total_write_buffer_size =
config.max_write_buffer_number.saturating_mul(write_buffer_size);
// Combine the cache and the write buffers to get the memory limit
let total_memory_limit = total_write_buffer_size + block_cache_size;
info!(target: TARGET, "Memory manager: total memory limit: {total_memory_limit}");
// Set the block cache size in bytes
info!(target: TARGET, "Memory manager: block cache size: {block_cache_size}B");
// Configure the in-memory cache options
let cache = Cache::new_lru_cache(config.block_cache_size);
// Create a new write buffer manager with the cache
let write_buffer_manager = WriteBufferManager::new_write_buffer_manager_with_cache(
total_memory_limit,
true,
cache.clone(),
);
// The options take their own handle on the write buffer manager, and an
// opened database keeps a clone of the options, so it outlives both.
opts.set_write_buffer_manager(&write_buffer_manager);
// Set the row cache in the options
opts.set_row_cache(&cache);
// Build the manager and apply its per-CF settings to `opts`
let manager = Self {
cache,
};
manager.apply_to_cf_options(opts, config);
// Continue
Ok(manager)
}
/// Apply the column-family-level memory settings to `target`. Called
/// once on the main `opts` during [`Self::configure`] so the implicit
/// default CF is configured correctly, and also called on an
/// explicit [`rocksdb::ColumnFamilyDescriptor`]'s options when
/// versioning is enabled (see `Datastore::new`) so the explicit
/// default CF receives the same memtable and block-cache settings
/// rather than RocksDB's fresh defaults.
pub(super) fn apply_to_cf_options(&self, target: &mut Options, config: &RocksDbConfig) {
let write_buffer_size = config.write_buffer_size;
// Get the minimum number of write buffers to merge
let requested_write_buffers_to_merge = config.min_write_buffer_number_to_merge.max(1);
// Get the maximum number of write buffers
let max_write_buffer_number = config.max_write_buffer_number.min(i32::MAX as usize) as i32;
// Clamp the merge target to the maximum number of write buffers.
// RocksDB cannot merge more memtables than it is allowed to keep
// in memory, so if `min_write_buffer_number_to_merge` exceeds
// `max_write_buffer_number` (for example when an operator lowers
// the max to 1 on a constrained instance) writers would stall
// indefinitely waiting for a merge that can never happen. Clamp
// to a valid value and warn so the misconfiguration is visible.
let write_buffers_to_merge =
requested_write_buffers_to_merge.min(config.max_write_buffer_number.max(1));
// Check if the number of write buffers exceeds the maximum number allowed
if write_buffers_to_merge != requested_write_buffers_to_merge {
warn!(target: TARGET,
"Memory manager: min_write_buffer_number_to_merge ({requested_write_buffers_to_merge}) exceeds \
max_write_buffer_number ({}); clamping to {write_buffers_to_merge} to avoid \
stalling writers",
config.max_write_buffer_number,
);
}
// Get the adjusted minimum number of write buffers to merge
let min_write_buffers_to_merge = write_buffers_to_merge.min(i32::MAX as usize) as i32;
// Set the amount of data to build up in memory
info!(target: TARGET, "Memory manager: write buffer size: {write_buffer_size}B");
target.set_write_buffer_size(write_buffer_size);
// Set the maximum number of write buffers
info!(target: TARGET, "Memory manager: maximum write buffers: {max_write_buffer_number}");
target.set_max_write_buffer_number(max_write_buffer_number);
// Set minimum number of write buffers to merge
info!(target: TARGET, "Memory manager: minimum write buffers to merge: {min_write_buffers_to_merge}");
target.set_min_write_buffer_number_to_merge(min_write_buffers_to_merge);
// Configure the block based file options
let mut block = BlockBasedOptions::default();
block.set_pin_l0_filter_and_index_blocks_in_cache(true);
block.set_pin_top_level_index_and_filter(true);
block.set_bloom_filter(10.0, false);
// Configure the target block size
info!(target: TARGET, "Target block size: {}", config.block_size);
block.set_block_size(config.block_size);
// Configure the block cache
info!(target: TARGET, "Block cache size: {}", config.block_cache_size);
block.set_block_cache(&self.cache);
// Configure the index type and partition filters
info!(target: TARGET, "Configuring two-level index search");
block.set_index_type(BlockBasedIndexType::TwoLevelIndexSearch);
// Configure the partition filters for SST files
info!(target: TARGET, "Use partitioned filters for each SST file");
block.set_partition_filters(true);
// Configure the metadata block size
info!(target: TARGET, "Block size for partitioned metadata: 4096 B");
block.set_metadata_block_size(4096);
// Set the initial size for implicit iterator auto-readahead
info!(target: TARGET, "Initial auto-readahead size: {}", config.initial_auto_readahead_size);
block.set_initial_auto_readahead_size(config.initial_auto_readahead_size);
// Set the maximum size for implicit iterator auto-readahead
info!(target: TARGET, "Maximum auto-readahead size: {}", config.max_auto_readahead_size);
block.set_max_auto_readahead_size(config.max_auto_readahead_size);
// Set the number of sequential file reads before triggering auto-readahead
info!(target: TARGET, "Number of file reads for auto-readahead: {}", config.file_reads_for_auto_readahead);
block.set_num_file_reads_for_auto_readahead(config.file_reads_for_auto_readahead);
// When the prefix extractor is enabled the SST bloom filter is
// keyed on table+category prefixes. `whole_key_filtering=true`
// additionally adds whole keys (better for point lookups, larger
// filter); `whole_key_filtering=false` keeps the filter tight and
// focuses all bits on the prefix (better for scan-heavy workloads).
// The setting only matters when a prefix extractor is configured —
// without one, `whole_key_filtering` is effectively always true.
let whole_key_filtering = if config.prefix_extractor_enabled {
config.whole_key_filtering
} else {
true
};
info!(target: TARGET, "Memory manager: whole key filtering: {whole_key_filtering}");
block.set_whole_key_filtering(whole_key_filtering);
// Configure the database with the cache
target.set_block_based_table_factory(&block);
target.set_blob_cache(&self.cache);
}
// Register the memory manager with the global allocator tracker
#[allow(clippy::clone_on_ref_ptr)] // Arc::clone would not coerce to `Arc<dyn MemoryReporter>`
pub(super) fn register_with_allocator_tracker(self: &Arc<Self>) {
// Downgrade the memory manager to a memory reporter
let reporter: Arc<dyn MemoryReporter> = self.clone();
// Register with the global allocator tracker
register_memory_reporter("rocksdb", Arc::downgrade(&reporter));
}
/// Shutdown the memory manager
pub fn shutdown(&self) -> Result<()> {
// Clean up the memory manager
cleanup_memory_reporters();
// All good
Ok(())
}
}