typed_store/rocks/
options.rs1use std::{collections::BTreeMap, env};
6
7use iota_macros::nondeterministic;
8use rocksdb::{BlockBasedOptions, Cache, ReadOptions};
9use sysinfo::{MemoryRefreshKind, RefreshKind, System};
10use tap::TapFallible;
11use tracing::{debug, info, warn};
12
13const ENV_VAR_DB_WRITE_BUFFER_SIZE: &str = "DB_WRITE_BUFFER_SIZE_MB";
16const DEFAULT_DB_WRITE_BUFFER_SIZE: usize = 1024;
17
18const ENV_VAR_DB_WAL_SIZE: &str = "DB_WAL_SIZE_MB";
21const DEFAULT_DB_WAL_SIZE: usize = 1024;
22
23const ENV_VAR_L0_NUM_FILES_COMPACTION_TRIGGER: &str = "L0_NUM_FILES_COMPACTION_TRIGGER";
26const DEFAULT_L0_NUM_FILES_COMPACTION_TRIGGER: usize = 4;
27const DEFAULT_UNIVERSAL_COMPACTION_L0_NUM_FILES_COMPACTION_TRIGGER: usize = 80;
28const ENV_VAR_MAX_WRITE_BUFFER_SIZE_MB: &str = "MAX_WRITE_BUFFER_SIZE_MB";
29const DEFAULT_MAX_WRITE_BUFFER_SIZE_MB: usize = 256;
30const ENV_VAR_MAX_WRITE_BUFFER_NUMBER: &str = "MAX_WRITE_BUFFER_NUMBER";
31const DEFAULT_MAX_WRITE_BUFFER_NUMBER: usize = 6;
32const ENV_VAR_TARGET_FILE_SIZE_BASE_MB: &str = "TARGET_FILE_SIZE_BASE_MB";
33const DEFAULT_TARGET_FILE_SIZE_BASE_MB: usize = 128;
34
35const ENV_VAR_DISABLE_BLOB_STORAGE: &str = "DISABLE_BLOB_STORAGE";
37const ENV_VAR_DB_PARALLELISM: &str = "DB_PARALLELISM";
38
39#[derive(Clone, Debug, Default)]
40pub struct ReadWriteOptions {
41 pub log_value_hash: bool,
44}
45
46impl ReadWriteOptions {
47 pub fn readopts(&self) -> ReadOptions {
48 ReadOptions::default()
49 }
50
51 pub fn set_log_value_hash(mut self, log_value_hash: bool) -> Self {
52 self.log_value_hash = log_value_hash;
53 self
54 }
55}
56
57#[derive(Default, Clone)]
58pub struct DBOptions {
59 pub options: rocksdb::Options,
60 pub rw_options: ReadWriteOptions,
61}
62
63pub fn bulk_ingestion_write_options() -> rocksdb::WriteOptions {
66 let mut opts = rocksdb::WriteOptions::default();
67 opts.disable_wal(true);
68 opts
69}
70
71pub struct BulkIngestionOptions {
75 pub db_options: rocksdb::Options,
76 pub column_family_options: DBOptions,
77 pub batch_size_limit: usize,
78}
79
80pub fn bulk_ingestion_options() -> BulkIngestionOptions {
81 let total_memory_bytes = available_memory_bytes();
82 let num_cpus = num_cpus::get();
83
84 let mut db_options = rocksdb::Options::default();
85
86 db_options.set_unordered_write(true);
90
91 db_options.set_max_background_jobs(num_cpus as i32);
93
94 let db_write_buffer_size = (total_memory_bytes as f64 * 0.8) as usize;
98 db_options.set_db_write_buffer_size(db_write_buffer_size);
99
100 let mut column_family_options = default_db_options();
104 column_family_options
105 .options
106 .set_disable_auto_compactions(true);
107
108 column_family_options
114 .options
115 .set_level_zero_file_num_compaction_trigger(-1);
116 column_family_options
117 .options
118 .set_level_zero_slowdown_writes_trigger(-1);
119 column_family_options
120 .options
121 .set_level_zero_stop_writes_trigger(i32::MAX);
122
123 let cf_memory_budget = (total_memory_bytes as f64 * 0.25) as usize;
124 const MIN_BUFFER_SIZE: usize = 64 * 1024 * 1024; let target_buffer_count = num_cpus.max(2);
128 let buffer_size = (cf_memory_budget / target_buffer_count).max(MIN_BUFFER_SIZE);
129 let buffer_count = (cf_memory_budget / buffer_size).clamp(2, target_buffer_count) as i32;
130 column_family_options
131 .options
132 .set_write_buffer_size(buffer_size);
133 column_family_options
134 .options
135 .set_max_write_buffer_number(buffer_count);
136
137 let batch_size_limit = (buffer_size / 2).min(1 << 27);
140
141 debug!(
142 total_memory_bytes,
143 num_cpus,
144 db_write_buffer_size,
145 buffer_size,
146 buffer_count,
147 batch_size_limit,
148 "configured bulk ingestion options"
149 );
150
151 BulkIngestionOptions {
152 db_options,
153 column_family_options,
154 batch_size_limit,
155 }
156}
157
158fn available_memory_bytes() -> u64 {
161 nondeterministic!({
165 let mut sys = System::new_with_specifics(
168 RefreshKind::nothing().with_memory(MemoryRefreshKind::everything()),
169 );
170 sys.refresh_memory();
171
172 if let Some(cgroup_limits) = sys.cgroup_limits() {
173 if cgroup_limits.total_memory > 0 {
175 debug!(
176 limit = cgroup_limits.total_memory,
177 "using cgroup memory limit"
178 );
179 return cgroup_limits.total_memory;
180 }
181 }
182
183 let total = sys.total_memory();
184 debug!(total, "using total system memory");
185 total
186 })
187}
188
189#[derive(Clone)]
190pub struct DBMapTableConfigMap(BTreeMap<String, DBOptions>);
191impl DBMapTableConfigMap {
192 pub fn new(map: BTreeMap<String, DBOptions>) -> Self {
193 Self(map)
194 }
195
196 pub fn to_map(&self) -> BTreeMap<String, DBOptions> {
197 self.0.clone()
198 }
199}
200
201impl DBOptions {
202 pub fn optimize_for_point_lookup(mut self, block_cache_size_mb: usize) -> DBOptions {
206 self.options
208 .optimize_for_point_lookup(block_cache_size_mb as u64);
209 self
210 }
211
212 pub fn optimize_for_large_values_no_scan(mut self, min_blob_size: u64) -> DBOptions {
215 if env::var(ENV_VAR_DISABLE_BLOB_STORAGE).is_ok() {
216 info!("Large value blob storage optimization is disabled via env var.");
217 return self;
218 }
219
220 self.options.set_enable_blob_files(true);
222 self.options
223 .set_blob_compression_type(rocksdb::DBCompressionType::Lz4);
224 self.options.set_enable_blob_gc(true);
225 self.options.set_min_blob_size(min_blob_size);
229
230 let write_buffer_size = read_size_from_env(ENV_VAR_MAX_WRITE_BUFFER_SIZE_MB)
232 .unwrap_or(DEFAULT_MAX_WRITE_BUFFER_SIZE_MB)
233 * 1024
234 * 1024;
235 self.options.set_write_buffer_size(write_buffer_size);
236 let target_file_size_base = 64 << 20;
239 self.options
240 .set_target_file_size_base(target_file_size_base);
241 let max_level_zero_file_num = read_size_from_env(ENV_VAR_L0_NUM_FILES_COMPACTION_TRIGGER)
243 .unwrap_or(DEFAULT_L0_NUM_FILES_COMPACTION_TRIGGER);
244 self.options
245 .set_max_bytes_for_level_base(target_file_size_base * max_level_zero_file_num as u64);
246
247 self
248 }
249
250 pub fn optimize_for_read(mut self, block_cache_size_mb: usize) -> DBOptions {
252 self.options
253 .set_block_based_table_factory(&get_block_options(block_cache_size_mb, 16 << 10));
254 self
255 }
256
257 pub fn optimize_db_for_write_throughput(mut self, db_max_write_buffer_gb: u64) -> DBOptions {
259 self.options
260 .set_db_write_buffer_size(db_max_write_buffer_gb as usize * 1024 * 1024 * 1024);
261 self.options
262 .set_max_total_wal_size(db_max_write_buffer_gb * 1024 * 1024 * 1024);
263 self
264 }
265
266 pub fn optimize_for_write_throughput(mut self) -> DBOptions {
268 let write_buffer_size = read_size_from_env(ENV_VAR_MAX_WRITE_BUFFER_SIZE_MB)
270 .unwrap_or(DEFAULT_MAX_WRITE_BUFFER_SIZE_MB)
271 * 1024
272 * 1024;
273 self.options.set_write_buffer_size(write_buffer_size);
274 let max_write_buffer_number = read_size_from_env(ENV_VAR_MAX_WRITE_BUFFER_NUMBER)
276 .unwrap_or(DEFAULT_MAX_WRITE_BUFFER_NUMBER);
277 self.options
278 .set_max_write_buffer_number(max_write_buffer_number.try_into().unwrap());
279 self.options
281 .set_max_write_buffer_size_to_maintain((write_buffer_size).try_into().unwrap());
282
283 let max_level_zero_file_num = read_size_from_env(ENV_VAR_L0_NUM_FILES_COMPACTION_TRIGGER)
285 .unwrap_or(DEFAULT_L0_NUM_FILES_COMPACTION_TRIGGER);
286 self.options.set_level_zero_file_num_compaction_trigger(
287 max_level_zero_file_num.try_into().unwrap(),
288 );
289 self.options.set_level_zero_slowdown_writes_trigger(
290 (max_level_zero_file_num * 12).try_into().unwrap(),
291 );
292 self.options
293 .set_level_zero_stop_writes_trigger((max_level_zero_file_num * 16).try_into().unwrap());
294
295 self.options.set_target_file_size_base(
297 read_size_from_env(ENV_VAR_TARGET_FILE_SIZE_BASE_MB)
298 .unwrap_or(DEFAULT_TARGET_FILE_SIZE_BASE_MB) as u64
299 * 1024
300 * 1024,
301 );
302
303 self.options
305 .set_max_bytes_for_level_base((write_buffer_size * max_level_zero_file_num) as u64);
306
307 self
308 }
309
310 pub fn optimize_for_write_throughput_no_deletion(mut self) -> DBOptions {
314 let write_buffer_size = read_size_from_env(ENV_VAR_MAX_WRITE_BUFFER_SIZE_MB)
316 .unwrap_or(DEFAULT_MAX_WRITE_BUFFER_SIZE_MB)
317 * 1024
318 * 1024;
319 self.options.set_write_buffer_size(write_buffer_size);
320 let max_write_buffer_number = read_size_from_env(ENV_VAR_MAX_WRITE_BUFFER_NUMBER)
322 .unwrap_or(DEFAULT_MAX_WRITE_BUFFER_NUMBER);
323 self.options
324 .set_max_write_buffer_number(max_write_buffer_number.try_into().unwrap());
325 self.options
327 .set_max_write_buffer_size_to_maintain((write_buffer_size).try_into().unwrap());
328
329 self.options
331 .set_compaction_style(rocksdb::DBCompactionStyle::Universal);
332 let mut compaction_options = rocksdb::UniversalCompactOptions::default();
333 compaction_options.set_max_size_amplification_percent(10000);
334 compaction_options.set_stop_style(rocksdb::UniversalCompactionStopStyle::Similar);
335 self.options
336 .set_universal_compaction_options(&compaction_options);
337
338 let max_level_zero_file_num = read_size_from_env(ENV_VAR_L0_NUM_FILES_COMPACTION_TRIGGER)
339 .unwrap_or(DEFAULT_UNIVERSAL_COMPACTION_L0_NUM_FILES_COMPACTION_TRIGGER);
340 self.options.set_level_zero_file_num_compaction_trigger(
341 max_level_zero_file_num.try_into().unwrap(),
342 );
343 self.options.set_level_zero_slowdown_writes_trigger(
344 (max_level_zero_file_num * 12).try_into().unwrap(),
345 );
346 self.options
347 .set_level_zero_stop_writes_trigger((max_level_zero_file_num * 16).try_into().unwrap());
348
349 self.options.set_target_file_size_base(
351 read_size_from_env(ENV_VAR_TARGET_FILE_SIZE_BASE_MB)
352 .unwrap_or(DEFAULT_TARGET_FILE_SIZE_BASE_MB) as u64
353 * 1024
354 * 1024,
355 );
356
357 self.options
359 .set_max_bytes_for_level_base((write_buffer_size * max_level_zero_file_num) as u64);
360
361 self
362 }
363
364 pub fn set_block_options(
366 mut self,
367 block_cache_size_mb: usize,
368 block_size_bytes: usize,
369 ) -> DBOptions {
370 self.options
371 .set_block_based_table_factory(&get_block_options(
372 block_cache_size_mb,
373 block_size_bytes,
374 ));
375 self
376 }
377
378 pub fn disable_write_throttling(mut self) -> DBOptions {
380 self.options.set_soft_pending_compaction_bytes_limit(0);
381 self.options.set_hard_pending_compaction_bytes_limit(0);
382 self
383 }
384}
385
386pub fn default_db_options() -> DBOptions {
389 let mut opt = rocksdb::Options::default();
390
391 if let Some(limit) = fdlimit::raise_fd_limit() {
395 opt.set_max_open_files((limit / 8) as i32);
397 }
398
399 opt.set_table_cache_num_shard_bits(10);
402
403 opt.set_compression_type(rocksdb::DBCompressionType::Lz4);
405 opt.set_bottommost_compression_type(rocksdb::DBCompressionType::Zstd);
406 opt.set_bottommost_zstd_max_train_bytes(1024 * 1024, true);
407
408 opt.set_db_write_buffer_size(
420 read_size_from_env(ENV_VAR_DB_WRITE_BUFFER_SIZE).unwrap_or(DEFAULT_DB_WRITE_BUFFER_SIZE)
421 * 1024
422 * 1024,
423 );
424 opt.set_max_total_wal_size(
425 read_size_from_env(ENV_VAR_DB_WAL_SIZE).unwrap_or(DEFAULT_DB_WAL_SIZE) as u64 * 1024 * 1024,
426 );
427
428 opt.increase_parallelism(read_size_from_env(ENV_VAR_DB_PARALLELISM).unwrap_or(8) as i32);
430
431 opt.set_enable_pipelined_write(true);
432
433 opt.set_block_based_table_factory(&get_block_options(128, 16 << 10));
436
437 opt.set_memtable_prefix_bloom_ratio(0.02);
439
440 DBOptions {
441 options: opt,
442 rw_options: ReadWriteOptions::default(),
443 }
444}
445
446fn get_block_options(block_cache_size_mb: usize, block_size_bytes: usize) -> BlockBasedOptions {
447 let mut block_options = BlockBasedOptions::default();
452 block_options.set_block_size(block_size_bytes);
454 block_options.set_block_cache(&Cache::new_lru_cache(block_cache_size_mb << 20));
456 block_options.set_bloom_filter(10.0, false);
458 block_options.set_pin_l0_filter_and_index_blocks_in_cache(true);
460 block_options
461}
462
463pub fn list_tables(path: std::path::PathBuf) -> eyre::Result<Vec<String>> {
464 const DB_DEFAULT_CF_NAME: &str = "default";
465
466 let opts = rocksdb::Options::default();
467 rocksdb::DBWithThreadMode::<rocksdb::MultiThreaded>::list_cf(&opts, path)
468 .map_err(|e| e.into())
469 .map(|q| {
470 q.iter()
471 .filter_map(|s| {
472 if s != DB_DEFAULT_CF_NAME {
474 Some(s.clone())
475 } else {
476 None
477 }
478 })
479 .collect()
480 })
481}
482
483pub fn read_size_from_env(var_name: &str) -> Option<usize> {
484 env::var(var_name)
485 .ok()?
486 .parse::<usize>()
487 .tap_err(|e| {
488 warn!(
489 "Env var {} does not contain valid usize integer: {}",
490 var_name, e
491 )
492 })
493 .ok()
494}