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 {
write_buffer_manager: WriteBufferManager,
cache: Cache,
}
impl MemoryReporter for MemoryManager {
fn memory_allocated(&self) -> usize {
self.write_buffer_manager.get_usage() + self.cache.get_usage()
}
}
impl MemoryManager {
pub(super) fn configure(opts: &mut Options, config: &RocksDbConfig) -> Result<Self> {
let block_cache_size = config.block_cache_size;
let write_buffer_size = config.write_buffer_size;
let total_write_buffer_size =
config.max_write_buffer_number.saturating_mul(write_buffer_size);
let total_memory_limit = total_write_buffer_size + block_cache_size;
info!(target: TARGET, "Memory manager: total memory limit: {total_memory_limit}");
info!(target: TARGET, "Memory manager: block cache size: {block_cache_size}B");
let cache = Cache::new_lru_cache(config.block_cache_size);
let write_buffer_manager = WriteBufferManager::new_write_buffer_manager_with_cache(
total_memory_limit,
true,
cache.clone(),
);
opts.set_write_buffer_manager(&write_buffer_manager);
opts.set_row_cache(&cache);
let manager = Self {
write_buffer_manager,
cache,
};
manager.apply_to_cf_options(opts, config);
Ok(manager)
}
pub(super) fn apply_to_cf_options(&self, target: &mut Options, config: &RocksDbConfig) {
let write_buffer_size = config.write_buffer_size;
let requested_write_buffers_to_merge = config.min_write_buffer_number_to_merge.max(1);
let max_write_buffer_number = config.max_write_buffer_number.min(i32::MAX as usize) as i32;
let write_buffers_to_merge =
requested_write_buffers_to_merge.min(config.max_write_buffer_number.max(1));
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,
);
}
let min_write_buffers_to_merge = write_buffers_to_merge.min(i32::MAX as usize) as i32;
info!(target: TARGET, "Memory manager: write buffer size: {write_buffer_size}B");
target.set_write_buffer_size(write_buffer_size);
info!(target: TARGET, "Memory manager: maximum write buffers: {max_write_buffer_number}");
target.set_max_write_buffer_number(max_write_buffer_number);
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);
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);
info!(target: TARGET, "Target block size: {}", config.block_size);
block.set_block_size(config.block_size);
info!(target: TARGET, "Block cache size: {}", config.block_cache_size);
block.set_block_cache(&self.cache);
info!(target: TARGET, "Configuring two-level index search");
block.set_index_type(BlockBasedIndexType::TwoLevelIndexSearch);
info!(target: TARGET, "Use partitioned filters for each SST file");
block.set_partition_filters(true);
info!(target: TARGET, "Block size for partitioned metadata: 4096 B");
block.set_metadata_block_size(4096);
info!(target: TARGET, "Initial auto-readahead size: {}", config.initial_auto_readahead_size);
block.set_initial_auto_readahead_size(config.initial_auto_readahead_size);
info!(target: TARGET, "Maximum auto-readahead size: {}", config.max_auto_readahead_size);
block.set_max_auto_readahead_size(config.max_auto_readahead_size);
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);
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);
target.set_block_based_table_factory(&block);
target.set_blob_cache(&self.cache);
}
#[allow(clippy::clone_on_ref_ptr)] pub(super) fn register_with_allocator_tracker(self: &Arc<Self>) {
let reporter: Arc<dyn MemoryReporter> = self.clone();
register_memory_reporter("rocksdb", Arc::downgrade(&reporter));
}
pub fn shutdown(&self) -> Result<()> {
cleanup_memory_reporters();
Ok(())
}
}