risingwave_meta/
table_refill.rs1use std::collections::HashMap;
16
17use risingwave_meta_model::{FragmentId, WorkerId};
18use risingwave_pb::meta::{PbTableCacheRefillPolicies, PbTableRefillRuntimeConfig};
19
20use crate::controller::fragment::FragmentServingInfo;
21use crate::manager::NotificationVersion;
22use crate::serving::{ServingVnodeMappingRef, to_pb_serving_table_vnode_mappings};
23
24pub fn build_hummock_table_refill_runtime_config(
25 serving_vnode_mapping: &ServingVnodeMappingRef,
26 worker_id: WorkerId,
27 table_cache_refill_policies: PbTableCacheRefillPolicies,
28 fragment_serving_infos: &HashMap<FragmentId, FragmentServingInfo>,
29 version: NotificationVersion,
30) -> PbTableRefillRuntimeConfig {
31 let table_vnode_mapping = serving_vnode_mapping
32 .table_vnode_mappings_by_worker([worker_id], fragment_serving_infos)
33 .remove(&worker_id)
34 .expect("requested worker must have a table vnode mapping");
35
36 PbTableRefillRuntimeConfig {
37 table_cache_refill_policies: Some(table_cache_refill_policies),
38 serving_table_vnode_mappings: Some(to_pb_serving_table_vnode_mappings(
39 &table_vnode_mapping,
40 )),
41 version,
42 }
43}
44
45#[cfg(test)]
46mod tests {
47 use std::sync::Arc;
48
49 use risingwave_common::bitmap::Bitmap;
50 use risingwave_common::hash::{VirtualNode, WorkerSlotId};
51 use risingwave_meta_model::TableId;
52 use risingwave_pb::common::{WorkerNode, WorkerType, worker_node};
53 use risingwave_pb::meta::table_cache_refill_policies::PbTableCacheRefillPolicy;
54 use risingwave_pb::meta::table_cache_refill_policies::table_cache_refill_policy::PbCacheRefillPolicy;
55 use risingwave_pb::meta::table_fragments::fragment::FragmentDistributionType;
56
57 use super::*;
58 use crate::serving::ServingVnodeMapping;
59
60 fn serving_worker(id: u32) -> WorkerNode {
61 WorkerNode {
62 id: id.into(),
63 r#type: WorkerType::ComputeNode as i32,
64 property: Some(worker_node::Property {
65 is_serving: true,
66 parallelism: 1,
67 ..Default::default()
68 }),
69 ..Default::default()
70 }
71 }
72
73 #[test]
74 fn test_build_hummock_table_refill_runtime_config_snapshot() {
75 let worker1 = serving_worker(1);
76 let worker2 = serving_worker(2);
77 let fragment_id = FragmentId::new(233);
78 let result_table_id = TableId::new(234);
79 let internal_table_id = TableId::new(235);
80 let fragment_serving_infos = HashMap::from([(
81 fragment_id,
82 FragmentServingInfo {
83 result_table_id: Some(result_table_id),
84 distribution_type: FragmentDistributionType::Hash,
85 vnode_count: VirtualNode::COUNT_FOR_TEST,
86 },
87 )]);
88 let serving_vnode_mapping = Arc::new(ServingVnodeMapping::default());
89 let (fragment_mappings, failed) = serving_vnode_mapping.upsert(
90 &fragment_serving_infos,
91 &[worker1.clone(), worker2],
92 None,
93 );
94 assert!(failed.is_empty());
95 let expected_bitmap =
96 fragment_mappings[&fragment_id].to_bitmaps()[&WorkerSlotId::new(worker1.id, 0)].clone();
97
98 let policies = PbTableCacheRefillPolicies {
99 table_policies: vec![PbTableCacheRefillPolicy {
100 table_id: result_table_id.as_raw_id(),
101 policy: PbCacheRefillPolicy::Serving as i32,
102 }],
103 internal_table_policies: vec![PbTableCacheRefillPolicy {
104 table_id: internal_table_id.as_raw_id(),
105 policy: PbCacheRefillPolicy::Streaming as i32,
106 }],
107 };
108 let config = build_hummock_table_refill_runtime_config(
109 &serving_vnode_mapping,
110 worker1.id,
111 policies.clone(),
112 &fragment_serving_infos,
113 42,
114 );
115
116 assert_eq!(config.version, 42);
117 assert_eq!(config.table_cache_refill_policies, Some(policies));
118 let mappings = config.serving_table_vnode_mappings.unwrap().mappings;
119 assert_eq!(mappings.len(), 1);
120 assert_eq!(mappings[0].table_id, result_table_id.as_raw_id());
121 assert_eq!(
122 Bitmap::from(mappings[0].bitmap.clone().unwrap()),
123 expected_bitmap
124 );
125 }
126}