1use 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 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 {
119 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 .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
186fn 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
206pub 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 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 heap_profiler.start();
301
302 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 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 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 if config.server.metrics_level > MetricLevel::Disabled {
380 MetricsManager::boot_metrics_service(opts.prometheus_listener_addr.clone());
381 }
382
383 meta_client.activate(&advertise_addr).await.unwrap();
385
386 shutdown.cancelled().await;
388 meta_client.try_unregister().await;
390}
391
392pub 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 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 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 if config.server.metrics_level > MetricLevel::Disabled {
490 MetricsManager::boot_metrics_service(opts.prometheus_listener_addr.clone());
491 }
492
493 shutdown.cancelled().await;
495
496 }