Skip to main content

risingwave_compactor/
server.rs

1// Copyright 2022 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 std::net::SocketAddr;
16use std::sync::Arc;
17use std::time::Duration;
18
19use risingwave_common::config::{
20    AsyncStackTraceOption, MetricLevel, RwConfig, extract_storage_memory_config, load_config,
21};
22use risingwave_common::monitor::{GLOBAL_METRICS_REGISTRY, RouterExt, TcpConfig};
23use risingwave_common::system_param::local_manager::LocalSystemParamsManager;
24use risingwave_common::system_param::reader::{SystemParamsRead, SystemParamsReader};
25use risingwave_common::telemetry::manager::TelemetryManager;
26use risingwave_common::telemetry::telemetry_env_enabled;
27use risingwave_common::util::addr::HostAddr;
28use risingwave_common::util::resource_util::memory::system_memory_available_bytes;
29use risingwave_common::util::tokio_util::sync::CancellationToken;
30use risingwave_common::{GIT_SHA, RW_VERSION};
31use risingwave_common_heap_profiling::HeapProfiler;
32use risingwave_common_service::{MetricsManager, ObserverManager};
33use risingwave_object_store::object::build_remote_object_store;
34use risingwave_object_store::object::object_metrics::GLOBAL_OBJECT_STORE_METRICS;
35use risingwave_pb::common::WorkerType;
36use risingwave_pb::common::worker_node::Property;
37use risingwave_pb::compactor::compactor_service_server::CompactorServiceServer;
38use risingwave_pb::configured_monitor_service_server;
39use risingwave_pb::monitor_service::monitor_service_server::MonitorServiceServer;
40use risingwave_rpc_client::{GrpcCompactorProxyClient, MetaClient};
41use risingwave_storage::compaction_catalog_manager::{
42    CompactionCatalogManager, RemoteTableAccessor,
43};
44use risingwave_storage::hummock::compactor::{
45    CompactionAwaitTreeRegRef, CompactionExecutor, CompactorContext,
46    new_compaction_await_tree_reg_ref,
47};
48use risingwave_storage::hummock::hummock_meta_client::MonitoredHummockMetaClient;
49use risingwave_storage::hummock::utils::HummockMemoryCollector;
50use risingwave_storage::hummock::{MemoryLimiter, ObjectIdManager, SstableStore};
51use risingwave_storage::monitor::{
52    CompactorMetrics, GLOBAL_COMPACTOR_METRICS, GLOBAL_HUMMOCK_METRICS, monitor_cache,
53};
54use risingwave_storage::opts::StorageOpts;
55use tokio::sync::mpsc;
56use tracing::info;
57
58use super::compactor_observer::observer_manager::CompactorObserverNode;
59use crate::rpc::{CompactorServiceImpl, MonitorServiceImpl};
60use crate::telemetry::CompactorTelemetryCreator;
61use crate::{
62    CompactorMode, CompactorOpts, default_rpc_max_decoding_message_size_bytes,
63    default_rpc_max_encoding_message_size_bytes,
64};
65
66pub async fn prepare_start_parameters(
67    compactor_opts: &CompactorOpts,
68    config: RwConfig,
69    system_params_reader: SystemParamsReader,
70) -> (
71    Arc<SstableStore>,
72    Arc<MemoryLimiter>,
73    HeapProfiler,
74    Option<CompactionAwaitTreeRegRef>,
75    Arc<StorageOpts>,
76    Arc<CompactorMetrics>,
77) {
78    // Boot compactor
79    let object_metrics = Arc::new(GLOBAL_OBJECT_STORE_METRICS.clone());
80    let compactor_metrics = Arc::new(GLOBAL_COMPACTOR_METRICS.clone());
81
82    let state_store_url = system_params_reader.state_store();
83    let state_store_url = state_store_url.expose();
84
85    let storage_memory_config = extract_storage_memory_config(&config);
86    let storage_opts: Arc<StorageOpts> = Arc::new(StorageOpts::from((
87        &config,
88        &system_params_reader,
89        &storage_memory_config,
90    )));
91    let non_reserved_memory_bytes = (compactor_opts.compactor_total_memory_bytes as f64
92        * config.storage.compactor_memory_available_proportion)
93        as usize;
94    let meta_cache_capacity_bytes = compactor_opts.compactor_meta_cache_memory_bytes;
95    let mut compactor_memory_limit_bytes = match config.storage.compactor_memory_limit_mb {
96        Some(compactor_memory_limit_mb) => compactor_memory_limit_mb * (1 << 20),
97        None => non_reserved_memory_bytes,
98    };
99
100    compactor_memory_limit_bytes = compactor_memory_limit_bytes.checked_sub(compactor_opts.compactor_meta_cache_memory_bytes).unwrap_or_else(|| {
101        panic!(
102            "compactor_memory_limit_bytes{} is too small to hold compactor_meta_cache_memory_bytes {}",
103            compactor_memory_limit_bytes,
104            meta_cache_capacity_bytes
105        );
106    });
107
108    tracing::info!(
109        "Compactor non_reserved_memory_bytes {} meta_cache_capacity_bytes {} compactor_memory_limit_bytes {} sstable_size_bytes {} block_size_bytes {}",
110        non_reserved_memory_bytes,
111        meta_cache_capacity_bytes,
112        compactor_memory_limit_bytes,
113        storage_opts.sstable_size_mb * (1 << 20),
114        storage_opts.block_size_kb * (1 << 10),
115    );
116
117    // check memory config
118    {
119        // This is a similar logic to SstableBuilder memory detection, to ensure that we can find
120        // configuration problems as quickly as possible
121        let min_compactor_memory_limit_bytes = (storage_opts.sstable_size_mb * (1 << 20)
122            + storage_opts.block_size_kb * (1 << 10))
123            as u64;
124
125        assert!(compactor_memory_limit_bytes > min_compactor_memory_limit_bytes as usize * 2);
126    }
127
128    let object_store = build_remote_object_store(
129        state_store_url
130            .strip_prefix("hummock+")
131            .expect("object store must be hummock for compactor server"),
132        object_metrics,
133        "Hummock",
134        Arc::new(config.storage.object_store.clone()),
135    )
136    .await;
137
138    let object_store = Arc::new(object_store);
139    let sstable_store = Arc::new(
140        SstableStore::for_compactor(
141            object_store,
142            storage_opts.data_directory.clone(),
143            0,
144            meta_cache_capacity_bytes,
145            system_params_reader.use_new_object_prefix_strategy(),
146        )
147        .await
148        // FIXME(MrCroxx): Handle this error.
149        .unwrap(),
150    );
151
152    let memory_limiter = Arc::new(MemoryLimiter::new(compactor_memory_limit_bytes as u64));
153    let storage_memory_config = extract_storage_memory_config(&config);
154    let memory_collector = Arc::new(HummockMemoryCollector::new(
155        sstable_store.clone(),
156        memory_limiter.clone(),
157        storage_memory_config,
158    ));
159
160    let heap_profiler = HeapProfiler::new(
161        system_memory_available_bytes(),
162        config.server.heap_profiling.clone(),
163    );
164
165    monitor_cache(memory_collector);
166
167    let await_tree_config = match &config.streaming.async_stack_trace {
168        AsyncStackTraceOption::Off => None,
169        c => await_tree::ConfigBuilder::default()
170            .verbose(c.is_verbose().unwrap())
171            .build()
172            .ok(),
173    };
174    let await_tree_reg = await_tree_config.map(new_compaction_await_tree_reg_ref);
175
176    (
177        sstable_store,
178        memory_limiter,
179        heap_profiler,
180        await_tree_reg,
181        storage_opts,
182        compactor_metrics,
183    )
184}
185
186/// Resolves the memory budget dedicated to Iceberg compaction.
187fn resolve_iceberg_compaction_memory_budget(
188    configured_limit_mb: Option<usize>,
189    compactor_total_memory_bytes: usize,
190    available_proportion: f64,
191) -> usize {
192    const MB: usize = 1 << 20;
193    let budget = match configured_limit_mb {
194        Some(limit_mb) => limit_mb
195            .checked_mul(MB)
196            .expect("Iceberg compaction memory limit overflows usize"),
197        None => (compactor_total_memory_bytes as f64 * available_proportion) as usize,
198    };
199    assert!(
200        budget > 0,
201        "Iceberg compaction memory limit must be positive"
202    );
203    budget
204}
205
206/// Fetches and runs compaction tasks.
207///
208/// Returns when the `shutdown` token is triggered.
209pub async fn compactor_serve(
210    listen_addr: SocketAddr,
211    advertise_addr: HostAddr,
212    opts: CompactorOpts,
213    shutdown: CancellationToken,
214    compactor_mode: CompactorMode,
215) {
216    let config = load_config(&opts.config_path, &opts);
217    info!("Starting compactor node",);
218    info!("> config: {:?}", config);
219    info!(
220        "> debug assertions: {}",
221        if cfg!(debug_assertions) { "on" } else { "off" }
222    );
223    info!("> version: {} ({})", RW_VERSION, GIT_SHA);
224
225    let is_iceberg_compactor = matches!(
226        compactor_mode,
227        CompactorMode::DedicatedIceberg | CompactorMode::SharedIceberg
228    );
229    if is_iceberg_compactor && config.server.metrics_level > MetricLevel::Disabled {
230        iceberg_storage_opendal::install_prometheus_metrics(&GLOBAL_METRICS_REGISTRY)
231            .expect("failed to install Iceberg OpenDAL metrics");
232    }
233
234    let compaction_executor = Arc::new(CompactionExecutor::new(
235        opts.compaction_worker_threads_number,
236    ));
237
238    let max_task_parallelism: u32 = (compaction_executor.worker_num() as f32
239        * config.storage.compactor_max_task_multiplier)
240        .ceil() as u32;
241
242    // Register to the cluster.
243    let (meta_client, system_params_reader) = MetaClient::register_new(
244        opts.meta_address.clone(),
245        WorkerType::Compactor,
246        &advertise_addr,
247        Property {
248            is_iceberg_compactor,
249            parallelism: max_task_parallelism,
250            ..Default::default()
251        },
252        Arc::new(config.meta.clone()),
253    )
254    .await;
255
256    info!("Assigned compactor id {}", meta_client.worker_id());
257
258    let hummock_metrics = Arc::new(GLOBAL_HUMMOCK_METRICS.clone());
259
260    let hummock_meta_client = Arc::new(MonitoredHummockMetaClient::new(
261        meta_client.clone(),
262        hummock_metrics.clone(),
263    ));
264
265    let (
266        sstable_store,
267        memory_limiter,
268        heap_profiler,
269        await_tree_reg,
270        storage_opts,
271        compactor_metrics,
272    ) = Box::pin(prepare_start_parameters(
273        &opts,
274        config.clone(),
275        system_params_reader.clone(),
276    ))
277    .await;
278    let iceberg_memory_budget_bytes = matches!(compactor_mode, CompactorMode::DedicatedIceberg)
279        .then(|| {
280            resolve_iceberg_compaction_memory_budget(
281                config.storage.iceberg_compaction_memory_limit_mb,
282                opts.compactor_total_memory_bytes,
283                config.storage.compactor_memory_available_proportion,
284            )
285        });
286
287    let compaction_catalog_manager_ref = Arc::new(CompactionCatalogManager::new(Box::new(
288        RemoteTableAccessor::new(meta_client.clone()),
289    )));
290
291    let system_params_manager = Arc::new(LocalSystemParamsManager::new(system_params_reader));
292    let compactor_observer_node = CompactorObserverNode::new(
293        compaction_catalog_manager_ref.clone(),
294        system_params_manager.clone(),
295    );
296    let observer_manager =
297        ObserverManager::new_with_meta_client(meta_client.clone(), compactor_observer_node).await;
298
299    // Run a background heap profiler
300    heap_profiler.start();
301
302    // use half of limit because any memory which would hold in meta-cache will be allocate by
303    // limited at first.
304    let _observer_join_handle = observer_manager.start().await;
305
306    let object_id_manager = Arc::new(ObjectIdManager::new(
307        hummock_meta_client.clone(),
308        storage_opts.sstable_id_remote_fetch_number,
309    ));
310
311    let compactor_context = CompactorContext {
312        storage_opts,
313        sstable_store: sstable_store.clone(),
314        compactor_metrics,
315        is_share_buffer_compact: false,
316        compaction_executor,
317        memory_limiter,
318        task_progress_manager: Default::default(),
319        await_tree_reg: await_tree_reg.clone(),
320    };
321
322    // TODO(shutdown): don't collect sub-tasks as there's no need to gracefully shutdown them.
323    let mut sub_tasks = vec![
324        MetaClient::start_heartbeat_loop(
325            meta_client.clone(),
326            Duration::from_millis(config.server.heartbeat_interval_ms as u64),
327        ),
328        match compactor_mode {
329            CompactorMode::Dedicated => risingwave_storage::hummock::compactor::start_compactor(
330                compactor_context.clone(),
331                hummock_meta_client.clone(),
332                object_id_manager.clone(),
333                compaction_catalog_manager_ref,
334            ),
335            CompactorMode::Shared => unreachable!(),
336            CompactorMode::DedicatedIceberg => {
337                risingwave_storage::hummock::compactor::start_iceberg_compactor(
338                    compactor_context.clone(),
339                    hummock_meta_client.clone(),
340                    iceberg_memory_budget_bytes
341                        .expect("Iceberg memory budget must be resolved for dedicated startup"),
342                )
343            }
344            CompactorMode::SharedIceberg => unreachable!(),
345        },
346    ];
347
348    let telemetry_manager = TelemetryManager::new(
349        Arc::new(meta_client.clone()),
350        Arc::new(CompactorTelemetryCreator::new()),
351    );
352    // if the toml config file or env variable disables telemetry, do not watch system params change
353    // because if any of configs disable telemetry, we should never start it
354    if config.server.telemetry_enabled && telemetry_env_enabled() {
355        sub_tasks.push(telemetry_manager.start().await);
356    } else {
357        tracing::info!("Telemetry didn't start due to config");
358    }
359
360    let compactor_srv = CompactorServiceImpl::default();
361    let monitor_srv = MonitorServiceImpl::new(await_tree_reg, config.server.clone());
362    let server = tonic::transport::Server::builder()
363        .add_service(CompactorServiceServer::new(compactor_srv))
364        .add_service(configured_monitor_service_server(
365            MonitorServiceServer::new(monitor_srv),
366        ))
367        .monitored_serve_with_shutdown(
368            listen_addr,
369            "grpc-compactor-node-service",
370            TcpConfig {
371                tcp_nodelay: true,
372                keepalive_duration: None,
373            },
374            shutdown.clone().cancelled_owned(),
375        );
376    let _server_handle = tokio::spawn(server);
377
378    // Boot metrics service.
379    if config.server.metrics_level > MetricLevel::Disabled {
380        MetricsManager::boot_metrics_service(opts.prometheus_listener_addr.clone());
381    }
382
383    // All set, let the meta service know we're ready.
384    meta_client.activate(&advertise_addr).await.unwrap();
385
386    // Wait for the shutdown signal.
387    shutdown.cancelled().await;
388    // Run shutdown logic.
389    meta_client.try_unregister().await;
390}
391
392/// Fetches and runs compaction tasks under shared mode.
393///
394/// Returns when the `shutdown` token is triggered.
395pub async fn shared_compactor_serve(
396    listen_addr: SocketAddr,
397    opts: CompactorOpts,
398    shutdown: CancellationToken,
399) {
400    let config = load_config(&opts.config_path, &opts);
401    info!("Starting shared compactor node",);
402    info!("> config: {:?}", config);
403    info!(
404        "> debug assertions: {}",
405        if cfg!(debug_assertions) { "on" } else { "off" }
406    );
407    info!("> version: {} ({})", RW_VERSION, GIT_SHA);
408
409    let grpc_proxy_client = GrpcCompactorProxyClient::new(opts.proxy_rpc_endpoint.clone()).await;
410    let system_params_response = grpc_proxy_client
411        .get_system_params()
412        .await
413        .expect("Fail to get system params, the compactor pod cannot be started.");
414    let system_params = system_params_response.into_inner().params.unwrap();
415
416    let (
417        sstable_store,
418        memory_limiter,
419        heap_profiler,
420        await_tree_reg,
421        storage_opts,
422        compactor_metrics,
423    ) = Box::pin(prepare_start_parameters(
424        &opts,
425        config.clone(),
426        system_params.into(),
427    ))
428    .await;
429    let (sender, receiver) = mpsc::unbounded_channel();
430    let compactor_srv: CompactorServiceImpl = CompactorServiceImpl::new(sender);
431
432    let monitor_srv = MonitorServiceImpl::new(await_tree_reg.clone(), config.server.clone());
433
434    // Run a background heap profiler
435    heap_profiler.start();
436
437    let compaction_executor = Arc::new(CompactionExecutor::new(
438        opts.compaction_worker_threads_number,
439    ));
440    let compactor_context = CompactorContext {
441        storage_opts,
442        sstable_store,
443        compactor_metrics,
444        is_share_buffer_compact: false,
445        compaction_executor,
446        memory_limiter,
447        task_progress_manager: Default::default(),
448        await_tree_reg,
449    };
450
451    // TODO(shutdown): don't collect there's no need to gracefully shutdown them.
452    // Hold the join handle and tx to keep the compactor running.
453    let _compactor_handle = risingwave_storage::hummock::compactor::start_shared_compactor(
454        grpc_proxy_client,
455        receiver,
456        compactor_context,
457    );
458
459    let rpc_max_encoding_message_size_bytes = opts
460        .rpc_max_encoding_message_size_bytes
461        .unwrap_or(default_rpc_max_encoding_message_size_bytes());
462
463    let rpc_max_decoding_message_size_bytes = opts
464        .rpc_max_decoding_message_size_bytes
465        .unwrap_or(default_rpc_max_decoding_message_size_bytes());
466
467    let server = tonic::transport::Server::builder()
468        .add_service(
469            CompactorServiceServer::new(compactor_srv)
470                .max_decoding_message_size(rpc_max_decoding_message_size_bytes)
471                .max_encoding_message_size(rpc_max_encoding_message_size_bytes),
472        )
473        .add_service(configured_monitor_service_server(
474            MonitorServiceServer::new(monitor_srv),
475        ))
476        .monitored_serve_with_shutdown(
477            listen_addr,
478            "grpc-compactor-node-service",
479            TcpConfig {
480                tcp_nodelay: true,
481                keepalive_duration: None,
482            },
483            shutdown.clone().cancelled_owned(),
484        );
485
486    let _server_handle = tokio::spawn(server);
487
488    // Boot metrics service.
489    if config.server.metrics_level > MetricLevel::Disabled {
490        MetricsManager::boot_metrics_service(opts.prometheus_listener_addr.clone());
491    }
492
493    // Wait for the shutdown signal.
494    shutdown.cancelled().await;
495
496    // TODO(shutdown): shall we notify the proxy that we are shutting down?
497}