1use std::collections::hash_map::Entry;
16use std::collections::{HashMap, HashSet};
17use std::error::Error;
18use std::fmt::{Debug, Formatter};
19use std::future::poll_fn;
20use std::sync::Arc;
21use std::task::{Context, Poll};
22use std::time::Duration;
23
24use anyhow::anyhow;
25use fail::fail_point;
26use futures::future::{BoxFuture, join_all};
27use futures::{FutureExt, StreamExt};
28use itertools::Itertools;
29use risingwave_common::bail;
30use risingwave_common::catalog::{DatabaseId, FragmentTypeFlag, TableId};
31use risingwave_common::id::JobId;
32use risingwave_common::util::epoch::Epoch;
33use risingwave_common::util::retry::exponential_backoff;
34use risingwave_common::util::stream_graph_visitor::visit_stream_node_cont;
35use risingwave_common::util::tracing::TracingContext;
36use risingwave_connector::source::SplitImpl;
37use risingwave_meta_model::WorkerId;
38use risingwave_pb::common::{HostAddress, WorkerNode};
39use risingwave_pb::hummock::HummockVersionStats;
40use risingwave_pb::id::PartialGraphId;
41use risingwave_pb::source::{PbCdcTableSnapshotSplits, PbCdcTableSnapshotSplitsWithGeneration};
42use risingwave_pb::stream_plan::barrier_mutation::Mutation;
43use risingwave_pb::stream_plan::stream_node::NodeBody;
44use risingwave_pb::stream_plan::{
45 AddMutation, Barrier, BarrierMutation, IcebergPkIndexCompactionContext,
46};
47use risingwave_pb::stream_service::inject_barrier_request::build_actor_info::UpstreamActors;
48use risingwave_pb::stream_service::inject_barrier_request::{
49 BuildActorInfo, FragmentBuildActorInfo,
50};
51use risingwave_pb::stream_service::streaming_control_stream_request::{
52 CreatePartialGraphRequest, PbCreatePartialGraphRequest, PbInitRequest,
53 RemovePartialGraphRequest, ResetPartialGraphsRequest,
54};
55use risingwave_pb::stream_service::{
56 InjectBarrierRequest, StreamingControlStreamRequest, streaming_control_stream_request,
57 streaming_control_stream_response,
58};
59use risingwave_rpc_client::StreamingControlHandle;
60use thiserror_ext::AsReport;
61use tokio::time::{Instant, sleep};
62use tracing::{debug, error, info, warn};
63use uuid::Uuid;
64
65use super::{BarrierKind, TracedEpoch};
66use crate::barrier::BackfillOrderState;
67use crate::barrier::backfill_order_control::get_nodes_with_backfill_dependencies;
68use crate::barrier::cdc_progress::CdcTableBackfillTracker;
69use crate::barrier::checkpoint::{
70 BarrierWorkerState, BatchRefreshJobCheckpointControl, BatchRefreshRenderResult,
71 CreatingStreamingJobControl, DatabaseCheckpointControl, DatabaseCheckpointControlMetrics,
72 IndependentCheckpointJobControl,
73};
74use crate::barrier::context::{GlobalBarrierWorkerContext, GlobalBarrierWorkerContextImpl};
75use crate::barrier::edge_builder::{EdgeBuilderFragmentInfo, FragmentEdgeBuilder};
76use crate::barrier::info::{
77 BarrierInfo, CreateStreamingJobStatus, InflightDatabaseInfo, InflightStreamingJobInfo,
78 SubscriberType,
79};
80use crate::barrier::partial_graph::PartialGraphRecoverer;
81use crate::barrier::progress::CreateMviewProgressTracker;
82use crate::barrier::utils::NodeToCollect;
83use crate::controller::fragment::InflightFragmentInfo;
84use crate::controller::utils::StreamingJobExtraInfo;
85use crate::manager::MetaSrvEnv;
86use crate::model::{
87 ActorId, FragmentDownstreamRelation, FragmentId, StreamActor, StreamJobActorsToCreate,
88 SubscriptionId,
89};
90use crate::stream::cdc::{
91 CdcTableSnapshotSplits, is_parallelized_backfill_enabled_cdc_scan_fragment,
92};
93use crate::stream::{
94 ExtendedFragmentBackfillOrder, StreamFragmentGraph, UserDefinedFragmentBackfillOrder,
95 build_actor_connector_splits,
96};
97use crate::{MetaError, MetaResult};
98
99pub(crate) fn to_partial_graph_id(
100 database_id: DatabaseId,
101 creating_job_id: Option<JobId>,
102) -> PartialGraphId {
103 let raw_job_id = creating_job_id
104 .map(|job_id| {
105 assert_ne!(job_id, u32::MAX);
106 job_id.as_raw_id()
107 })
108 .unwrap_or(u32::MAX);
109 (((database_id.as_raw_id() as u64) << 32) | (raw_job_id as u64)).into()
110}
111
112pub(super) fn from_partial_graph_id(
113 partial_graph_id: PartialGraphId,
114) -> (DatabaseId, Option<JobId>) {
115 let id = partial_graph_id.as_raw_id();
116 let database_id = (id >> 32) as u32;
117 let raw_creating_job_id = (id & ((1 << 32) - 1)) as u32;
118 let creating_job_id = if raw_creating_job_id == u32::MAX {
119 None
120 } else {
121 Some(JobId::new(raw_creating_job_id))
122 };
123 (database_id.into(), creating_job_id)
124}
125
126pub(super) fn build_locality_fragment_state_table_mapping(
127 fragment_infos: &HashMap<FragmentId, InflightFragmentInfo>,
128) -> HashMap<FragmentId, Vec<TableId>> {
129 let mut mapping = HashMap::new();
130
131 for (fragment_id, fragment_info) in fragment_infos {
132 let mut state_table_ids = Vec::new();
133 visit_stream_node_cont(&fragment_info.nodes, |stream_node| {
134 if let Some(NodeBody::LocalityProvider(locality_provider)) =
135 stream_node.node_body.as_ref()
136 {
137 let state_table_id = locality_provider
138 .state_table
139 .as_ref()
140 .expect("must have state table")
141 .id;
142 state_table_ids.push(state_table_id);
143 false
144 } else {
145 true
146 }
147 });
148 if !state_table_ids.is_empty() {
149 mapping.insert(*fragment_id, state_table_ids);
150 }
151 }
152
153 mapping
154}
155
156pub(super) fn database_partial_graphs<'a>(
157 database_id: DatabaseId,
158 creating_jobs: impl Iterator<Item = JobId> + Sized + 'a,
159) -> impl Iterator<Item = PartialGraphId> + 'a {
160 creating_jobs
161 .map(Some)
162 .chain([None])
163 .map(move |creating_job_id| to_partial_graph_id(database_id, creating_job_id))
164}
165
166struct ControlStreamNode {
167 worker_id: WorkerId,
168 host: HostAddress,
169 handle: StreamingControlHandle,
170}
171
172enum WorkerNodeState {
173 Connected {
174 control_stream: ControlStreamNode,
175 removed: bool,
176 },
177 Reconnecting(BoxFuture<'static, StreamingControlHandle>),
178}
179
180pub(super) struct ControlStreamManager {
181 workers: HashMap<WorkerId, (WorkerNode, WorkerNodeState)>,
182 pub env: MetaSrvEnv,
183}
184
185impl ControlStreamManager {
186 pub(super) fn new(env: MetaSrvEnv) -> Self {
187 Self {
188 workers: Default::default(),
189 env,
190 }
191 }
192
193 pub(super) fn host_addr(&self, worker_id: WorkerId) -> HostAddress {
194 self.workers[&worker_id].0.host.clone().unwrap()
195 }
196
197 pub(super) async fn add_worker(
198 &mut self,
199 node: WorkerNode,
200 partial_graphs: impl Iterator<Item = PartialGraphId>,
201 term_id: &String,
202 context: Arc<impl GlobalBarrierWorkerContext>,
203 ) {
204 let node_id = node.id;
205 if let Entry::Occupied(entry) = self.workers.entry(node_id) {
206 let (existing_node, worker_state) = entry.get();
207 assert_eq!(existing_node.host, node.host);
208 warn!(id = %node.id, host = ?node.host, "node already exists");
209 match worker_state {
210 WorkerNodeState::Connected { .. } => {
211 warn!(id = %node.id, host = ?node.host, "new node already connected");
212 return;
213 }
214 WorkerNodeState::Reconnecting(_) => {
215 warn!(id = %node.id, host = ?node.host, "remove previous pending worker connect request and reconnect");
216 entry.remove();
217 }
218 }
219 }
220 let node_host = node.host.clone().unwrap();
221 let mut backoff =
222 exponential_backoff(Duration::from_millis(100), 5, Duration::from_secs(3));
223 const MAX_RETRY: usize = 5;
224 for i in 1..=MAX_RETRY {
225 match context
226 .new_control_stream(
227 &node,
228 &PbInitRequest {
229 term_id: term_id.clone(),
230 },
231 )
232 .await
233 {
234 Ok(mut handle) => {
235 WorkerNodeConnected {
236 handle: &mut handle,
237 node: &node,
238 }
239 .initialize(partial_graphs);
240 info!(?node_host, "add control stream worker");
241 assert!(
242 self.workers
243 .insert(
244 node_id,
245 (
246 node,
247 WorkerNodeState::Connected {
248 control_stream: ControlStreamNode {
249 worker_id: node_id as _,
250 host: node_host,
251 handle,
252 },
253 removed: false
254 }
255 )
256 )
257 .is_none()
258 );
259 return;
260 }
261 Err(e) => {
262 let delay = backoff.next().unwrap();
265 error!(
266 attempt = i,
267 backoff_delay = ?delay,
268 err = %e.as_report(),
269 ?node_host,
270 "failed to resolve the worker node address",
271 );
272 sleep(delay).await;
273 }
274 }
275 }
276 error!(?node_host, "failed to create the worker node after retries");
277 assert!(
278 self.workers
279 .insert(
280 node_id,
281 (
282 node.clone(),
283 WorkerNodeState::Reconnecting(ControlStreamManager::retry_connect(
284 node,
285 term_id.to_owned(),
286 context,
287 ))
288 )
289 )
290 .is_none()
291 );
292 }
293
294 pub(super) fn remove_worker(&mut self, node: WorkerNode) {
295 if let Entry::Occupied(mut entry) = self.workers.entry(node.id) {
296 let (_, worker_state) = entry.get_mut();
297 match worker_state {
298 WorkerNodeState::Connected { removed, .. } => {
299 info!(worker_id = %node.id, "mark connected worker as removed");
300 *removed = true;
301 }
302 WorkerNodeState::Reconnecting(_) => {
303 info!(worker_id = %node.id, "remove worker");
304 entry.remove();
305 }
306 }
307 }
308 }
309
310 fn retry_connect(
311 node: WorkerNode,
312 term_id: String,
313 context: Arc<impl GlobalBarrierWorkerContext>,
314 ) -> BoxFuture<'static, StreamingControlHandle> {
315 async move {
316 let mut attempt = 0;
317 let backoff = exponential_backoff(
318 Duration::from_millis(100),
319 5,
320 Duration::from_mins(1),
321 );
322 let init_request = PbInitRequest { term_id };
323 for delay in backoff {
324 attempt += 1;
325 sleep(delay).await;
326 match context.new_control_stream(&node, &init_request).await {
327 Ok(handle) => {
328 return handle;
329 }
330 Err(e) => {
331 warn!(e = %e.as_report(), ?node, attempt, "failed to create the control stream worker");
332 }
333 }
334 }
335 unreachable!("end of retry backoff")
336 }.boxed()
337 }
338
339 pub(super) async fn recover(
340 env: MetaSrvEnv,
341 nodes: &HashMap<WorkerId, WorkerNode>,
342 term_id: &str,
343 context: Arc<impl GlobalBarrierWorkerContext>,
344 ) -> Self {
345 let reset_start_time = Instant::now();
346 let init_request = PbInitRequest {
347 term_id: term_id.to_owned(),
348 };
349 let init_request = &init_request;
350 let nodes = join_all(nodes.iter().map(|(worker_id, node)| async {
351 let result = context.new_control_stream(node, init_request).await;
352 (*worker_id, node.clone(), result)
353 }))
354 .await;
355 let mut unconnected_workers = HashSet::new();
356 let mut workers = HashMap::new();
357 for (worker_id, node, result) in nodes {
358 match result {
359 Ok(handle) => {
360 let control_stream = ControlStreamNode {
361 worker_id: node.id,
362 host: node.host.clone().unwrap(),
363 handle,
364 };
365 assert!(
366 workers
367 .insert(
368 worker_id,
369 (
370 node,
371 WorkerNodeState::Connected {
372 control_stream,
373 removed: false
374 }
375 )
376 )
377 .is_none()
378 );
379 }
380 Err(e) => {
381 unconnected_workers.insert(worker_id);
382 warn!(
383 e = %e.as_report(),
384 %worker_id,
385 ?node,
386 "failed to connect to node"
387 );
388 assert!(
389 workers
390 .insert(
391 worker_id,
392 (
393 node.clone(),
394 WorkerNodeState::Reconnecting(Self::retry_connect(
395 node,
396 term_id.to_owned(),
397 context.clone()
398 ))
399 )
400 )
401 .is_none()
402 );
403 }
404 }
405 }
406
407 info!(elapsed=?reset_start_time.elapsed(), ?unconnected_workers, "control stream reset");
408
409 Self { workers, env }
410 }
411
412 pub(super) fn clear(&mut self) {
414 *self = Self::new(self.env.clone());
415 }
416}
417
418pub(super) struct WorkerNodeConnected<'a> {
419 node: &'a WorkerNode,
420 handle: &'a mut StreamingControlHandle,
421}
422
423impl<'a> WorkerNodeConnected<'a> {
424 pub(super) fn initialize(self, partial_graphs: impl Iterator<Item = PartialGraphId>) {
425 for partial_graph_id in partial_graphs {
426 if let Err(e) = self.handle.send_request(StreamingControlStreamRequest {
427 request: Some(
428 streaming_control_stream_request::Request::CreatePartialGraph(
429 PbCreatePartialGraphRequest { partial_graph_id },
430 ),
431 ),
432 }) {
433 warn!(e = %e.as_report(), node = ?self.node, "failed to send initial partial graph request");
434 }
435 }
436 }
437}
438
439pub(super) enum WorkerNodeEvent<'a> {
440 Response(MetaResult<streaming_control_stream_response::Response>),
441 Connected(WorkerNodeConnected<'a>),
442}
443
444impl ControlStreamManager {
445 fn poll_next_event<'a>(
446 this_opt: &mut Option<&'a mut Self>,
447 cx: &mut Context<'_>,
448 term_id: &str,
449 context: &Arc<impl GlobalBarrierWorkerContext>,
450 poll_reconnect: bool,
451 ) -> Poll<(WorkerId, WorkerNodeEvent<'a>)> {
452 let this = this_opt.as_mut().expect("Future polled after completion");
453 if this.workers.is_empty() {
454 return Poll::Pending;
455 }
456 {
457 for (&worker_id, (node, worker_state)) in &mut this.workers {
458 let control_stream = match worker_state {
459 WorkerNodeState::Connected { control_stream, .. } => control_stream,
460 WorkerNodeState::Reconnecting(_) if !poll_reconnect => {
461 continue;
462 }
463 WorkerNodeState::Reconnecting(join_handle) => {
464 match join_handle.poll_unpin(cx) {
465 Poll::Ready(handle) => {
466 info!(id=%node.id, host=?node.host, "reconnected to worker");
467 *worker_state = WorkerNodeState::Connected {
468 control_stream: ControlStreamNode {
469 worker_id: node.id,
470 host: node.host.clone().unwrap(),
471 handle,
472 },
473 removed: false,
474 };
475 let this = this_opt.take().expect("should exist");
476 let (node, worker_state) =
477 this.workers.get_mut(&worker_id).expect("should exist");
478 let WorkerNodeState::Connected { control_stream, .. } =
479 worker_state
480 else {
481 unreachable!()
482 };
483 return Poll::Ready((
484 worker_id,
485 WorkerNodeEvent::Connected(WorkerNodeConnected {
486 node,
487 handle: &mut control_stream.handle,
488 }),
489 ));
490 }
491 Poll::Pending => {
492 continue;
493 }
494 }
495 }
496 };
497 match control_stream.handle.response_stream.poll_next_unpin(cx) {
498 Poll::Ready(result) => {
499 {
500 let result = result
501 .ok_or_else(|| (false, anyhow!("end of stream").into()))
502 .and_then(|result| {
503 result.map_err(|err| -> (bool, MetaError) { (false, err.into()) }).and_then(|resp| {
504 match resp
505 .response
506 .ok_or_else(|| (false, anyhow!("empty response").into()))?
507 {
508 streaming_control_stream_response::Response::Shutdown(_) => Err((true, anyhow!(
509 "worker node {worker_id} is shutting down"
510 )
511 .into())),
512 streaming_control_stream_response::Response::Init(_) => {
513 Err((false, anyhow!("get unexpected init response").into()))
515 }
516 resp => {
517 if let streaming_control_stream_response::Response::CompleteBarrier(barrier_resp) = &resp {
518 assert_eq!(worker_id, barrier_resp.worker_id);
519 }
520 Ok(resp)
521 }
522 }
523 })
524 });
525 let result = match result {
526 Ok(resp) => Ok(resp),
527 Err((shutdown, err)) => {
528 warn!(worker_id = %node.id, host = ?node.host, err = %err.as_report(), "get error from response stream");
529 let WorkerNodeState::Connected { removed, .. } = worker_state
530 else {
531 unreachable!("checked connected")
532 };
533 if *removed || shutdown {
534 this.workers.remove(&worker_id);
535 } else {
536 *worker_state = WorkerNodeState::Reconnecting(
537 ControlStreamManager::retry_connect(
538 node.clone(),
539 term_id.to_owned(),
540 context.clone(),
541 ),
542 );
543 }
544 Err(err)
545 }
546 };
547 return Poll::Ready((worker_id, WorkerNodeEvent::Response(result)));
548 }
549 }
550 Poll::Pending => {
551 continue;
552 }
553 }
554 }
555 };
556
557 Poll::Pending
558 }
559
560 #[await_tree::instrument("control_stream_next_event")]
561 pub(super) async fn next_event<'a>(
562 &'a mut self,
563 term_id: &str,
564 context: &Arc<impl GlobalBarrierWorkerContext>,
565 ) -> (WorkerId, WorkerNodeEvent<'a>) {
566 let mut this = Some(self);
567 poll_fn(|cx| Self::poll_next_event(&mut this, cx, term_id, context, true)).await
568 }
569
570 #[await_tree::instrument("control_stream_next_response")]
571 pub(super) async fn next_response(
572 &mut self,
573 term_id: &str,
574 context: &Arc<impl GlobalBarrierWorkerContext>,
575 ) -> (
576 WorkerId,
577 MetaResult<streaming_control_stream_response::Response>,
578 ) {
579 let mut this = Some(self);
580 let (worker_id, event) =
581 poll_fn(|cx| Self::poll_next_event(&mut this, cx, term_id, context, false)).await;
582 match event {
583 WorkerNodeEvent::Response(result) => (worker_id, result),
584 WorkerNodeEvent::Connected(_) => {
585 unreachable!("set poll_reconnect=false")
586 }
587 }
588 }
589}
590
591pub(super) struct DatabaseInitialBarrierCollector {
592 pub(super) database_id: DatabaseId,
593 pub(super) initializing_partial_graphs: HashSet<PartialGraphId>,
594 pub(super) database: DatabaseCheckpointControl,
595}
596
597impl Debug for DatabaseInitialBarrierCollector {
598 fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
599 f.debug_struct("DatabaseInitialBarrierCollector")
600 .field("database_id", &self.database_id)
601 .field("initializing_graphs", &self.initializing_partial_graphs)
602 .finish()
603 }
604}
605
606impl DatabaseInitialBarrierCollector {
607 pub(super) fn is_collected(&self) -> bool {
608 self.initializing_partial_graphs.is_empty()
609 }
610
611 pub(super) fn partial_graph_initialized(&mut self, partial_graph_id: PartialGraphId) {
612 assert!(self.initializing_partial_graphs.remove(&partial_graph_id));
613 }
614
615 pub(super) fn all_partial_graphs(&self) -> impl Iterator<Item = PartialGraphId> + '_ {
616 database_partial_graphs(
617 self.database_id,
618 self.database
619 .independent_checkpoint_job_controls
620 .keys()
621 .copied(),
622 )
623 }
624
625 pub(super) fn finish(self) -> DatabaseCheckpointControl {
626 assert!(self.is_collected());
627 self.database
628 }
629
630 pub(super) fn is_valid_after_worker_err(&self, worker_id: WorkerId) -> bool {
631 self.database.is_valid_after_worker_err(worker_id)
632 }
633}
634
635impl PartialGraphRecoverer<'_> {
636 #[expect(clippy::too_many_arguments)]
638 pub(super) fn inject_database_initial_barrier(
639 &mut self,
640 database_id: DatabaseId,
641 jobs: HashMap<JobId, HashMap<FragmentId, InflightFragmentInfo>>,
642 job_extra_info: &HashMap<JobId, StreamingJobExtraInfo>,
643 state_table_committed_epochs: &mut HashMap<TableId, u64>,
644 state_table_log_epochs: &mut HashMap<TableId, Vec<(Vec<u64>, u64)>>,
645 fragment_relations: &FragmentDownstreamRelation,
646 stream_actors: &HashMap<ActorId, StreamActor>,
647 source_splits: &mut HashMap<ActorId, Vec<SplitImpl>>,
648 creating_jobs: &mut HashSet<JobId>,
649 mv_depended_subscriptions: &mut HashMap<TableId, HashMap<SubscriptionId, u64>>,
650 is_paused: bool,
651 hummock_version_stats: &HummockVersionStats,
652 cdc_table_snapshot_splits: &mut HashMap<JobId, CdcTableSnapshotSplits>,
653 batch_refresh: HashMap<JobId, BatchRefreshRenderResult>,
654 ) -> MetaResult<DatabaseCheckpointControl> {
655 fn collect_source_splits(
656 fragment_infos: impl Iterator<Item = &InflightFragmentInfo>,
657 source_splits: &mut HashMap<ActorId, Vec<SplitImpl>>,
658 ) -> HashMap<ActorId, Vec<SplitImpl>> {
659 fragment_infos
660 .flat_map(|info| info.actors.keys())
661 .filter_map(|actor_id| {
662 let actor_id = *actor_id as ActorId;
663 source_splits
664 .remove(&actor_id)
665 .map(|splits| (actor_id, splits))
666 })
667 .collect()
668 }
669 fn build_mutation(
670 splits: &HashMap<ActorId, Vec<SplitImpl>>,
671 cdc_table_snapshot_split_assignment: HashMap<ActorId, PbCdcTableSnapshotSplits>,
672 backfill_orders: &ExtendedFragmentBackfillOrder,
673 is_paused: bool,
674 ) -> Mutation {
675 let backfill_nodes_to_pause = get_nodes_with_backfill_dependencies(backfill_orders)
676 .into_iter()
677 .collect();
678 Mutation::Add(AddMutation {
679 actor_dispatchers: Default::default(),
681 added_actors: Default::default(),
682 actor_splits: build_actor_connector_splits(splits),
683 actor_cdc_table_snapshot_splits: Some(PbCdcTableSnapshotSplitsWithGeneration {
684 splits: cdc_table_snapshot_split_assignment,
685 }),
686 pause: is_paused,
687 subscriptions_to_add: Default::default(),
688 backfill_nodes_to_pause,
689 new_upstream_sinks: Default::default(),
690 dropped_actors: Default::default(),
691 sink_log_store_flush: Default::default(),
692 })
693 }
694
695 fn resolve_jobs_committed_epoch(
696 state_table_committed_epochs: &mut HashMap<TableId, u64>,
697 table_ids: impl Iterator<Item = TableId>,
698 ) -> u64 {
699 let mut epochs = table_ids.map(|table_id| {
700 (
701 table_id,
702 state_table_committed_epochs
703 .remove(&table_id)
704 .expect("should exist"),
705 )
706 });
707 let (first_table_id, prev_epoch) = epochs.next().expect("non-empty");
708 for (table_id, epoch) in epochs {
709 assert_eq!(
710 prev_epoch, epoch,
711 "{} has different committed epoch to {}",
712 first_table_id, table_id
713 );
714 }
715 prev_epoch
716 }
717 fn job_backfill_orders(
718 job_extra_info: &HashMap<JobId, StreamingJobExtraInfo>,
719 job_id: JobId,
720 ) -> UserDefinedFragmentBackfillOrder {
721 UserDefinedFragmentBackfillOrder::new(
722 job_extra_info
723 .get(&job_id)
724 .and_then(|info| info.backfill_orders.clone())
725 .map_or_else(HashMap::new, |orders| orders.0),
726 )
727 }
728
729 let mut subscribers: HashMap<_, HashMap<_, _>> = jobs
730 .keys()
731 .filter_map(|job_id| {
732 mv_depended_subscriptions
733 .remove(&job_id.as_mv_table_id())
734 .map(|subscriptions| {
735 (
736 job_id.as_mv_table_id(),
737 subscriptions
738 .into_iter()
739 .map(|(subscription_id, retention)| {
740 (
741 subscription_id.as_subscriber_id(),
742 SubscriberType::Subscription(retention),
743 )
744 })
745 .collect(),
746 )
747 })
748 })
749 .collect();
750
751 for (job_id, render_result) in &batch_refresh {
754 let snapshot_backfill_info = StreamFragmentGraph::collect_snapshot_backfill_info_impl(
755 render_result
756 .fragment_infos
757 .values()
758 .map(|fragment| (&fragment.nodes, fragment.fragment_type_mask)),
759 )?
760 .0
761 .ok_or_else(|| anyhow!("batch refresh job {} has no snapshot backfill info", job_id))?;
762
763 for upstream_table_id in snapshot_backfill_info
764 .upstream_mv_table_id_to_backfill_epoch
765 .keys()
766 {
767 subscribers
768 .entry(*upstream_table_id)
769 .or_default()
770 .try_insert(job_id.as_subscriber_id(), SubscriberType::SnapshotBackfill)
771 .expect("non-duplicate");
772 }
773 }
774
775 let mut database_jobs = HashMap::new();
776 let mut snapshot_backfill_jobs = HashMap::new();
777
778 for (job_id, job_fragments) in jobs {
779 if creating_jobs.remove(&job_id) {
780 if job_fragments.values().any(|fragment| {
781 fragment
782 .fragment_type_mask
783 .contains(FragmentTypeFlag::SnapshotBackfillStreamScan)
784 }) {
785 debug!(%job_id, "recovered snapshot backfill job");
786 snapshot_backfill_jobs.insert(job_id, job_fragments);
787 } else {
788 database_jobs.insert(job_id, (job_fragments, true));
789 }
790 } else {
791 database_jobs.insert(job_id, (job_fragments, false));
792 }
793 }
794
795 let database_job_log_epochs: HashMap<_, _> = database_jobs
796 .keys()
797 .filter_map(|job_id| {
798 state_table_log_epochs
799 .remove(&job_id.as_mv_table_id())
800 .map(|epochs| (job_id.as_mv_table_id(), epochs))
801 })
802 .collect();
803
804 let prev_epoch = resolve_jobs_committed_epoch(
805 state_table_committed_epochs,
806 InflightFragmentInfo::existing_table_ids(
807 database_jobs.values().flat_map(|(job, _)| job.values()),
808 ),
809 );
810 let prev_epoch = TracedEpoch::new(Epoch(prev_epoch));
811 let curr_epoch = prev_epoch.next();
813 let barrier_info = BarrierInfo {
814 prev_epoch,
815 curr_epoch,
816 kind: BarrierKind::Initial,
817 };
818
819 let mut ongoing_snapshot_backfill_jobs: HashMap<JobId, _> = HashMap::new();
820 for (job_id, fragment_infos) in snapshot_backfill_jobs {
821 let committed_epoch = resolve_jobs_committed_epoch(
822 state_table_committed_epochs,
823 InflightFragmentInfo::existing_table_ids(fragment_infos.values()),
824 );
825 if committed_epoch == barrier_info.prev_epoch() {
826 info!(
827 "recovered creating snapshot backfill job {} catch up with upstream already",
828 job_id
829 );
830 database_jobs
831 .try_insert(job_id, (fragment_infos, true))
832 .expect("non-duplicate");
833 continue;
834 }
835 let snapshot_backfill_info = StreamFragmentGraph::collect_snapshot_backfill_info_impl(
836 fragment_infos
837 .values()
838 .map(|fragment| (&fragment.nodes, fragment.fragment_type_mask)),
839 )?
840 .0
841 .ok_or_else(|| {
842 anyhow!(
843 "recovered snapshot backfill job {} has no snapshot backfill info",
844 job_id
845 )
846 })?;
847 let mut snapshot_epoch = None;
848 let upstream_table_ids: HashSet<_> = snapshot_backfill_info
849 .upstream_mv_table_id_to_backfill_epoch
850 .keys()
851 .cloned()
852 .collect();
853 for (upstream_table_id, epoch) in
854 snapshot_backfill_info.upstream_mv_table_id_to_backfill_epoch
855 {
856 let epoch = epoch.ok_or_else(|| anyhow!("recovered snapshot backfill job {} to upstream {} has not set snapshot epoch", job_id, upstream_table_id))?;
857 let snapshot_epoch = snapshot_epoch.get_or_insert(epoch);
858 if *snapshot_epoch != epoch {
859 return Err(anyhow!("snapshot epoch {} to upstream {} different to snapshot epoch {} to previous upstream", epoch, upstream_table_id, snapshot_epoch).into());
860 }
861 }
862 let snapshot_epoch = snapshot_epoch.ok_or_else(|| {
863 anyhow!(
864 "snapshot backfill job {} has not set snapshot epoch",
865 job_id
866 )
867 })?;
868 for upstream_table_id in &upstream_table_ids {
869 subscribers
870 .entry(*upstream_table_id)
871 .or_default()
872 .try_insert(job_id.as_subscriber_id(), SubscriberType::SnapshotBackfill)
873 .expect("non-duplicate");
874 }
875 ongoing_snapshot_backfill_jobs
876 .try_insert(
877 job_id,
878 (
879 fragment_infos,
880 upstream_table_ids,
881 committed_epoch,
882 snapshot_epoch,
883 ),
884 )
885 .expect("non-duplicated");
886 }
887
888 let mut cdc_table_snapshot_split_assignment: HashMap<ActorId, PbCdcTableSnapshotSplits> =
889 HashMap::new();
890
891 let database_jobs: HashMap<JobId, InflightStreamingJobInfo> = {
892 database_jobs
893 .into_iter()
894 .map(|(job_id, (fragment_infos, is_creating))| {
895 let status = if is_creating {
896 let backfill_ordering = job_backfill_orders(job_extra_info, job_id);
897 let backfill_ordering = StreamFragmentGraph::extend_fragment_backfill_ordering_with_locality_backfill(
898 backfill_ordering,
899 fragment_relations,
900 || fragment_infos.iter().map(|(fragment_id, fragment)| {
901 (*fragment_id, fragment.fragment_type_mask, &fragment.nodes)
902 }));
903 let locality_fragment_state_table_mapping =
904 build_locality_fragment_state_table_mapping(&fragment_infos);
905 let backfill_order_state = BackfillOrderState::recover_from_fragment_infos(
906 &backfill_ordering,
907 &fragment_infos,
908 locality_fragment_state_table_mapping,
909 );
910 CreateStreamingJobStatus::Creating {
911 tracker: CreateMviewProgressTracker::recover(
912 job_id,
913 &fragment_infos,
914 backfill_order_state,
915 hummock_version_stats,
916 ),
917 }
918 } else {
919 CreateStreamingJobStatus::Created
920 };
921 let cdc_table_backfill_tracker =
922 if let Some(splits) = cdc_table_snapshot_splits.remove(&job_id) {
923 let cdc_fragment = fragment_infos
924 .values()
925 .find(|fragment| {
926 is_parallelized_backfill_enabled_cdc_scan_fragment(
927 fragment.fragment_type_mask,
928 &fragment.nodes,
929 )
930 .is_some()
931 })
932 .expect("should have parallel cdc fragment");
933 let cdc_actors = cdc_fragment.actors.keys().copied().collect();
934 let mut tracker =
935 CdcTableBackfillTracker::restore(cdc_fragment.fragment_id, splits);
936 cdc_table_snapshot_split_assignment
937 .extend(tracker.reassign_splits(cdc_actors)?);
938 Some(tracker)
939 } else {
940 None
941 };
942 Ok((
943 job_id,
944 InflightStreamingJobInfo {
945 job_id,
946 fragment_infos,
947 subscribers: subscribers
948 .remove(&job_id.as_mv_table_id())
949 .unwrap_or_default(),
950 status,
951 cdc_table_backfill_tracker,
952 },
953 ))
954 })
955 .try_collect::<_, _, MetaError>()
956 }?;
957
958 let control_stream_manager = self.control_stream_manager();
959 let mut builder = FragmentEdgeBuilder::new(
960 database_jobs
961 .values()
962 .flat_map(|job| {
963 let partial_graph_id = to_partial_graph_id(database_id, None);
964 job.fragment_infos().map(move |info| {
965 (
966 info.fragment_id,
967 EdgeBuilderFragmentInfo::from_inflight(
968 info,
969 partial_graph_id,
970 control_stream_manager,
971 ),
972 )
973 })
974 })
975 .chain(ongoing_snapshot_backfill_jobs.iter().flat_map(
976 |(job_id, (fragments, ..))| {
977 let partial_graph_id = to_partial_graph_id(database_id, Some(*job_id));
978 fragments.values().map(move |fragment| {
979 (
980 fragment.fragment_id,
981 EdgeBuilderFragmentInfo::from_inflight(
982 fragment,
983 partial_graph_id,
984 control_stream_manager,
985 ),
986 )
987 })
988 },
989 )),
990 );
991 builder.add_relations(fragment_relations);
992 let mut edges = builder.build();
993
994 {
995 let new_actors =
996 edges.collect_actors_to_create(database_jobs.values().flat_map(move |job| {
997 job.fragment_infos.values().map(move |fragment_infos| {
998 (
999 fragment_infos.fragment_id,
1000 &fragment_infos.nodes,
1001 fragment_infos.actors.iter().map(move |(actor_id, actor)| {
1002 (
1003 stream_actors.get(actor_id).expect("should exist"),
1004 actor.worker_id,
1005 )
1006 }),
1007 job.subscribers.keys().copied(),
1008 )
1009 })
1010 }));
1011
1012 let nodes_actors =
1013 InflightFragmentInfo::actor_ids_to_collect(database_jobs.values().flatten());
1014 let database_job_source_splits =
1015 collect_source_splits(database_jobs.values().flatten(), source_splits);
1016 let database_backfill_orders =
1017 UserDefinedFragmentBackfillOrder::merge(database_jobs.values().map(|job| {
1018 if matches!(job.status, CreateStreamingJobStatus::Creating { .. }) {
1019 job_backfill_orders(job_extra_info, job.job_id)
1020 } else {
1021 UserDefinedFragmentBackfillOrder::default()
1022 }
1023 }));
1024 let database_backfill_orders =
1025 StreamFragmentGraph::extend_fragment_backfill_ordering_with_locality_backfill(
1026 database_backfill_orders,
1027 fragment_relations,
1028 || {
1029 database_jobs.values().flat_map(|job_fragments| {
1030 job_fragments
1031 .fragment_infos
1032 .iter()
1033 .map(|(fragment_id, fragment)| {
1034 (*fragment_id, fragment.fragment_type_mask, &fragment.nodes)
1035 })
1036 })
1037 },
1038 );
1039 let mutation = build_mutation(
1040 &database_job_source_splits,
1041 cdc_table_snapshot_split_assignment,
1042 &database_backfill_orders,
1043 is_paused,
1044 );
1045
1046 let partial_graph_id = to_partial_graph_id(database_id, None);
1047 self.recover_graph(
1048 partial_graph_id,
1049 mutation,
1050 &barrier_info,
1051 &nodes_actors,
1052 InflightFragmentInfo::existing_table_ids(database_jobs.values().flatten()),
1053 new_actors,
1054 DatabaseCheckpointControlMetrics::new(database_id),
1055 )?;
1056 debug!(
1057 %database_id,
1058 "inject initial barrier"
1059 );
1060 };
1061
1062 let mut independent_checkpoint_job_controls: HashMap<
1063 JobId,
1064 IndependentCheckpointJobControl,
1065 > = HashMap::new();
1066 for (job_id, (info, upstream_table_ids, committed_epoch, snapshot_epoch)) in
1067 ongoing_snapshot_backfill_jobs
1068 {
1069 let node_actors = edges.collect_actors_to_create(info.values().map(|fragment_infos| {
1070 (
1071 fragment_infos.fragment_id,
1072 &fragment_infos.nodes,
1073 fragment_infos.actors.iter().map(move |(actor_id, actor)| {
1074 (
1075 stream_actors.get(actor_id).expect("should exist"),
1076 actor.worker_id,
1077 )
1078 }),
1079 vec![], )
1081 }));
1082
1083 let database_job_source_splits =
1084 collect_source_splits(database_jobs.values().flatten(), source_splits);
1085 assert!(
1086 !cdc_table_snapshot_splits.contains_key(&job_id),
1087 "snapshot backfill job {job_id} should not have cdc backfill"
1088 );
1089 if is_paused {
1090 bail!("should not pause when having snapshot backfill job {job_id}");
1091 }
1092 let job_backfill_orders = job_backfill_orders(job_extra_info, job_id);
1093 let job_backfill_orders =
1094 StreamFragmentGraph::extend_fragment_backfill_ordering_with_locality_backfill(
1095 job_backfill_orders,
1096 fragment_relations,
1097 || {
1098 info.iter().map(|(fragment_id, fragment)| {
1099 (*fragment_id, fragment.fragment_type_mask, &fragment.nodes)
1100 })
1101 },
1102 );
1103 let mutation = build_mutation(
1104 &database_job_source_splits,
1105 Default::default(), &job_backfill_orders,
1107 false,
1108 );
1109
1110 let job = CreatingStreamingJobControl::recover(
1111 database_id,
1112 job_id,
1113 upstream_table_ids,
1114 &database_job_log_epochs,
1115 snapshot_epoch,
1116 committed_epoch,
1117 &barrier_info,
1118 info,
1119 job_backfill_orders,
1120 fragment_relations,
1121 hummock_version_stats,
1122 node_actors,
1123 mutation.clone(),
1124 self,
1125 )?;
1126 independent_checkpoint_job_controls.insert(
1127 job_id,
1128 IndependentCheckpointJobControl::CreatingStreamingJob(job),
1129 );
1130 }
1131
1132 for (job_id, render_result) in batch_refresh {
1135 creating_jobs.remove(&job_id);
1136 debug!(%job_id, "recovered batch refresh job");
1137
1138 let committed_epoch = resolve_jobs_committed_epoch(
1140 state_table_committed_epochs,
1141 InflightFragmentInfo::existing_table_ids(render_result.fragment_infos.values()),
1142 );
1143
1144 let snapshot_backfill_info = StreamFragmentGraph::collect_snapshot_backfill_info_impl(
1145 render_result
1146 .fragment_infos
1147 .values()
1148 .map(|fragment| (&fragment.nodes, fragment.fragment_type_mask)),
1149 )?
1150 .0
1151 .ok_or_else(|| anyhow!("batch refresh job {} has no snapshot backfill info", job_id))?;
1152
1153 let upstream_table_ids: HashSet<TableId> = snapshot_backfill_info
1154 .upstream_mv_table_id_to_backfill_epoch
1155 .keys()
1156 .copied()
1157 .collect();
1158 let snapshot_epoch = snapshot_backfill_info
1159 .upstream_mv_table_id_to_backfill_epoch
1160 .values()
1161 .find_map(|e| *e)
1162 .unwrap_or(committed_epoch);
1163
1164 let job_backfill_orders = job_backfill_orders(job_extra_info, job_id);
1165 let job_backfill_orders =
1166 StreamFragmentGraph::extend_fragment_backfill_ordering_with_locality_backfill(
1167 job_backfill_orders,
1168 fragment_relations,
1169 || {
1170 render_result
1171 .fragment_infos
1172 .iter()
1173 .map(|(fid, f)| (*fid, f.fragment_type_mask, &f.nodes))
1174 },
1175 );
1176 let mutation = build_mutation(
1177 &Default::default(), Default::default(),
1179 &job_backfill_orders,
1180 false,
1181 );
1182
1183 let refresh_interval_sec = job_extra_info
1184 .get(&job_id)
1185 .and_then(|info| info.refresh_interval_sec)
1186 .expect("batch refresh job should have refresh_interval_sec in job extra info");
1187
1188 let job = BatchRefreshJobCheckpointControl::recover(
1189 database_id,
1190 job_id,
1191 upstream_table_ids,
1192 snapshot_epoch,
1193 committed_epoch,
1194 job_backfill_orders,
1195 hummock_version_stats,
1196 mutation,
1197 render_result,
1198 self,
1199 refresh_interval_sec,
1200 )?;
1201 independent_checkpoint_job_controls
1202 .insert(job_id, IndependentCheckpointJobControl::BatchRefresh(job));
1203 }
1204
1205 self.control_stream_manager()
1206 .env
1207 .shared_actor_infos()
1208 .recover_database(
1209 database_id,
1210 database_jobs
1211 .values()
1212 .flat_map(|info| {
1213 info.fragment_infos()
1214 .map(move |fragment| (fragment, info.job_id))
1215 })
1216 .chain(
1217 independent_checkpoint_job_controls
1218 .iter()
1219 .flat_map(|(job_id, job)| {
1220 let job_id = *job_id;
1221 job.fragment_infos()
1222 .into_iter()
1223 .flat_map(move |infos| infos.values().map(move |f| (f, job_id)))
1224 }),
1225 ),
1226 );
1227
1228 let committed_epoch = barrier_info.prev_epoch();
1229 let new_epoch = barrier_info.curr_epoch;
1230 let database_info = InflightDatabaseInfo::recover(
1231 database_id,
1232 database_jobs.into_values(),
1233 self.control_stream_manager()
1234 .env
1235 .shared_actor_infos()
1236 .clone(),
1237 );
1238 let database_state = BarrierWorkerState::recovery(new_epoch, is_paused);
1239 Ok(DatabaseCheckpointControl::recovery(
1240 database_id,
1241 database_state,
1242 committed_epoch,
1243 database_info,
1244 independent_checkpoint_job_controls,
1245 ))
1246 }
1247}
1248
1249impl ControlStreamManager {
1250 fn connected_workers(&self) -> impl Iterator<Item = (WorkerId, &ControlStreamNode)> + '_ {
1251 self.workers
1252 .iter()
1253 .filter_map(|(worker_id, (_, worker_state))| match worker_state {
1254 WorkerNodeState::Connected { control_stream, .. } => {
1255 Some((*worker_id, control_stream))
1256 }
1257 WorkerNodeState::Reconnecting(_) => None,
1258 })
1259 }
1260
1261 pub(super) fn inject_barrier(
1262 &mut self,
1263 partial_graph_id: PartialGraphId,
1264 mutation: Option<Mutation>,
1265 iceberg_pk_index_compaction: Option<IcebergPkIndexCompactionContext>,
1266 barrier_info: &BarrierInfo,
1267 node_actors: &HashMap<WorkerId, HashSet<ActorId>>,
1268 table_ids_to_sync: impl Iterator<Item = TableId>,
1269 nodes_to_sync_table: impl Iterator<Item = WorkerId>,
1270 mut new_actors: Option<StreamJobActorsToCreate>,
1271 ) -> MetaResult<NodeToCollect> {
1272 fail_point!("inject_barrier_err", |_| risingwave_common::bail!(
1273 "inject_barrier_err"
1274 ));
1275
1276 let nodes_to_sync_table: HashSet<_> = nodes_to_sync_table.collect();
1277
1278 nodes_to_sync_table.iter().for_each(|worker_id| {
1279 assert!(node_actors.contains_key(worker_id), "worker_id {worker_id} in nodes_to_sync_table {nodes_to_sync_table:?} but not in node_actors {node_actors:?}");
1280 });
1281
1282 let mut node_need_collect = NodeToCollect::new();
1283 let table_ids_to_sync = table_ids_to_sync.collect_vec();
1284
1285 node_actors.iter()
1286 .try_for_each(|(worker_id, actor_ids_to_collect)| {
1287 assert!(!actor_ids_to_collect.is_empty(), "empty actor_ids_to_collect on worker {worker_id} in node_actors {node_actors:?}");
1288 let table_ids_to_sync = if nodes_to_sync_table.contains(worker_id) {
1289 table_ids_to_sync.clone()
1290 } else {
1291 vec![]
1292 };
1293
1294 let node = if let Some((_, worker_state)) = self.workers.get(worker_id)
1295 &&
1296 let WorkerNodeState::Connected { control_stream, .. } = worker_state
1297 {
1298 control_stream
1299 } else {
1300 return Err(anyhow!("unconnected worker node {}", worker_id).into());
1301 };
1302
1303 {
1304 let mutation = mutation.clone();
1305 let barrier = Barrier {
1306 epoch: Some(risingwave_pb::data::Epoch {
1307 curr: barrier_info.curr_epoch(),
1308 prev: barrier_info.prev_epoch(),
1309 }),
1310 mutation: mutation.clone().map(|_| BarrierMutation { mutation }),
1311 tracing_context: TracingContext::from_span(barrier_info.curr_epoch.span())
1312 .to_protobuf(),
1313 kind: barrier_info.kind.to_protobuf() as i32,
1314 iceberg_pk_index_compaction,
1315 };
1316
1317 node.handle
1318 .request_sender
1319 .send(StreamingControlStreamRequest {
1320 request: Some(
1321 streaming_control_stream_request::Request::InjectBarrier(
1322 InjectBarrierRequest {
1323 request_id: Uuid::new_v4().to_string(),
1324 barrier: Some(barrier),
1325 actor_ids_to_collect: actor_ids_to_collect.iter().copied().collect(),
1326 table_ids_to_sync,
1327 partial_graph_id,
1328 actors_to_build: new_actors
1329 .as_mut()
1330 .map(|new_actors| new_actors.remove(worker_id))
1331 .into_iter()
1332 .flatten()
1333 .flatten()
1334 .map(|(fragment_id, (node, actors, initial_subscriber_ids))| {
1335 FragmentBuildActorInfo {
1336 fragment_id,
1337 node: Some(node),
1338 actors: actors
1339 .into_iter()
1340 .map(|(actor, upstreams, dispatchers)| {
1341 BuildActorInfo {
1342 actor_id: actor.actor_id,
1343 fragment_upstreams: upstreams
1344 .into_iter()
1345 .map(|(fragment_id, upstreams)| {
1346 (
1347 fragment_id,
1348 UpstreamActors {
1349 actors: upstreams
1350 .into_values()
1351 .collect(),
1352 },
1353 )
1354 })
1355 .collect(),
1356 dispatchers,
1357 vnode_bitmap: actor.vnode_bitmap.map(|bitmap| bitmap.to_protobuf()),
1358 mview_definition: actor.mview_definition,
1359 expr_context: actor.expr_context,
1360 config_override: actor.config_override.to_string(),
1361 initial_subscriber_ids: initial_subscriber_ids.iter().copied().collect(),
1362 }
1363 })
1364 .collect(),
1365 }
1366 })
1367 .collect(),
1368 },
1369 ),
1370 ),
1371 })
1372 .map_err(|_| {
1373 MetaError::from(anyhow!(
1374 "failed to send request to {} {:?}",
1375 node.worker_id,
1376 node.host
1377 ))
1378 })?;
1379
1380 node_need_collect.insert(*worker_id);
1381 Result::<_, MetaError>::Ok(())
1382 }
1383 })
1384 .inspect_err(|e| {
1385 use risingwave_pb::meta::event_log;
1387 let event = event_log::EventInjectBarrierFail {
1388 prev_epoch: barrier_info.prev_epoch(),
1389 cur_epoch: barrier_info.curr_epoch(),
1390 error: e.to_report_string(),
1391 };
1392 self.env
1393 .event_log_manager_ref()
1394 .add_event_logs(vec![event_log::Event::InjectBarrierFail(event)]);
1395 })?;
1396 Ok(node_need_collect)
1397 }
1398
1399 pub(super) fn add_partial_graph(&mut self, partial_graph_id: PartialGraphId) {
1400 self.connected_workers().for_each(|(_, node)| {
1401 if node
1402 .handle
1403 .request_sender
1404 .send(StreamingControlStreamRequest {
1405 request: Some(
1406 streaming_control_stream_request::Request::CreatePartialGraph(
1407 CreatePartialGraphRequest {
1408 partial_graph_id,
1409 },
1410 ),
1411 ),
1412 }).is_err() {
1413 let (database_id, creating_job_id) = from_partial_graph_id(partial_graph_id);
1414 warn!(%database_id, ?creating_job_id, worker_id = %node.worker_id, "failed to add the partial graph to the worker")
1415 }
1416 });
1417 }
1418
1419 pub(super) fn remove_partial_graphs(&mut self, partial_graph_ids: Vec<PartialGraphId>) {
1420 self.connected_workers().for_each(|(_, node)| {
1421 if node.handle
1422 .request_sender
1423 .send(StreamingControlStreamRequest {
1424 request: Some(
1425 streaming_control_stream_request::Request::RemovePartialGraph(
1426 RemovePartialGraphRequest {
1427 partial_graph_ids: partial_graph_ids.clone(),
1428 },
1429 ),
1430 ),
1431 })
1432 .is_err()
1433 {
1434 warn!(worker_id = %node.worker_id,node = ?node.host,"failed to send remove partial graph request");
1435 }
1436 })
1437 }
1438
1439 pub(super) fn reset_partial_graphs(
1440 &mut self,
1441 partial_graph_ids: Vec<PartialGraphId>,
1442 ) -> HashSet<WorkerId> {
1443 self.connected_workers()
1444 .filter_map(|(worker_id, node)| {
1445 if node
1446 .handle
1447 .request_sender
1448 .send(StreamingControlStreamRequest {
1449 request: Some(
1450 streaming_control_stream_request::Request::ResetPartialGraphs(
1451 ResetPartialGraphsRequest {
1452 partial_graph_ids: partial_graph_ids.clone(),
1453 },
1454 ),
1455 ),
1456 })
1457 .is_err()
1458 {
1459 warn!(%worker_id, node = ?node.host,"failed to send reset database request");
1460 None
1461 } else {
1462 Some(worker_id)
1463 }
1464 })
1465 .collect()
1466 }
1467}
1468
1469impl GlobalBarrierWorkerContextImpl {
1470 pub(super) async fn new_control_stream_impl(
1471 &self,
1472 node: &WorkerNode,
1473 init_request: &PbInitRequest,
1474 ) -> MetaResult<StreamingControlHandle> {
1475 let handle = self
1476 .env
1477 .stream_client_pool()
1478 .get(node)
1479 .await?
1480 .start_streaming_control(init_request.clone())
1481 .await?;
1482 Ok(handle)
1483 }
1484}
1485
1486pub(super) fn merge_node_rpc_errors<E: Error + Send + Sync + 'static>(
1487 message: &str,
1488 errors: impl IntoIterator<Item = (WorkerId, E)>,
1489) -> MetaError {
1490 use std::fmt::Write;
1491
1492 use risingwave_common::error::error_request_copy;
1493 use risingwave_common::error::tonic::extra::Score;
1494
1495 let errors = errors.into_iter().collect_vec();
1496
1497 if errors.is_empty() {
1498 return anyhow!(message.to_owned()).into();
1499 }
1500
1501 let single_error = |(worker_id, e)| {
1503 anyhow::Error::from(e)
1504 .context(format!("{message}, in worker node {worker_id}"))
1505 .into()
1506 };
1507
1508 if errors.len() == 1 {
1509 return single_error(errors.into_iter().next().unwrap());
1510 }
1511
1512 let max_score = errors
1514 .iter()
1515 .filter_map(|(_, e)| error_request_copy::<Score>(e))
1516 .max();
1517
1518 if let Some(max_score) = max_score {
1519 let mut errors = errors;
1520 let max_scored = errors
1521 .extract_if(.., |(_, e)| {
1522 error_request_copy::<Score>(e) == Some(max_score)
1523 })
1524 .next()
1525 .unwrap();
1526
1527 return single_error(max_scored);
1528 }
1529
1530 let concat: String = errors
1532 .into_iter()
1533 .fold(format!("{message}: "), |mut s, (w, e)| {
1534 write!(&mut s, " in worker node {}, {};", w, e.as_report()).unwrap();
1535 s
1536 });
1537 anyhow!(concat).into()
1538}
1539
1540#[cfg(test)]
1541mod test_partial_graph_id {
1542 use crate::barrier::rpc::{from_partial_graph_id, to_partial_graph_id};
1543
1544 #[test]
1545 fn test_partial_graph_id_conversion() {
1546 let database_id = 233.into();
1547 let job_id = 233.into();
1548 assert_eq!(
1549 (database_id, None),
1550 from_partial_graph_id(to_partial_graph_id(database_id, None))
1551 );
1552 assert_eq!(
1553 (database_id, Some(job_id)),
1554 from_partial_graph_id(to_partial_graph_id(database_id, Some(job_id)))
1555 );
1556 }
1557}