Skip to main content

risingwave_storage/
opts.rs

1// Copyright 2023 RisingWave Labs
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use risingwave_common::config::storage::FileCacheRuntimeConfig;
16use risingwave_common::config::streaming::CacheRefillPolicy;
17use risingwave_common::config::{
18    EvictionConfig, ObjectStoreConfig, RwConfig, StorageMemoryConfig, extract_storage_memory_config,
19};
20use risingwave_common::system_param::reader::{SystemParamsRead, SystemParamsReader};
21use risingwave_common::system_param::system_params_for_test;
22
23#[derive(Clone, Debug)]
24pub struct StorageOpts {
25    /// The size of parallel task for one compact/flush job.
26    pub parallel_compact_size_mb: u32,
27    /// Target size of the Sstable.
28    pub sstable_size_mb: u32,
29    /// Minimal target size of the Sstable to store data of different state-table in independent files as soon as possible.
30    pub min_sstable_size_mb: u32,
31    /// Size of each block in bytes in SST.
32    pub block_size_kb: u32,
33    /// Deprecated and ignored by SST filter builders; kept for backward compatibility.
34    pub bloom_false_positive: f64,
35    /// parallelism while syncing share buffers into L0 SST. Should NOT be 0.
36    pub share_buffers_sync_parallelism: u32,
37    /// Worker threads number of dedicated tokio runtime for share buffer compaction. 0 means use
38    /// tokio's default value (number of CPU core).
39    pub share_buffer_compaction_worker_threads_number: u32,
40    /// Maximum shared buffer size, writes attempting to exceed the capacity will stall until there
41    /// is enough space.
42    pub shared_buffer_capacity_mb: usize,
43    /// The shared buffer will start flushing data to object when the ratio of memory usage to the
44    /// shared buffer capacity exceed such ratio.
45    pub shared_buffer_flush_ratio: f32,
46    /// The minimum total flush size of shared buffer spill. When a shared buffer spill is trigger,
47    /// the total flush size across multiple epochs should be at least higher than this size.
48    pub shared_buffer_min_batch_flush_size_mb: usize,
49    /// Remote directory for storing data and metadata objects.
50    pub data_directory: String,
51    /// Whether to enable write conflict detection
52    pub write_conflict_detection_enabled: bool,
53    /// Capacity of sstable block cache.
54    pub block_cache_capacity_mb: usize,
55    /// the number of block-cache shard. Less shard means that more concurrent-conflict.
56    pub block_cache_shard_num: usize,
57    /// Eviction config for block cache.
58    pub block_cache_eviction_config: EvictionConfig,
59    /// Capacity of sstable meta cache.
60    pub meta_cache_capacity_mb: usize,
61    /// the number of meta-cache shard. Less shard means that more concurrent-conflict.
62    pub meta_cache_shard_num: usize,
63    /// Eviction config for meta cache.
64    pub meta_cache_eviction_config: EvictionConfig,
65    /// max memory usage for large query.
66    pub prefetch_buffer_capacity_mb: usize,
67
68    pub max_cached_recent_versions_number: usize,
69
70    pub max_prefetch_block_number: usize,
71
72    pub disable_remote_compactor: bool,
73    /// Number of tasks shared buffer can upload in parallel.
74    pub share_buffer_upload_concurrency: usize,
75    /// Capacity of sstable meta cache.
76    pub compactor_memory_limit_mb: usize,
77    /// compactor streaming iterator recreate timeout.
78    /// deprecated
79    pub compact_iter_recreate_timeout_ms: u64,
80    /// Number of SST ids fetched from meta per RPC
81    pub sstable_id_remote_fetch_number: u32,
82    /// Whether to enable streaming upload for sstable.
83    pub min_sst_size_for_streaming_upload: u64,
84    pub max_concurrent_compaction_task_number: u64,
85    pub max_version_pinning_duration_sec: u64,
86    pub compactor_iter_max_io_retry_times: usize,
87
88    /// If set, block metadata keys will be shortened when their length exceeds this threshold.
89    pub shorten_block_meta_key_threshold: Option<usize>,
90
91    pub data_file_cache_dir: String,
92    pub data_file_cache_direct_io: bool,
93    pub data_file_cache_capacity_mb: usize,
94    pub data_file_cache_file_capacity_mb: usize,
95    pub data_file_cache_flushers: usize,
96    pub data_file_cache_reclaimers: usize,
97    pub data_file_cache_recover_mode: foyer::RecoverMode,
98    pub data_file_cache_recover_concurrency: usize,
99    pub data_file_cache_indexer_shards: usize,
100    pub data_file_cache_compression: foyer::Compression,
101    pub data_file_cache_flush_buffer_threshold_mb: usize,
102    pub data_file_cache_submit_queue_size_threshold_mb: usize,
103    pub data_file_cache_fifo_probation_ratio: f64,
104    pub data_file_cache_blob_index_size_kb: usize,
105    pub data_file_cache_runtime_config: FileCacheRuntimeConfig,
106    pub data_file_cache_throttle: foyer::Throttle,
107
108    pub cache_refill_data_refill_levels: Vec<u32>,
109    pub cache_refill_timeout_ms: u64,
110    pub cache_refill_meta_refill_concurrency: usize,
111    pub cache_refill_concurrency: usize,
112    pub cache_refill_recent_filter_shards: usize,
113    pub cache_refill_recent_filter_layers: usize,
114    pub cache_refill_recent_filter_rotate_interval_ms: usize,
115    pub cache_refill_unit: usize,
116    pub cache_refill_threshold: f64,
117    pub cache_refill_skip_recent_filter: bool,
118    pub cache_refill_skip_inheritance_filter: bool,
119    pub cache_refill_table_cache_refill_default_policy: CacheRefillPolicy,
120
121    pub meta_file_cache_dir: String,
122    pub meta_file_cache_direct_io: bool,
123    pub meta_file_cache_capacity_mb: usize,
124    pub meta_file_cache_file_capacity_mb: usize,
125    pub meta_file_cache_flushers: usize,
126    pub meta_file_cache_reclaimers: usize,
127    pub meta_file_cache_recover_mode: foyer::RecoverMode,
128    pub meta_file_cache_recover_concurrency: usize,
129    pub meta_file_cache_indexer_shards: usize,
130    pub meta_file_cache_compression: foyer::Compression,
131    pub meta_file_cache_flush_buffer_threshold_mb: usize,
132    pub meta_file_cache_submit_queue_size_threshold_mb: usize,
133    pub meta_file_cache_fifo_probation_ratio: f64,
134    pub meta_file_cache_blob_index_size_kb: usize,
135    pub meta_file_cache_runtime_config: FileCacheRuntimeConfig,
136    pub meta_file_cache_throttle: foyer::Throttle,
137    pub sst_skip_bloom_filter_in_serde: bool,
138
139    pub vector_file_block_size_kb: usize,
140    pub vector_block_cache_capacity_mb: usize,
141    pub vector_block_cache_shard_num: usize,
142    pub vector_block_cache_eviction_config: EvictionConfig,
143    pub vector_meta_cache_capacity_mb: usize,
144    pub vector_meta_cache_shard_num: usize,
145    pub vector_meta_cache_eviction_config: EvictionConfig,
146
147    /// The storage url for storing backups.
148    pub backup_storage_url: String,
149    /// The storage directory for storing backups.
150    pub backup_storage_directory: String,
151    /// max time which wait for preload. 0 represent do not do any preload.
152    pub max_preload_wait_time_mill: u64,
153
154    pub compactor_max_sst_key_count: u64,
155    pub compactor_max_task_multiplier: f32,
156    pub compactor_max_sst_size: u64,
157    /// enable `FastCompactorRunner`.
158    pub enable_fast_compaction: bool,
159    pub check_compaction_result: bool,
160    pub compactor_fast_max_compact_delete_ratio: u32,
161    pub compactor_fast_max_compact_task_size: u64,
162
163    pub mem_table_spill_threshold: usize,
164
165    pub compactor_concurrent_uploading_sst_count: Option<usize>,
166
167    pub compactor_max_overlap_sst_count: usize,
168
169    /// The maximum number of meta files that can be preloaded.
170    pub compactor_max_preload_meta_file_count: usize,
171
172    pub object_store_config: ObjectStoreConfig,
173    pub time_travel_version_cache_capacity: u64,
174    pub table_change_log_cache_capacity: u64,
175
176    pub iceberg_compaction_enable_validate: bool,
177    pub iceberg_compaction_max_record_batch_rows: usize,
178    pub iceberg_compaction_write_parquet_max_row_group_rows: usize,
179    pub iceberg_compaction_min_size_per_partition_mb: u32,
180    pub iceberg_compaction_max_file_count_per_partition: u32,
181    pub iceberg_compaction_target_binpack_group_size_mb: Option<u64>,
182    pub iceberg_compaction_min_group_size_mb: Option<u64>,
183    pub iceberg_compaction_min_group_file_count: Option<usize>,
184
185    /// The ratio of iceberg compaction max parallelism to the number of CPU cores
186    pub iceberg_compaction_task_parallelism_ratio: f32,
187    /// Whether to enable heuristic output parallelism in iceberg compaction.
188    pub iceberg_compaction_enable_heuristic_output_parallelism: bool,
189    /// Maximum number of concurrent file close operations
190    pub iceberg_compaction_max_concurrent_closes: usize,
191    /// Whether to enable dynamic size estimation for iceberg compaction.
192    pub iceberg_compaction_enable_dynamic_size_estimation: bool,
193    /// The smoothing factor for size estimation in iceberg compaction.(default: 0.3)
194    pub iceberg_compaction_size_estimation_smoothing_factor: f64,
195    /// Multiplier for pending waiting parallelism budget for iceberg compaction task queue.
196    pub iceberg_compaction_pending_parallelism_budget_multiplier: f32,
197    /// Maximum number of Iceberg compaction tasks requested in one pull.
198    pub iceberg_compaction_max_pull_task_count: u32,
199    /// Pull interval for iceberg compaction task requests in milliseconds.
200    pub iceberg_compaction_pull_interval_ms: u64,
201    /// Whether to enable prefetch for iceberg compaction.
202    pub iceberg_compaction_enable_prefetch: bool,
203}
204
205impl Default for StorageOpts {
206    fn default() -> Self {
207        let c = RwConfig::default();
208        let p = system_params_for_test();
209        let s = extract_storage_memory_config(&c);
210        Self::from((&c, &p.into(), &s))
211    }
212}
213
214impl From<(&RwConfig, &SystemParamsReader, &StorageMemoryConfig)> for StorageOpts {
215    fn from((c, p, s): (&RwConfig, &SystemParamsReader, &StorageMemoryConfig)) -> Self {
216        let mut data_file_cache_throttle = c.storage.data_file_cache.throttle.clone();
217        if data_file_cache_throttle.write_throughput.is_none() {
218            data_file_cache_throttle = data_file_cache_throttle.with_write_throughput(
219                c.storage.data_file_cache.insert_rate_limit_mb * 1024 * 1024,
220            );
221        }
222        let mut meta_file_cache_throttle = c.storage.meta_file_cache.throttle.clone();
223        if meta_file_cache_throttle.write_throughput.is_none() {
224            meta_file_cache_throttle = meta_file_cache_throttle.with_write_throughput(
225                c.storage.meta_file_cache.insert_rate_limit_mb * 1024 * 1024,
226            );
227        }
228
229        Self {
230            parallel_compact_size_mb: p.parallel_compact_size_mb(),
231            sstable_size_mb: p.sstable_size_mb(),
232            min_sstable_size_mb: c.storage.min_sstable_size_mb,
233            block_size_kb: p.block_size_kb(),
234            bloom_false_positive: p.bloom_false_positive(),
235            share_buffers_sync_parallelism: c.storage.share_buffers_sync_parallelism,
236            share_buffer_compaction_worker_threads_number: c
237                .storage
238                .share_buffer_compaction_worker_threads_number,
239            shared_buffer_capacity_mb: s.shared_buffer_capacity_mb,
240            shared_buffer_flush_ratio: c.storage.shared_buffer_flush_ratio,
241            shared_buffer_min_batch_flush_size_mb: c.storage.shared_buffer_min_batch_flush_size_mb,
242            data_directory: p.data_directory().to_owned(),
243            write_conflict_detection_enabled: c.storage.write_conflict_detection_enabled,
244            block_cache_capacity_mb: s.block_cache_capacity_mb,
245            block_cache_shard_num: s.block_cache_shard_num,
246            block_cache_eviction_config: s.block_cache_eviction_config.clone(),
247            meta_cache_capacity_mb: s.meta_cache_capacity_mb,
248            meta_cache_shard_num: s.meta_cache_shard_num,
249            meta_cache_eviction_config: s.meta_cache_eviction_config.clone(),
250            prefetch_buffer_capacity_mb: s.prefetch_buffer_capacity_mb,
251            max_cached_recent_versions_number: c.storage.max_cached_recent_versions_number,
252            max_prefetch_block_number: c.storage.max_prefetch_block_number,
253            disable_remote_compactor: c.storage.disable_remote_compactor,
254            share_buffer_upload_concurrency: c.storage.share_buffer_upload_concurrency,
255            compactor_memory_limit_mb: s.compactor_memory_limit_mb,
256            sstable_id_remote_fetch_number: c.storage.sstable_id_remote_fetch_number,
257            min_sst_size_for_streaming_upload: c.storage.min_sst_size_for_streaming_upload,
258            max_concurrent_compaction_task_number: c.storage.max_concurrent_compaction_task_number,
259            max_version_pinning_duration_sec: c.storage.max_version_pinning_duration_sec,
260            data_file_cache_dir: c.storage.data_file_cache.dir.clone(),
261            data_file_cache_direct_io: c.storage.data_file_cache.direct_io,
262            data_file_cache_capacity_mb: c.storage.data_file_cache.capacity_mb,
263            data_file_cache_file_capacity_mb: c.storage.data_file_cache.file_capacity_mb,
264            data_file_cache_flushers: c.storage.data_file_cache.flushers,
265            data_file_cache_reclaimers: c.storage.data_file_cache.reclaimers,
266            data_file_cache_recover_mode: c.storage.data_file_cache.recover_mode,
267            data_file_cache_recover_concurrency: c.storage.data_file_cache.recover_concurrency,
268            data_file_cache_indexer_shards: c.storage.data_file_cache.indexer_shards,
269            data_file_cache_compression: c.storage.data_file_cache.compression,
270            data_file_cache_flush_buffer_threshold_mb: s.block_file_cache_flush_buffer_threshold_mb,
271            data_file_cache_submit_queue_size_threshold_mb: c
272                .storage
273                .data_file_cache
274                .submit_queue_size_threshold_mb,
275            data_file_cache_fifo_probation_ratio: c.storage.data_file_cache.fifo_probation_ratio,
276            data_file_cache_blob_index_size_kb: c.storage.data_file_cache.blob_index_size_kb,
277            data_file_cache_runtime_config: c.storage.data_file_cache.runtime_config.clone(),
278            data_file_cache_throttle,
279            meta_file_cache_dir: c.storage.meta_file_cache.dir.clone(),
280            meta_file_cache_direct_io: c.storage.meta_file_cache.direct_io,
281            meta_file_cache_capacity_mb: c.storage.meta_file_cache.capacity_mb,
282            meta_file_cache_file_capacity_mb: c.storage.meta_file_cache.file_capacity_mb,
283            meta_file_cache_flushers: c.storage.meta_file_cache.flushers,
284            meta_file_cache_reclaimers: c.storage.meta_file_cache.reclaimers,
285            meta_file_cache_recover_mode: c.storage.meta_file_cache.recover_mode,
286            meta_file_cache_recover_concurrency: c.storage.meta_file_cache.recover_concurrency,
287            meta_file_cache_indexer_shards: c.storage.meta_file_cache.indexer_shards,
288            meta_file_cache_compression: c.storage.meta_file_cache.compression,
289            meta_file_cache_flush_buffer_threshold_mb: s.meta_file_cache_flush_buffer_threshold_mb,
290            meta_file_cache_submit_queue_size_threshold_mb: c
291                .storage
292                .meta_file_cache
293                .submit_queue_size_threshold_mb,
294            meta_file_cache_fifo_probation_ratio: c.storage.meta_file_cache.fifo_probation_ratio,
295            meta_file_cache_blob_index_size_kb: c.storage.meta_file_cache.blob_index_size_kb,
296            meta_file_cache_runtime_config: c.storage.meta_file_cache.runtime_config.clone(),
297            meta_file_cache_throttle,
298            sst_skip_bloom_filter_in_serde: c.storage.sst_skip_bloom_filter_in_serde,
299            cache_refill_data_refill_levels: c.storage.cache_refill.data_refill_levels.clone(),
300            cache_refill_timeout_ms: c.storage.cache_refill.timeout_ms,
301            cache_refill_meta_refill_concurrency: c.storage.cache_refill.meta_refill_concurrency,
302            cache_refill_concurrency: c.storage.cache_refill.concurrency,
303            cache_refill_recent_filter_shards: c.storage.cache_refill.recent_filter_shards,
304            cache_refill_recent_filter_layers: c.storage.cache_refill.recent_filter_layers,
305            cache_refill_recent_filter_rotate_interval_ms: c
306                .storage
307                .cache_refill
308                .recent_filter_rotate_interval_ms,
309            cache_refill_unit: c.storage.cache_refill.unit,
310            cache_refill_threshold: c.storage.cache_refill.threshold,
311            cache_refill_skip_recent_filter: c.storage.cache_refill.skip_recent_filter,
312            cache_refill_skip_inheritance_filter: c.storage.cache_refill.skip_inheritance_filter,
313            cache_refill_table_cache_refill_default_policy: c
314                .streaming
315                .developer
316                .cache_refill_policy,
317            max_preload_wait_time_mill: c.storage.max_preload_wait_time_mill,
318            compact_iter_recreate_timeout_ms: c.storage.compact_iter_recreate_timeout_ms,
319
320            backup_storage_url: p.backup_storage_url().to_owned(),
321            backup_storage_directory: p.backup_storage_directory().to_owned(),
322            compactor_max_sst_key_count: c.storage.compactor_max_sst_key_count,
323            compactor_max_task_multiplier: c.storage.compactor_max_task_multiplier,
324            compactor_max_sst_size: c.storage.compactor_max_sst_size,
325            enable_fast_compaction: c.storage.enable_fast_compaction,
326            check_compaction_result: c.storage.check_compaction_result,
327            mem_table_spill_threshold: c.storage.mem_table_spill_threshold,
328            object_store_config: c.storage.object_store.clone(),
329            compactor_fast_max_compact_delete_ratio: c
330                .storage
331                .compactor_fast_max_compact_delete_ratio,
332            compactor_fast_max_compact_task_size: c.storage.compactor_fast_max_compact_task_size,
333            compactor_iter_max_io_retry_times: c.storage.compactor_iter_max_io_retry_times,
334            shorten_block_meta_key_threshold: c.storage.shorten_block_meta_key_threshold,
335            compactor_concurrent_uploading_sst_count: c
336                .storage
337                .compactor_concurrent_uploading_sst_count,
338            time_travel_version_cache_capacity: c.storage.time_travel_version_cache_capacity,
339            table_change_log_cache_capacity: c.storage.table_change_log_cache_capacity,
340            compactor_max_overlap_sst_count: c.storage.compactor_max_overlap_sst_count,
341            compactor_max_preload_meta_file_count: c.storage.compactor_max_preload_meta_file_count,
342
343            iceberg_compaction_enable_validate: c.storage.iceberg_compaction_enable_validate,
344            iceberg_compaction_max_record_batch_rows: c
345                .storage
346                .iceberg_compaction_max_record_batch_rows,
347            #[expect(deprecated)]
348            iceberg_compaction_write_parquet_max_row_group_rows: c
349                .storage
350                .iceberg_compaction_write_parquet_max_row_group_rows,
351            iceberg_compaction_min_size_per_partition_mb: c
352                .storage
353                .iceberg_compaction_min_size_per_partition_mb,
354            iceberg_compaction_max_file_count_per_partition: c
355                .storage
356                .iceberg_compaction_max_file_count_per_partition,
357            iceberg_compaction_task_parallelism_ratio: c
358                .storage
359                .iceberg_compaction_task_parallelism_ratio,
360            iceberg_compaction_enable_heuristic_output_parallelism: c
361                .storage
362                .iceberg_compaction_enable_heuristic_output_parallelism,
363            iceberg_compaction_max_concurrent_closes: c
364                .storage
365                .iceberg_compaction_max_concurrent_closes,
366            iceberg_compaction_enable_dynamic_size_estimation: c
367                .storage
368                .iceberg_compaction_enable_dynamic_size_estimation,
369            iceberg_compaction_size_estimation_smoothing_factor: c
370                .storage
371                .iceberg_compaction_size_estimation_smoothing_factor,
372            iceberg_compaction_pending_parallelism_budget_multiplier: c
373                .storage
374                .iceberg_compaction_pending_parallelism_budget_multiplier,
375            iceberg_compaction_max_pull_task_count: c
376                .storage
377                .iceberg_compaction_max_pull_task_count,
378            iceberg_compaction_pull_interval_ms: c.storage.iceberg_compaction_pull_interval_ms,
379            iceberg_compaction_enable_prefetch: c.storage.iceberg_compaction_enable_prefetch,
380            iceberg_compaction_target_binpack_group_size_mb: c
381                .storage
382                .iceberg_compaction_target_binpack_group_size_mb,
383            iceberg_compaction_min_group_size_mb: c.storage.iceberg_compaction_min_group_size_mb,
384            iceberg_compaction_min_group_file_count: c
385                .storage
386                .iceberg_compaction_min_group_file_count,
387            vector_file_block_size_kb: c.storage.vector_file_block_size_kb,
388            vector_block_cache_capacity_mb: s.vector_block_cache_capacity_mb,
389            vector_block_cache_shard_num: s.vector_block_cache_shard_num,
390            vector_block_cache_eviction_config: s.vector_block_cache_eviction_config.clone(),
391            vector_meta_cache_capacity_mb: s.vector_meta_cache_capacity_mb,
392            vector_meta_cache_shard_num: s.vector_meta_cache_shard_num,
393            vector_meta_cache_eviction_config: s.vector_meta_cache_eviction_config.clone(),
394        }
395    }
396}
397
398#[cfg(test)]
399mod tests {
400    use super::*;
401
402    #[test]
403    fn test_file_cache_submit_queue_size_threshold() {
404        let mut config = RwConfig::default();
405        config
406            .storage
407            .data_file_cache
408            .submit_queue_size_threshold_mb = 256;
409        config
410            .storage
411            .meta_file_cache
412            .submit_queue_size_threshold_mb = 32;
413        let system_params = system_params_for_test();
414        let storage_memory_config = extract_storage_memory_config(&config);
415
416        let opts = StorageOpts::from((&config, &system_params.into(), &storage_memory_config));
417
418        assert_eq!(opts.data_file_cache_submit_queue_size_threshold_mb, 256);
419        assert_eq!(opts.meta_file_cache_submit_queue_size_threshold_mb, 32);
420    }
421}