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