1use std::net::SocketAddr;
16use std::sync::Arc;
17use std::time::Duration;
18
19use risingwave_batch::monitor::{
20 GLOBAL_BATCH_EXECUTOR_METRICS, GLOBAL_BATCH_MANAGER_METRICS, GLOBAL_BATCH_SPILL_METRICS,
21};
22use risingwave_batch::rpc::service::task_service::BatchServiceImpl;
23use risingwave_batch::spill::spill_op::SpillOp;
24use risingwave_batch::task::{BatchEnvironment, BatchManager};
25use risingwave_common::config::{
26 AsyncStackTraceOption, MAX_CONNECTION_WINDOW_SIZE, MetricLevel, STREAM_WINDOW_SIZE,
27 StorageMemoryConfig, load_config,
28};
29use risingwave_common::license::LicenseManager;
30use risingwave_common::lru::init_global_sequencer_args;
31use risingwave_common::monitor::{RouterExt, TcpConfig};
32use risingwave_common::secret::LocalSecretManager;
33use risingwave_common::system_param::local_manager::LocalSystemParamsManager;
34use risingwave_common::system_param::reader::SystemParamsRead;
35use risingwave_common::telemetry::manager::TelemetryManager;
36use risingwave_common::telemetry::telemetry_env_enabled;
37use risingwave_common::util::addr::HostAddr;
38use risingwave_common::util::pretty_bytes::convert;
39use risingwave_common::util::tokio_util::sync::CancellationToken;
40use risingwave_common::{DATA_DIRECTORY, GIT_SHA, RW_VERSION, STATE_STORE_URL};
41use risingwave_common_heap_profiling::HeapProfiler;
42use risingwave_common_service::{MetricsManager, ObserverManager, TracingExtractLayer};
43use risingwave_connector::source::iceberg::GLOBAL_ICEBERG_SCAN_METRICS;
44use risingwave_connector::source::monitor::GLOBAL_SOURCE_METRICS;
45use risingwave_dml::dml_manager::DmlManager;
46use risingwave_pb::common::WorkerType;
47use risingwave_pb::common::worker_node::Property;
48use risingwave_pb::compute::config_service_server::ConfigServiceServer;
49use risingwave_pb::configured_monitor_service_server;
50use risingwave_pb::health::health_server::HealthServer;
51use risingwave_pb::monitor_service::monitor_service_server::MonitorServiceServer;
52use risingwave_pb::stream_service::stream_service_server::StreamServiceServer;
53use risingwave_pb::task_service::batch_exchange_service_server::BatchExchangeServiceServer;
54use risingwave_pb::task_service::stream_exchange_service_server::StreamExchangeServiceServer;
55use risingwave_pb::task_service::task_service_server::TaskServiceServer;
56use risingwave_rpc_client::{ComputeClientPool, MetaClient};
57use risingwave_storage::StateStoreImpl;
58use risingwave_storage::hummock::MemoryLimiter;
59use risingwave_storage::hummock::compactor::{
60 CompactionExecutor, CompactorContext, new_compaction_await_tree_reg_ref, start_compactor,
61};
62use risingwave_storage::hummock::hummock_meta_client::MonitoredHummockMetaClient;
63use risingwave_storage::hummock::utils::HummockMemoryCollector;
64use risingwave_storage::monitor::{
65 GLOBAL_COMPACTOR_METRICS, GLOBAL_HUMMOCK_METRICS, GLOBAL_OBJECT_STORE_METRICS,
66 global_hummock_state_store_metrics, global_storage_metrics, monitor_cache,
67};
68use risingwave_storage::opts::StorageOpts;
69use risingwave_stream::executor::monitor::global_streaming_metrics;
70use risingwave_stream::task::{LocalStreamManager, StreamEnvironment};
71use tokio::sync::oneshot::Sender;
72use tokio::task::JoinHandle;
73use tower::Layer;
74
75use crate::ComputeNodeOpts;
76use crate::memory::config::{
77 MIN_COMPUTE_MEMORY_MB, batch_mem_limit, reserve_memory_bytes, storage_memory_config,
78};
79use crate::memory::manager::{MemoryManager, MemoryManagerConfig};
80use crate::observer::observer_manager::ComputeObserverNode;
81use crate::rpc::service::batch_exchange_service::BatchExchangeServiceImpl;
82use crate::rpc::service::config_service::ConfigServiceImpl;
83use crate::rpc::service::health_service::HealthServiceImpl;
84use crate::rpc::service::monitor_service::{AwaitTreeMiddlewareLayer, MonitorServiceImpl};
85use crate::rpc::service::stream_exchange_service::{
86 GLOBAL_STREAM_EXCHANGE_SERVICE_METRICS, StreamExchangeServiceImpl,
87};
88use crate::rpc::service::stream_service::StreamServiceImpl;
89use crate::telemetry::ComputeTelemetryCreator;
90pub async fn compute_node_serve(
94 listen_addr: SocketAddr,
95 advertise_addr: HostAddr,
96 opts: Arc<ComputeNodeOpts>,
97 shutdown: CancellationToken,
98) {
99 let config = Arc::new(load_config(&opts.config_path, &*opts));
101 info!("Starting compute node",);
102 info!("> config: {:?}", &*config);
103 info!(
104 "> debug assertions: {}",
105 if cfg!(debug_assertions) { "on" } else { "off" }
106 );
107 info!("> version: {} ({})", RW_VERSION, GIT_SHA);
108
109 let stream_config = Arc::new(config.streaming.clone());
111 let batch_config = Arc::new(config.batch.clone());
112
113 init_global_sequencer_args(
115 config
116 .streaming
117 .developer
118 .memory_controller_sequence_tls_step,
119 config
120 .streaming
121 .developer
122 .memory_controller_sequence_tls_lag,
123 );
124
125 let (meta_client, system_params) = MetaClient::register_new(
127 opts.meta_address.clone(),
128 WorkerType::ComputeNode,
129 &advertise_addr,
130 Property {
131 parallelism: opts.parallelism as u32,
132 is_streaming: opts.role.for_streaming(),
133 is_serving: opts.role.for_serving(),
134 internal_rpc_host_addr: "".to_owned(),
135 resource_group: Some(opts.resource_group.clone()),
136 is_iceberg_compactor: false,
137 },
138 Arc::new(config.meta.clone()),
139 )
140 .await;
141 let mut sub_tasks: Vec<(JoinHandle<()>, Sender<()>)> = vec![];
143 sub_tasks.push(MetaClient::start_heartbeat_loop(
144 meta_client.clone(),
145 Duration::from_millis(config.server.heartbeat_interval_ms as u64),
146 ));
147
148 let state_store_url = system_params.state_store();
149 let state_store_url = state_store_url.expose();
150
151 let data_directory = system_params.data_directory();
152
153 let embedded_compactor_enabled =
154 embedded_compactor_enabled(state_store_url, config.storage.disable_remote_compactor);
155
156 let (reserved_memory_bytes, non_reserved_memory_bytes) = reserve_memory_bytes(&opts);
157 let storage_memory_config = storage_memory_config(
158 non_reserved_memory_bytes,
159 embedded_compactor_enabled,
160 &config.storage,
161 !opts.role.for_streaming(),
162 );
163
164 let storage_memory_bytes = total_storage_memory_limit_bytes(&storage_memory_config);
165 let compute_memory_bytes = validate_compute_node_memory_config(
166 opts.total_memory_bytes,
167 reserved_memory_bytes,
168 storage_memory_bytes,
169 );
170 print_memory_config(
171 opts.total_memory_bytes,
172 compute_memory_bytes,
173 storage_memory_bytes,
174 &storage_memory_config,
175 embedded_compactor_enabled,
176 reserved_memory_bytes,
177 );
178
179 let storage_opts = Arc::new(StorageOpts::from((
180 &*config,
181 &system_params,
182 &storage_memory_config,
183 )));
184
185 let worker_id = meta_client.worker_id();
186 info!("Assigned worker node id {}", worker_id);
187
188 let source_metrics = Arc::new(GLOBAL_SOURCE_METRICS.clone());
190 let hummock_metrics = Arc::new(GLOBAL_HUMMOCK_METRICS.clone());
191 let streaming_metrics = Arc::new(global_streaming_metrics(config.server.metrics_level));
192 let batch_executor_metrics = Arc::new(GLOBAL_BATCH_EXECUTOR_METRICS.clone());
193 let batch_manager_metrics = Arc::new(GLOBAL_BATCH_MANAGER_METRICS.clone());
194 let stream_exchange_srv_metrics = Arc::new(GLOBAL_STREAM_EXCHANGE_SERVICE_METRICS.clone());
195 let batch_spill_metrics = Arc::new(GLOBAL_BATCH_SPILL_METRICS.clone());
196 let iceberg_scan_metrics = Arc::new(GLOBAL_ICEBERG_SCAN_METRICS.clone());
197
198 let state_store_metrics = Arc::new(global_hummock_state_store_metrics(
200 config.server.metrics_level,
201 ));
202 let object_store_metrics = Arc::new(GLOBAL_OBJECT_STORE_METRICS.clone());
203 let storage_metrics = Arc::new(global_storage_metrics(config.server.metrics_level));
204 let compactor_metrics = Arc::new(GLOBAL_COMPACTOR_METRICS.clone());
205 let hummock_meta_client = Arc::new(MonitoredHummockMetaClient::new(
206 meta_client.clone(),
207 hummock_metrics.clone(),
208 ));
209
210 let await_tree_config = match &config.streaming.async_stack_trace {
211 AsyncStackTraceOption::Off => None,
212 c => await_tree::ConfigBuilder::default()
213 .verbose(c.is_verbose().unwrap())
214 .build()
215 .ok(),
216 };
217 if let Err(existing_url) = STATE_STORE_URL.set(state_store_url.to_owned()) {
220 assert_eq!(
221 existing_url, state_store_url,
222 "STATE_STORE_URL already set with different value"
223 );
224 }
225
226 if let Err(existing_dir) = DATA_DIRECTORY.set(data_directory.to_owned()) {
229 assert_eq!(
230 existing_dir, data_directory,
231 "DATA_DIRECTORY already set with different value"
232 );
233 }
234 LicenseManager::get().refresh(system_params.license_key());
235 let state_store = Box::pin(StateStoreImpl::new(
236 state_store_url,
237 opts.role,
238 storage_opts.clone(),
239 hummock_meta_client.clone(),
240 state_store_metrics.clone(),
241 object_store_metrics,
242 storage_metrics.clone(),
243 compactor_metrics.clone(),
244 await_tree_config.clone(),
245 system_params.use_new_object_prefix_strategy(),
246 ))
247 .await
248 .unwrap();
249
250 LocalSecretManager::init(
251 opts.temp_secret_file_dir.clone(),
252 meta_client.cluster_id().to_owned(),
253 worker_id,
254 );
255
256 let batch_client_pool = Arc::new(ComputeClientPool::new(
258 config.batch_exchange_connection_pool_size(),
259 config.batch.developer.compute_client_config.clone(),
260 ));
261 let system_params_manager = Arc::new(LocalSystemParamsManager::new(system_params.clone()));
262 let compute_observer_node =
263 ComputeObserverNode::new(system_params_manager.clone(), batch_client_pool.clone());
264 let observer_manager =
265 ObserverManager::new_with_meta_client(meta_client.clone(), compute_observer_node).await;
266 observer_manager.start().await;
267
268 if let Some(storage) = state_store.as_hummock() {
269 if embedded_compactor_enabled {
270 tracing::info!("start embedded compactor");
271 let memory_limiter = Arc::new(MemoryLimiter::new(
272 storage_opts.compactor_memory_limit_mb as u64 * 1024 * 1024 / 2,
273 ));
274
275 let compaction_executor = Arc::new(CompactionExecutor::new(Some(1)));
276 let compactor_context = CompactorContext {
277 storage_opts,
278 sstable_store: storage.sstable_store(),
279 compactor_metrics: compactor_metrics.clone(),
280 is_share_buffer_compact: false,
281 compaction_executor,
282 memory_limiter,
283
284 task_progress_manager: Default::default(),
285 await_tree_reg: await_tree_config
286 .clone()
287 .map(new_compaction_await_tree_reg_ref),
288 };
289
290 let (handle, shutdown_sender) = start_compactor(
291 compactor_context,
292 hummock_meta_client.clone(),
293 storage.object_id_manager().clone(),
294 storage.compaction_catalog_manager_ref(),
295 );
296 sub_tasks.push((handle, shutdown_sender));
297 }
298 let flush_limiter = storage.get_memory_limiter();
299 let memory_collector = Arc::new(HummockMemoryCollector::new(
300 storage.sstable_store(),
301 flush_limiter,
302 storage_memory_config,
303 ));
304 monitor_cache(memory_collector);
305 let backup_reader = storage.backup_reader();
306 let system_params_mgr = system_params_manager.clone();
307 tokio::spawn(async move {
308 backup_reader
309 .watch_config_change(system_params_mgr.watch_params())
310 .await;
311 });
312 }
313
314 let batch_mgr = Arc::new(BatchManager::new(
316 config.batch.clone(),
317 batch_manager_metrics,
318 batch_mem_limit(compute_memory_bytes, opts.role.for_serving()),
319 await_tree_config.clone(),
320 ));
321
322 let target_memory = if let Some(v) = opts.memory_manager_target_bytes {
323 v
324 } else {
325 compute_memory_bytes + storage_memory_bytes
326 };
327
328 let memory_mgr = MemoryManager::new(MemoryManagerConfig {
329 target_memory,
330 threshold_aggressive: config
331 .streaming
332 .developer
333 .memory_controller_threshold_aggressive,
334 threshold_graceful: config
335 .streaming
336 .developer
337 .memory_controller_threshold_graceful,
338 threshold_stable: config
339 .streaming
340 .developer
341 .memory_controller_threshold_stable,
342 eviction_factor_stable: config
343 .streaming
344 .developer
345 .memory_controller_eviction_factor_stable,
346 eviction_factor_graceful: config
347 .streaming
348 .developer
349 .memory_controller_eviction_factor_graceful,
350 eviction_factor_aggressive: config
351 .streaming
352 .developer
353 .memory_controller_eviction_factor_aggressive,
354 metrics: streaming_metrics.clone(),
355 });
356
357 tokio::spawn(
359 memory_mgr.clone().run(Duration::from_millis(
360 config
361 .streaming
362 .developer
363 .memory_controller_update_interval_ms as _,
364 )),
365 );
366
367 let heap_profiler = HeapProfiler::new(
368 opts.total_memory_bytes,
369 config.server.heap_profiling.clone(),
370 );
371 heap_profiler.start();
373
374 let dml_mgr = Arc::new(DmlManager::new(
375 worker_id,
376 config.streaming.developer.dml_channel_initial_permits,
377 ));
378
379 let batch_env = BatchEnvironment::new(
381 batch_mgr.clone(),
382 advertise_addr.clone(),
383 batch_config,
384 worker_id,
385 state_store.clone(),
386 batch_executor_metrics.clone(),
387 batch_client_pool,
388 dml_mgr.clone(),
389 source_metrics.clone(),
390 batch_spill_metrics.clone(),
391 iceberg_scan_metrics.clone(),
392 config.server.metrics_level,
393 );
394
395 let stream_client_pool = Arc::new(ComputeClientPool::new(
397 config.streaming_exchange_connection_pool_size(),
398 config.streaming.developer.compute_client_config.clone(),
399 ));
400 let stream_env = StreamEnvironment::new(
401 advertise_addr.clone(),
402 stream_config,
403 worker_id,
404 state_store.clone(),
405 dml_mgr,
406 system_params_manager.clone(),
407 source_metrics,
408 meta_client.clone(),
409 stream_client_pool,
410 );
411
412 let stream_mgr = LocalStreamManager::new(
413 stream_env.clone(),
414 streaming_metrics.clone(),
415 await_tree_config.clone(),
416 memory_mgr.get_watermark_sequence(),
417 );
418
419 let batch_srv = BatchServiceImpl::new(batch_mgr.clone(), batch_env);
421 let batch_exchange_srv = BatchExchangeServiceImpl::new(batch_mgr.clone());
422 let stream_exchange_srv =
423 StreamExchangeServiceImpl::new(stream_mgr.clone(), stream_exchange_srv_metrics);
424 let stream_srv = StreamServiceImpl::new(stream_mgr.clone(), stream_env.clone());
425 let (meta_cache, block_cache, hummock_storage) = if let Some(hummock) = state_store.as_hummock()
426 {
427 (
428 Some(hummock.sstable_store().meta_cache().clone()),
429 Some(hummock.sstable_store().block_cache().clone()),
430 Some(hummock.clone()),
431 )
432 } else {
433 (None, None, None)
434 };
435 let monitor_srv = MonitorServiceImpl::new(
436 stream_mgr.clone(),
437 batch_mgr.clone(),
438 config.server.clone(),
439 meta_cache.clone(),
440 block_cache.clone(),
441 hummock_storage,
442 );
443 let config_srv = ConfigServiceImpl::new(
444 batch_mgr.clone(),
445 stream_mgr.clone(),
446 meta_cache,
447 block_cache,
448 );
449 let health_srv = HealthServiceImpl::new();
450
451 let telemetry_manager = TelemetryManager::new(
452 Arc::new(meta_client.clone()),
453 Arc::new(ComputeTelemetryCreator::new()),
454 );
455
456 if config.server.telemetry_enabled && telemetry_env_enabled() {
459 sub_tasks.push(telemetry_manager.start().await);
460 } else {
461 tracing::info!("Telemetry didn't start due to config");
462 }
463
464 #[cfg(not(madsim))]
466 if config.batch.enable_spill {
467 SpillOp::clean_spill_directory().await.unwrap();
468 }
469
470 let server = tonic::transport::Server::builder()
471 .initial_connection_window_size(MAX_CONNECTION_WINDOW_SIZE)
472 .initial_stream_window_size(STREAM_WINDOW_SIZE)
473 .http2_max_pending_accept_reset_streams(Some(config.server.grpc_max_reset_stream as usize))
474 .layer(TracingExtractLayer::new())
475 .add_service({
477 let await_tree_reg = batch_mgr.await_tree_reg().cloned();
478 let srv = TaskServiceServer::new(batch_srv).max_decoding_message_size(usize::MAX);
479 #[cfg(madsim)]
480 {
481 srv
482 }
483 #[cfg(not(madsim))]
484 {
485 AwaitTreeMiddlewareLayer::new_optional(await_tree_reg).layer(srv)
486 }
487 })
488 .add_service(
489 BatchExchangeServiceServer::new(batch_exchange_srv)
490 .max_decoding_message_size(usize::MAX),
491 )
492 .add_service(
493 StreamExchangeServiceServer::new(stream_exchange_srv)
494 .max_decoding_message_size(usize::MAX),
495 )
496 .add_service({
497 let await_tree_reg = stream_srv.mgr.await_tree_reg().cloned();
498 let srv = StreamServiceServer::new(stream_srv).max_decoding_message_size(usize::MAX);
499 #[cfg(madsim)]
500 {
501 srv
502 }
503 #[cfg(not(madsim))]
504 {
505 AwaitTreeMiddlewareLayer::new_optional(await_tree_reg).layer(srv)
506 }
507 })
508 .add_service(configured_monitor_service_server(
509 MonitorServiceServer::new(monitor_srv),
510 ))
511 .add_service(ConfigServiceServer::new(config_srv))
512 .add_service(HealthServer::new(health_srv))
513 .monitored_serve_with_shutdown(
514 listen_addr,
515 "grpc-compute-node-service",
516 TcpConfig {
517 tcp_nodelay: true,
518 keepalive_duration: None,
519 },
520 shutdown.clone().cancelled_owned(),
521 );
522 let _server_handle = tokio::spawn(server);
523
524 if config.server.metrics_level > MetricLevel::Disabled {
526 MetricsManager::boot_metrics_service(opts.prometheus_listener_addr.clone());
527 }
528
529 meta_client.activate(&advertise_addr).await.unwrap();
531 shutdown.cancelled().await;
533
534 meta_client.try_unregister().await;
538 let _ = stream_mgr.shutdown().await;
540
541 }
546
547fn validate_compute_node_memory_config(
553 cn_total_memory_bytes: usize,
554 reserved_memory_bytes: usize,
555 storage_memory_bytes: usize,
556) -> usize {
557 if storage_memory_bytes > cn_total_memory_bytes {
558 tracing::warn!(
559 "The storage memory exceeds the total compute node memory:\nTotal compute node memory: {}\nStorage memory: {}\nWe recommend that at least 4 GiB memory should be reserved for RisingWave. Please increase the total compute node memory or decrease the storage memory in configurations.",
560 convert(cn_total_memory_bytes as _),
561 convert(storage_memory_bytes as _)
562 );
563 MIN_COMPUTE_MEMORY_MB << 20
564 } else if storage_memory_bytes + (MIN_COMPUTE_MEMORY_MB << 20) + reserved_memory_bytes
565 >= cn_total_memory_bytes
566 {
567 tracing::warn!(
568 "Not enough memory remains for compute workloads and other system usage:\nTotal compute node memory: {}\nStorage memory: {}\nWe recommend reserving at least 4 GiB of memory for RisingWave. Please increase the total compute node memory or decrease the configured storage memory.",
569 convert(cn_total_memory_bytes as _),
570 convert(storage_memory_bytes as _)
571 );
572 MIN_COMPUTE_MEMORY_MB << 20
573 } else {
574 cn_total_memory_bytes - storage_memory_bytes - reserved_memory_bytes
575 }
576}
577
578fn total_storage_memory_limit_bytes(storage_memory_config: &StorageMemoryConfig) -> usize {
581 let total_storage_memory_mb = storage_memory_config.block_cache_capacity_mb
582 + storage_memory_config.meta_cache_capacity_mb
583 + storage_memory_config.shared_buffer_capacity_mb
584 + storage_memory_config.compactor_memory_limit_mb;
585 total_storage_memory_mb << 20
586}
587
588fn embedded_compactor_enabled(state_store_url: &str, disable_remote_compactor: bool) -> bool {
590 state_store_url.starts_with("hummock+memory")
592 || state_store_url.starts_with("hummock+disk")
593 || disable_remote_compactor
594}
595
596fn print_memory_config(
598 cn_total_memory_bytes: usize,
599 compute_memory_bytes: usize,
600 storage_memory_bytes: usize,
601 storage_memory_config: &StorageMemoryConfig,
602 embedded_compactor_enabled: bool,
603 reserved_memory_bytes: usize,
604) {
605 let memory_config = format!(
606 "Memory outline:\n\
607 > total_memory: {}\n\
608 > storage_memory: {}\n\
609 > block_cache_capacity: {}\n\
610 > meta_cache_capacity: {}\n\
611 > shared_buffer_capacity: {}\n\
612 > compactor_memory_limit: {}\n\
613 > compute_memory: {}\n\
614 > reserved_memory: {}",
615 convert(cn_total_memory_bytes as _),
616 convert(storage_memory_bytes as _),
617 convert((storage_memory_config.block_cache_capacity_mb << 20) as _),
618 convert((storage_memory_config.meta_cache_capacity_mb << 20) as _),
619 convert((storage_memory_config.shared_buffer_capacity_mb << 20) as _),
620 if embedded_compactor_enabled {
621 convert((storage_memory_config.compactor_memory_limit_mb << 20) as _)
622 } else {
623 "Not enabled".to_owned()
624 },
625 convert(compute_memory_bytes as _),
626 convert(reserved_memory_bytes as _),
627 );
628 info!("{}", memory_config);
629}