1use std::collections::{HashMap, HashSet};
16use std::sync::Arc;
17
18use parking_lot::RwLock;
19use risingwave_common::bitmap::Bitmap;
20use risingwave_common::hash::WorkerSlotMapping;
21use risingwave_common::vnode_mapping::vnode_placement::place_vnode;
22use risingwave_meta_model::{TableId, WorkerId};
23use risingwave_pb::common::{WorkerNode, WorkerType};
24use risingwave_pb::meta::serving_table_vnode_mappings::PbServingTableVnodeMapping;
25use risingwave_pb::meta::subscribe_response::{Info, Operation};
26use risingwave_pb::meta::table_fragments::fragment::FragmentDistributionType;
27use risingwave_pb::meta::{
28 FragmentWorkerSlotMapping, FragmentWorkerSlotMappings, PbServingTableVnodeMappings,
29 PbTableRefillRuntimeConfig,
30};
31use tokio::sync::oneshot::Sender;
32use tokio::task::JoinHandle;
33
34use crate::MetaResult;
35use crate::controller::fragment::FragmentServingInfo;
36use crate::controller::session_params::SessionParamsControllerRef;
37use crate::manager::{LocalNotification, MetadataManager, NotificationManagerRef, WorkerKey};
38use crate::model::FragmentId;
39
40pub type ServingVnodeMappingRef = Arc<ServingVnodeMapping>;
41
42#[derive(Default)]
43pub struct ServingVnodeMapping {
44 serving_vnode_mappings: RwLock<HashMap<FragmentId, WorkerSlotMapping>>,
45}
46
47impl ServingVnodeMapping {
48 pub fn all(&self) -> HashMap<FragmentId, WorkerSlotMapping> {
49 self.serving_vnode_mappings.read().clone()
50 }
51
52 pub(crate) fn table_vnode_mappings_by_worker(
53 &self,
54 worker_ids: impl IntoIterator<Item = WorkerId>,
55 fragment_serving_infos: &HashMap<FragmentId, FragmentServingInfo>,
56 ) -> HashMap<WorkerId, HashMap<TableId, Bitmap>> {
57 let mut result = worker_ids
58 .into_iter()
59 .map(|worker_id| (worker_id, HashMap::new()))
60 .collect::<HashMap<_, _>>();
61 let serving_vnode_mappings = self.all();
62
63 for (fragment_id, info) in fragment_serving_infos {
64 let Some(table_id) = info.result_table_id else {
65 continue;
66 };
67 let Some(mapping) = serving_vnode_mappings.get(fragment_id) else {
68 continue;
69 };
70
71 for (worker_slot_id, bitmap) in mapping.to_bitmaps() {
72 let Some(table_mappings) = result.get_mut(&worker_slot_id.worker_id()) else {
73 continue;
74 };
75 table_mappings
76 .entry(table_id)
77 .and_modify(|current| *current |= &bitmap)
78 .or_insert(bitmap);
79 }
80 }
81 result
82 }
83
84 pub fn upsert(
87 &self,
88 fragment_serving_infos: &HashMap<FragmentId, FragmentServingInfo>,
89 workers: &[WorkerNode],
90 max_serving_parallelism: Option<u64>,
91 ) -> (HashMap<FragmentId, WorkerSlotMapping>, HashSet<FragmentId>) {
92 let mut serving_vnode_mappings = self.serving_vnode_mappings.write();
93 Self::upsert_locked(
94 &mut serving_vnode_mappings,
95 fragment_serving_infos,
96 workers,
97 max_serving_parallelism,
98 )
99 }
100
101 pub(crate) fn reconcile(
105 &self,
106 fragment_serving_infos: &HashMap<FragmentId, FragmentServingInfo>,
107 workers: &[WorkerNode],
108 max_serving_parallelism: Option<u64>,
109 ) -> (HashMap<FragmentId, WorkerSlotMapping>, HashSet<FragmentId>) {
110 let mut serving_vnode_mappings = self.serving_vnode_mappings.write();
111 serving_vnode_mappings
112 .retain(|fragment_id, _| fragment_serving_infos.contains_key(fragment_id));
113 Self::upsert_locked(
114 &mut serving_vnode_mappings,
115 fragment_serving_infos,
116 workers,
117 max_serving_parallelism,
118 )
119 }
120
121 fn upsert_locked(
122 serving_vnode_mappings: &mut HashMap<FragmentId, WorkerSlotMapping>,
123 fragment_serving_infos: &HashMap<FragmentId, FragmentServingInfo>,
124 workers: &[WorkerNode],
125 max_serving_parallelism: Option<u64>,
126 ) -> (HashMap<FragmentId, WorkerSlotMapping>, HashSet<FragmentId>) {
127 let mut upserted: HashMap<FragmentId, WorkerSlotMapping> = HashMap::default();
128 let mut failed: HashSet<FragmentId> = HashSet::default();
129 for (fragment_id, info) in fragment_serving_infos {
130 let new_mapping = {
131 let old_mapping = serving_vnode_mappings.get(fragment_id);
132 let max_parallelism = match info.distribution_type {
133 FragmentDistributionType::Unspecified => unreachable!(),
134 FragmentDistributionType::Single => Some(1),
135 FragmentDistributionType::Hash => None,
136 }
137 .or_else(|| max_serving_parallelism.map(|p| p as usize));
138 place_vnode(old_mapping, workers, max_parallelism, info.vnode_count)
139 };
140 match new_mapping {
141 None => {
142 serving_vnode_mappings.remove(fragment_id);
143 failed.insert(*fragment_id);
144 }
145 Some(mapping) => {
146 serving_vnode_mappings.insert(*fragment_id, mapping.clone());
147 upserted.insert(*fragment_id, mapping);
148 }
149 }
150 }
151 (upserted, failed)
152 }
153
154 fn remove(&self, fragment_ids: &[FragmentId]) {
155 let mut mappings = self.serving_vnode_mappings.write();
156 for fragment_id in fragment_ids {
157 mappings.remove(fragment_id);
158 }
159 }
160}
161
162pub(crate) fn to_fragment_worker_slot_mapping(
163 mappings: &HashMap<FragmentId, WorkerSlotMapping>,
164) -> Vec<FragmentWorkerSlotMapping> {
165 mappings
166 .iter()
167 .map(|(&fragment_id, mapping)| FragmentWorkerSlotMapping {
168 fragment_id,
169 mapping: Some(mapping.to_protobuf()),
170 })
171 .collect()
172}
173
174pub(crate) fn to_deleted_fragment_worker_slot_mapping(
175 fragment_ids: impl Iterator<Item = FragmentId>,
176) -> Vec<FragmentWorkerSlotMapping> {
177 fragment_ids
178 .map(|fragment_id| FragmentWorkerSlotMapping {
179 fragment_id,
180 mapping: None,
181 })
182 .collect()
183}
184
185pub async fn on_meta_start(
186 notification_manager: NotificationManagerRef,
187 metadata_manager: &MetadataManager,
188 serving_vnode_mapping: ServingVnodeMappingRef,
189 max_serving_parallelism: Option<u64>,
190) {
191 let (serving_compute_nodes, fragment_serving_infos) = fetch_serving_infos(metadata_manager)
192 .await
193 .expect("fail to fetch serving infos");
194 let (mappings, failed) = serving_vnode_mapping.reconcile(
195 &fragment_serving_infos,
196 &serving_compute_nodes,
197 max_serving_parallelism,
198 );
199 tracing::debug!(
200 "Initialize serving vnode mapping snapshot for fragments {:?}.",
201 mappings.keys()
202 );
203 if !failed.is_empty() {
204 tracing::warn!(
205 "Fail to update serving vnode mapping for fragments {:?}.",
206 failed
207 );
208 }
209 notification_manager.notify_frontend_without_version(
210 Operation::Snapshot,
211 Info::ServingWorkerSlotMappings(FragmentWorkerSlotMappings {
212 mappings: to_fragment_worker_slot_mapping(&mappings),
213 }),
214 );
215}
216
217pub(crate) async fn fetch_serving_infos(
218 metadata_manager: &MetadataManager,
219) -> MetaResult<(Vec<WorkerNode>, HashMap<FragmentId, FragmentServingInfo>)> {
220 let fragment_serving_infos = metadata_manager
221 .catalog_controller
222 .fragment_serving_infos()
223 .await?;
224 let serving_compute_nodes = metadata_manager
225 .cluster_controller
226 .list_active_serving_workers()
227 .await?;
228 Ok((
229 serving_compute_nodes,
230 fragment_serving_infos
231 .into_iter()
232 .map(|(fragment_id, info)| (fragment_id as FragmentId, info))
233 .collect(),
234 ))
235}
236
237pub(crate) async fn sync_serving_table_vnode_mappings_to_hummock(
238 notification_manager: &NotificationManagerRef,
239 serving_vnode_mapping: &ServingVnodeMappingRef,
240 workers: &[WorkerNode],
241 fragment_serving_infos: &HashMap<FragmentId, FragmentServingInfo>,
242) {
243 let mut table_vnode_mappings = serving_vnode_mapping.table_vnode_mappings_by_worker(
244 workers.iter().map(|worker| worker.id),
245 fragment_serving_infos,
246 );
247 for worker in workers {
248 let Some(host) = worker.host.clone() else {
249 tracing::warn!(worker_id = %worker.id, "serving worker host not found");
250 continue;
251 };
252 let table_vnode_mapping = table_vnode_mappings
253 .remove(&worker.id)
254 .expect("requested worker must have a table vnode mapping");
255 let config = PbTableRefillRuntimeConfig {
256 serving_table_vnode_mappings: Some(to_pb_serving_table_vnode_mappings(
257 &table_vnode_mapping,
258 )),
259 ..Default::default()
260 };
261 notification_manager
262 .notify_hummock_targeted_update(WorkerKey(host), Info::TableRefillRuntimeConfig(config))
263 .await;
264 }
265}
266
267pub fn start_serving_vnode_mapping_worker(
268 notification_manager: NotificationManagerRef,
269 metadata_manager: MetadataManager,
270 serving_vnode_mapping: ServingVnodeMappingRef,
271 session_params: SessionParamsControllerRef,
272) -> (JoinHandle<()>, Sender<()>) {
273 let (local_notification_tx, mut local_notification_rx) = tokio::sync::mpsc::unbounded_channel();
274 let (shutdown_tx, mut shutdown_rx) = tokio::sync::oneshot::channel();
275 notification_manager.insert_local_sender(local_notification_tx);
276 let join_handle = tokio::spawn(async move {
277 let reset = || async {
278 let (workers, fragment_serving_infos) = fetch_serving_infos(&metadata_manager)
279 .await
280 .expect("fail to fetch serving infos");
281 let max_serving_parallelism = session_params
282 .get_params()
283 .await
284 .batch_parallelism()
285 .map(|p| p.get());
286 let (mappings, failed) = serving_vnode_mapping.reconcile(
287 &fragment_serving_infos,
288 &workers,
289 max_serving_parallelism,
290 );
291 tracing::debug!(
292 "Update serving vnode mapping snapshot for fragments {:?}.",
293 mappings.keys()
294 );
295 if !failed.is_empty() {
296 tracing::warn!(
297 "Fail to update serving vnode mapping for fragments {:?}.",
298 failed
299 );
300 }
301 notification_manager.notify_frontend_without_version(
302 Operation::Snapshot,
303 Info::ServingWorkerSlotMappings(FragmentWorkerSlotMappings {
304 mappings: to_fragment_worker_slot_mapping(&mappings),
305 }),
306 );
307 sync_serving_table_vnode_mappings_to_hummock(
308 ¬ification_manager,
309 &serving_vnode_mapping,
310 &workers,
311 &fragment_serving_infos,
312 )
313 .await;
314 };
315 loop {
316 tokio::select! {
317 notification = local_notification_rx.recv() => {
318 match notification {
319 Some(notification) => {
320 match notification {
321 LocalNotification::WorkerNodeActivated(w) | LocalNotification::WorkerNodeDeleted(w) => {
322 if w.r#type() != WorkerType::ComputeNode || !w.property.as_ref().is_some_and(|p| p.is_serving) {
323 continue;
324 }
325 reset().await;
326 }
327 LocalNotification::BatchParallelismChange => {
328 reset().await;
329 }
330 LocalNotification::ServingFragmentMappingsUpsert(fragment_ids) => {
331 if fragment_ids.is_empty() {
332 continue;
333 }
334 let (workers, fragment_serving_infos) = fetch_serving_infos(&metadata_manager)
335 .await
336 .expect("fail to fetch serving infos");
337 let filtered_fragment_serving_infos = fragment_ids.iter().filter_map(|frag_id| {
338 match fragment_serving_infos.get(frag_id) {
339 Some(info) => Some((*frag_id, info.clone())),
340 None => {
341 tracing::warn!(fragment_id = %frag_id, "fragment serving info not found");
342 None
343 }
344 }
345 }).collect();
346 let max_serving_parallelism = session_params
347 .get_params()
348 .await
349 .batch_parallelism()
350 .map(|p|p.get());
351 let (upserted, failed) = serving_vnode_mapping.upsert(&filtered_fragment_serving_infos, &workers, max_serving_parallelism);
352 if !upserted.is_empty() {
353 tracing::debug!("Update serving vnode mapping for fragments {:?}.", upserted.keys());
354 notification_manager.notify_frontend_without_version(Operation::Update, Info::ServingWorkerSlotMappings(FragmentWorkerSlotMappings{ mappings: to_fragment_worker_slot_mapping(&upserted) }));
355 }
356 if !failed.is_empty() {
357 tracing::warn!("Fail to update serving vnode mapping for fragments {:?}.", failed);
358 notification_manager.notify_frontend_without_version(Operation::Delete, Info::ServingWorkerSlotMappings(FragmentWorkerSlotMappings{ mappings: to_deleted_fragment_worker_slot_mapping(failed.iter().cloned())}));
359 }
360 sync_serving_table_vnode_mappings_to_hummock(
361 ¬ification_manager,
362 &serving_vnode_mapping,
363 &workers,
364 &fragment_serving_infos,
365 )
366 .await;
367 }
368 LocalNotification::ServingFragmentMappingsDelete(fragment_ids) => {
369 if fragment_ids.is_empty() {
370 continue;
371 }
372
373 tracing::debug!("Delete serving vnode mapping for fragments {:?}.", fragment_ids);
374
375 serving_vnode_mapping.remove(&fragment_ids);
376 notification_manager.notify_frontend_without_version(Operation::Delete, Info::ServingWorkerSlotMappings(FragmentWorkerSlotMappings{ mappings: to_deleted_fragment_worker_slot_mapping(fragment_ids.iter().cloned()) }));
377 let (workers, fragment_serving_infos) = fetch_serving_infos(&metadata_manager)
378 .await
379 .expect("fail to fetch serving infos");
380 sync_serving_table_vnode_mappings_to_hummock(
381 ¬ification_manager,
382 &serving_vnode_mapping,
383 &workers,
384 &fragment_serving_infos,
385 )
386 .await;
387 }
388 _ => {}
389 }
390 }
391 None => {
392 return;
393 }
394 }
395 }
396 _ = &mut shutdown_rx => {
397 return;
398 }
399 }
400 }
401 });
402 (join_handle, shutdown_tx)
403}
404
405pub(crate) fn to_pb_serving_table_vnode_mappings(
406 table_vnode_mapping: &HashMap<TableId, Bitmap>,
407) -> PbServingTableVnodeMappings {
408 PbServingTableVnodeMappings {
409 mappings: table_vnode_mapping
410 .iter()
411 .map(|(table_id, bitmap)| PbServingTableVnodeMapping {
412 table_id: table_id.as_raw_id(),
413 bitmap: Some(bitmap.to_protobuf()),
414 })
415 .collect(),
416 }
417}
418
419#[cfg(test)]
420mod tests {
421 use risingwave_common::hash::{VirtualNode, WorkerSlotId};
422 use risingwave_pb::common::{HostAddress, worker_node};
423 use risingwave_pb::meta::subscribe_response::Info;
424 use risingwave_pb::meta::{SubscribeResponse, SubscribeType};
425 use tokio::sync::mpsc;
426
427 use super::*;
428 use crate::controller::SqlMetaStore;
429 use crate::manager::NotificationManager;
430
431 fn serving_worker(id: u32) -> WorkerNode {
432 WorkerNode {
433 id: id.into(),
434 r#type: WorkerType::ComputeNode.into(),
435 host: Some(HostAddress {
436 host: format!("host{}", id),
437 port: id as i32,
438 }),
439 state: worker_node::State::Running as i32,
440 property: Some(worker_node::Property {
441 is_serving: true,
442 parallelism: 1,
443 ..Default::default()
444 }),
445 ..Default::default()
446 }
447 }
448
449 #[test]
450 fn test_table_vnode_mappings_merge_worker_slots() {
451 let worker1 = WorkerId::new(1);
452 let worker2 = WorkerId::new(2);
453 let worker3 = WorkerId::new(3);
454 let slot1 = WorkerSlotId::new(worker1, 0);
455 let slot2 = WorkerSlotId::new(worker1, 1);
456 let slot3 = WorkerSlotId::new(worker2, 0);
457 let fragment_id = FragmentId::new(233);
458 let table_id = TableId::new(234);
459 let mapping = WorkerSlotMapping::new_uniform(
460 [slot1, slot2, slot3].into_iter(),
461 VirtualNode::COUNT_FOR_TEST,
462 );
463 let slot_bitmaps = mapping.to_bitmaps();
464 let serving_vnode_mapping = ServingVnodeMapping {
465 serving_vnode_mappings: RwLock::new(HashMap::from([(fragment_id, mapping)])),
466 };
467 let fragment_serving_infos = HashMap::from([(
468 fragment_id,
469 FragmentServingInfo {
470 result_table_id: Some(table_id),
471 distribution_type: FragmentDistributionType::Hash,
472 vnode_count: VirtualNode::COUNT_FOR_TEST,
473 },
474 )]);
475
476 let mappings = serving_vnode_mapping
477 .table_vnode_mappings_by_worker([worker1, worker2, worker3], &fragment_serving_infos);
478 let mut worker1_bitmap = slot_bitmaps[&slot1].clone();
479 worker1_bitmap |= &slot_bitmaps[&slot2];
480 assert_eq!(mappings[&worker1][&table_id], worker1_bitmap);
481 assert_eq!(mappings[&worker2][&table_id], slot_bitmaps[&slot3]);
482 assert!(mappings[&worker3].is_empty());
483 }
484
485 #[tokio::test]
486 async fn test_sync_serving_table_vnode_mappings_to_hummock_targets_only_result_tables() {
487 let notification_manager =
488 Arc::new(NotificationManager::new(SqlMetaStore::for_test().await).await);
489 let worker1 = serving_worker(1);
490 let worker2 = serving_worker(2);
491 let workers = vec![worker1.clone(), worker2.clone()];
492 let worker_key1 = WorkerKey(worker1.host.clone().unwrap());
493 let worker_key2 = WorkerKey(worker2.host.clone().unwrap());
494 let (tx1, mut rx1) = mpsc::unbounded_channel();
495 let (tx2, mut rx2) = mpsc::unbounded_channel();
496 notification_manager.insert_sender(SubscribeType::Hummock, worker_key1, tx1);
497 notification_manager.insert_sender(SubscribeType::Hummock, worker_key2, tx2);
498
499 let table_id = TableId::new(233);
500 let fragment_id = FragmentId::new(234);
501 let private_fragment_id = FragmentId::new(235);
502 let serving_vnode_mapping = Arc::new(ServingVnodeMapping::default());
503 let fragment_serving_infos = HashMap::from([
504 (
505 fragment_id,
506 FragmentServingInfo {
507 result_table_id: Some(table_id),
508 distribution_type: FragmentDistributionType::Hash,
509 vnode_count: VirtualNode::COUNT_FOR_TEST,
510 },
511 ),
512 (
513 private_fragment_id,
514 FragmentServingInfo {
515 result_table_id: None,
516 distribution_type: FragmentDistributionType::Hash,
517 vnode_count: VirtualNode::COUNT_FOR_TEST,
518 },
519 ),
520 ]);
521 let (upserted, failed) =
522 serving_vnode_mapping.upsert(&fragment_serving_infos, &workers, None);
523 assert_eq!(upserted.len(), 2);
524 assert!(failed.is_empty());
525 let expected_bitmaps = upserted[&fragment_id]
526 .to_bitmaps()
527 .into_iter()
528 .map(|(worker_slot_id, bitmap)| (worker_slot_id.worker_id(), bitmap))
529 .collect::<HashMap<_, _>>();
530 serving_vnode_mapping.serving_vnode_mappings.write().insert(
533 private_fragment_id,
534 WorkerSlotMapping::new_single(WorkerSlotId::new(worker1.id, 0)),
535 );
536
537 sync_serving_table_vnode_mappings_to_hummock(
538 ¬ification_manager,
539 &serving_vnode_mapping,
540 &workers,
541 &fragment_serving_infos,
542 )
543 .await;
544
545 let response1 = rx1.recv().await.unwrap().unwrap();
546 let response2 = rx2.recv().await.unwrap().unwrap();
547 assert!(rx1.try_recv().is_err());
548 assert!(rx2.try_recv().is_err());
549
550 let serving_bitmap = |response: SubscribeResponse| {
551 let Some(Info::TableRefillRuntimeConfig(config)) = response.info else {
552 panic!("expect table refill runtime config");
553 };
554 let mappings = config.serving_table_vnode_mappings.unwrap().mappings;
555 assert_eq!(mappings.len(), 1);
556 assert_eq!(mappings[0].table_id, table_id.as_raw_id());
557 Bitmap::from(mappings[0].bitmap.clone().unwrap())
558 };
559 let bitmap1 = serving_bitmap(response1);
560 let bitmap2 = serving_bitmap(response2);
561 assert_eq!(bitmap1, expected_bitmaps[&worker1.id]);
562 assert_eq!(bitmap2, expected_bitmaps[&worker2.id]);
563 assert_eq!(
564 bitmap1.count_ones() + bitmap2.count_ones(),
565 VirtualNode::COUNT_FOR_TEST
566 );
567 }
568
569 fn hash_serving_info() -> FragmentServingInfo {
570 FragmentServingInfo {
571 result_table_id: None,
572 distribution_type: FragmentDistributionType::Hash,
573 vnode_count: VirtualNode::COUNT_FOR_TEST,
574 }
575 }
576
577 #[test]
578 fn test_reconcile_exactly_matches_full_snapshot_and_removes_failed_placements() {
579 let mapping = ServingVnodeMapping::default();
580 let worker = serving_worker(1);
581 let stale_fragment = FragmentId::new(1);
582 let retained_fragment = FragmentId::new(2);
583 let added_fragment = FragmentId::new(3);
584
585 let initial_snapshot = HashMap::from([
586 (stale_fragment, hash_serving_info()),
587 (retained_fragment, hash_serving_info()),
588 ]);
589 mapping.upsert(&initial_snapshot, std::slice::from_ref(&worker), None);
590 assert!(mapping.all().contains_key(&stale_fragment));
591 let retained_mapping = mapping.all()[&retained_fragment].clone();
592
593 let current_snapshot = HashMap::from([
594 (retained_fragment, hash_serving_info()),
595 (added_fragment, hash_serving_info()),
596 ]);
597 let (reconciled, failed) = mapping.reconcile(¤t_snapshot, &[worker], None);
598 assert!(failed.is_empty());
599 assert_eq!(reconciled.len(), 2);
600 assert!(reconciled.contains_key(&retained_fragment));
601 assert!(reconciled.contains_key(&added_fragment));
602
603 let mappings = mapping.all();
604 assert_eq!(mappings.len(), 2);
605 assert!(mappings.contains_key(&retained_fragment));
606 assert!(mappings.contains_key(&added_fragment));
607 assert_eq!(mappings[&retained_fragment], retained_mapping);
608
609 let final_snapshot = HashMap::from([(retained_fragment, hash_serving_info())]);
610 let (reconciled, failed) = mapping.reconcile(&final_snapshot, &[], None);
611 assert!(reconciled.is_empty());
612 assert_eq!(failed, HashSet::from([retained_fragment]));
613 assert!(mapping.all().is_empty());
614 }
615
616 #[test]
617 fn test_upsert_keeps_fragments_absent_from_incremental_update() {
618 let mapping = ServingVnodeMapping::default();
619 let worker = serving_worker(1);
620 let first_fragment = FragmentId::new(1);
621 let second_fragment = FragmentId::new(2);
622
623 let first_update = HashMap::from([(first_fragment, hash_serving_info())]);
624 mapping.upsert(&first_update, std::slice::from_ref(&worker), None);
625 let second_update = HashMap::from([(second_fragment, hash_serving_info())]);
626 mapping.upsert(&second_update, &[worker], None);
627
628 let mappings = mapping.all();
629 assert!(mappings.contains_key(&first_fragment));
630 assert!(mappings.contains_key(&second_fragment));
631 }
632}