1use std::collections::{BTreeMap, HashMap, HashSet};
16use std::fmt::{Debug, Formatter};
17
18use anyhow::anyhow;
19use itertools::Itertools;
20use risingwave_common::catalog::{DatabaseId, TableId, TableOption};
21use risingwave_common::id::JobId;
22use risingwave_meta_model::{SinkId, SourceId, WorkerId, refresh_job};
23use risingwave_pb::catalog::{PbSource, PbTable};
24use risingwave_pb::common::worker_node::{PbResource, Property as AddNodeProperty, State};
25use risingwave_pb::common::{HostAddress, PbWorkerNode, PbWorkerType, WorkerNode, WorkerType};
26use risingwave_pb::meta::list_rate_limits_response::RateLimitInfo;
27use risingwave_pb::stream_plan::{PbDispatcherType, PbStreamNode, PbStreamScanType};
28use sea_orm::TransactionTrait;
29use sea_orm::prelude::DateTime;
30use tokio::sync::mpsc::{UnboundedReceiver, unbounded_channel};
31use tokio::sync::oneshot;
32use tracing::warn;
33
34use crate::MetaResult;
35use crate::controller::catalog::CatalogControllerRef;
36use crate::controller::cluster::{ClusterControllerRef, StreamingClusterInfo, WorkerExtraInfo};
37use crate::controller::scale::find_fragment_no_shuffle_dags_detailed;
38use crate::manager::{LocalNotification, NotificationVersion};
39use crate::model::{ActorId, ClusterId, Fragment, FragmentId, StreamJobFragments, SubscriptionId};
40use crate::stream::SplitAssignment;
41use crate::telemetry::MetaTelemetryJobDesc;
42
43#[derive(Clone)]
44pub struct MetadataManager {
45 pub cluster_controller: ClusterControllerRef,
46 pub catalog_controller: CatalogControllerRef,
47}
48
49#[derive(Debug)]
50pub(crate) enum ActiveStreamingWorkerChange {
51 Add(WorkerNode),
52 Remove(WorkerNode),
53 Update(WorkerNode),
54}
55
56pub struct ActiveStreamingWorkerNodes {
57 worker_nodes: HashMap<WorkerId, WorkerNode>,
58 rx: UnboundedReceiver<LocalNotification>,
59 #[cfg_attr(not(debug_assertions), expect(dead_code))]
60 meta_manager: Option<MetadataManager>,
61}
62
63impl Debug for ActiveStreamingWorkerNodes {
64 fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
65 f.debug_struct("ActiveStreamingWorkerNodes")
66 .field("worker_nodes", &self.worker_nodes)
67 .finish()
68 }
69}
70
71impl ActiveStreamingWorkerNodes {
72 pub(crate) fn uninitialized() -> Self {
73 Self {
74 worker_nodes: Default::default(),
75 rx: unbounded_channel().1,
76 meta_manager: None,
77 }
78 }
79
80 #[cfg(test)]
81 pub(crate) fn for_test(worker_nodes: HashMap<WorkerId, WorkerNode>) -> Self {
82 let (tx, rx) = unbounded_channel();
83 let _join_handle = tokio::spawn(async move {
84 let _tx = tx;
85 std::future::pending::<()>().await
86 });
87 Self {
88 worker_nodes,
89 rx,
90 meta_manager: None,
91 }
92 }
93
94 pub(crate) async fn new_snapshot(meta_manager: MetadataManager) -> MetaResult<Self> {
96 let (nodes, rx) = meta_manager
97 .subscribe_active_streaming_compute_nodes()
98 .await?;
99 Ok(Self {
100 worker_nodes: nodes
101 .into_iter()
102 .filter_map(|node| {
103 let is_streaming = node.property.as_ref().is_some_and(|p| p.is_streaming);
104 is_streaming.then_some((node.id, node))
105 })
106 .collect(),
107 rx,
108 meta_manager: Some(meta_manager),
109 })
110 }
111
112 pub(crate) fn current(&self) -> &HashMap<WorkerId, WorkerNode> {
113 &self.worker_nodes
114 }
115
116 pub(crate) async fn changed(&mut self) -> ActiveStreamingWorkerChange {
117 loop {
118 let notification = self
119 .rx
120 .recv()
121 .await
122 .expect("notification stopped or uninitialized");
123 fn is_target_worker_node(worker: &WorkerNode) -> bool {
124 worker.r#type == WorkerType::ComputeNode as i32
125 && worker.property.as_ref().unwrap().is_streaming
126 }
127 match notification {
128 LocalNotification::WorkerNodeDeleted(worker) => {
129 let is_target_worker_node = is_target_worker_node(&worker);
130 let Some(prev_worker) = self.worker_nodes.remove(&worker.id) else {
131 if is_target_worker_node {
132 warn!(
133 ?worker,
134 "notify to delete an non-existing streaming compute worker"
135 );
136 }
137 continue;
138 };
139 if !is_target_worker_node {
140 warn!(
141 ?worker,
142 ?prev_worker,
143 "deleted worker has a different recent type"
144 );
145 }
146 if worker.state == State::Starting as i32 {
147 warn!(
148 id = %worker.id,
149 host = ?worker.host,
150 state = worker.state,
151 "a starting streaming worker is deleted"
152 );
153 }
154 break ActiveStreamingWorkerChange::Remove(prev_worker);
155 }
156 LocalNotification::WorkerNodeActivated(worker) => {
157 if !is_target_worker_node(&worker) {
158 if let Some(prev_worker) = self.worker_nodes.remove(&worker.id) {
159 warn!(
160 ?worker,
161 ?prev_worker,
162 "the type of a streaming worker is changed"
163 );
164 break ActiveStreamingWorkerChange::Remove(prev_worker);
165 } else {
166 continue;
167 }
168 }
169 assert_eq!(
170 worker.state,
171 State::Running as i32,
172 "not started worker added: {:?}",
173 worker
174 );
175 if let Some(prev_worker) = self.worker_nodes.insert(worker.id, worker.clone()) {
176 assert_eq!(prev_worker.host, worker.host);
177 assert_eq!(prev_worker.r#type, worker.r#type);
178 warn!(
179 ?prev_worker,
180 ?worker,
181 eq = prev_worker == worker,
182 "notify to update an existing active worker"
183 );
184 if prev_worker == worker {
185 continue;
186 } else {
187 break ActiveStreamingWorkerChange::Update(worker);
188 }
189 } else {
190 break ActiveStreamingWorkerChange::Add(worker);
191 }
192 }
193 _ => {
194 continue;
195 }
196 }
197 }
198 }
199
200 #[cfg(debug_assertions)]
201 pub(crate) async fn validate_change(&self) {
202 use risingwave_pb::common::WorkerNode;
203 use thiserror_ext::AsReport;
204 let Some(meta_manager) = &self.meta_manager else {
205 return;
206 };
207 match meta_manager.list_active_streaming_compute_nodes().await {
208 Ok(worker_nodes) => {
209 let ignore_irrelevant_info = |node: &WorkerNode| {
210 (
211 node.id,
212 WorkerNode {
213 id: node.id,
214 r#type: node.r#type,
215 host: node.host.clone(),
216 property: node.property.clone(),
217 resource: node.resource.clone(),
218 ..Default::default()
219 },
220 )
221 };
222 let worker_nodes: HashMap<_, _> =
223 worker_nodes.iter().map(ignore_irrelevant_info).collect();
224 let curr_worker_nodes: HashMap<_, _> = self
225 .current()
226 .values()
227 .map(ignore_irrelevant_info)
228 .collect();
229 if worker_nodes != curr_worker_nodes {
230 warn!(
231 ?worker_nodes,
232 ?curr_worker_nodes,
233 "different to global snapshot"
234 );
235 }
236 }
237 Err(e) => {
238 warn!(
239 e = ?e.as_report(),
240 "failed to list active streaming compute nodes for comparison with the local snapshot",
241 );
242 }
243 }
244 }
245}
246
247impl MetadataManager {
248 pub fn new(
249 cluster_controller: ClusterControllerRef,
250 catalog_controller: CatalogControllerRef,
251 ) -> Self {
252 Self {
253 cluster_controller,
254 catalog_controller,
255 }
256 }
257
258 pub async fn get_worker_by_id(&self, worker_id: WorkerId) -> MetaResult<Option<PbWorkerNode>> {
259 self.cluster_controller.get_worker_by_id(worker_id).await
260 }
261
262 pub async fn count_worker_node(&self) -> MetaResult<HashMap<WorkerType, u64>> {
263 let node_map = self.cluster_controller.count_worker_by_type().await?;
264 Ok(node_map
265 .into_iter()
266 .map(|(ty, cnt)| (ty.into(), cnt as u64))
267 .collect())
268 }
269
270 pub async fn get_worker_info_by_id(&self, worker_id: WorkerId) -> Option<WorkerExtraInfo> {
271 self.cluster_controller
272 .get_worker_info_by_id(worker_id as _)
273 .await
274 }
275
276 pub async fn add_worker_node(
277 &self,
278 r#type: PbWorkerType,
279 host_address: HostAddress,
280 property: AddNodeProperty,
281 resource: PbResource,
282 ) -> MetaResult<WorkerId> {
283 self.cluster_controller
284 .add_worker(r#type, host_address, property, resource)
285 .await
286 .map(|id| id as WorkerId)
287 }
288
289 pub async fn list_worker_node(
290 &self,
291 worker_type: Option<WorkerType>,
292 worker_state: Option<State>,
293 ) -> MetaResult<Vec<PbWorkerNode>> {
294 self.cluster_controller
295 .list_workers(worker_type.map(Into::into), worker_state.map(Into::into))
296 .await
297 }
298
299 pub async fn subscribe_active_streaming_compute_nodes(
300 &self,
301 ) -> MetaResult<(Vec<WorkerNode>, UnboundedReceiver<LocalNotification>)> {
302 self.cluster_controller
303 .subscribe_active_streaming_compute_nodes()
304 .await
305 }
306
307 pub async fn list_active_streaming_compute_nodes(&self) -> MetaResult<Vec<PbWorkerNode>> {
308 self.cluster_controller
309 .list_active_streaming_workers()
310 .await
311 }
312
313 pub async fn list_active_serving_compute_nodes(&self) -> MetaResult<Vec<PbWorkerNode>> {
314 self.cluster_controller.list_active_serving_workers().await
315 }
316
317 pub async fn list_active_database_ids(&self) -> MetaResult<HashSet<DatabaseId>> {
318 Ok(self
319 .catalog_controller
320 .list_fragment_database_ids(None)
321 .await?
322 .into_iter()
323 .map(|(_, database_id)| database_id)
324 .collect())
325 }
326
327 pub async fn split_fragment_map_by_database<T: Debug>(
328 &self,
329 fragment_map: HashMap<FragmentId, T>,
330 ) -> MetaResult<HashMap<DatabaseId, HashMap<FragmentId, T>>> {
331 let fragment_to_database_map: HashMap<_, _> = self
332 .catalog_controller
333 .list_fragment_database_ids(Some(
334 fragment_map
335 .keys()
336 .map(|fragment_id| *fragment_id as _)
337 .collect(),
338 ))
339 .await?
340 .into_iter()
341 .map(|(fragment_id, database_id)| (fragment_id as FragmentId, database_id))
342 .collect();
343 let mut ret: HashMap<_, HashMap<_, _>> = HashMap::new();
344 for (fragment_id, value) in fragment_map {
345 let database_id = *fragment_to_database_map
346 .get(&fragment_id)
347 .ok_or_else(|| anyhow!("cannot get database_id of fragment {fragment_id}"))?;
348 ret.entry(database_id)
349 .or_default()
350 .try_insert(fragment_id, value)
351 .expect("non duplicate");
352 }
353 Ok(ret)
354 }
355
356 pub async fn list_creating_jobs(&self) -> MetaResult<HashSet<JobId>> {
357 Ok(self
358 .catalog_controller
359 .list_creating_jobs(false, None)
360 .await?
361 .into_iter()
362 .map(|(job_id, _, _, _, _)| job_id)
363 .collect())
364 }
365
366 pub async fn list_sources(&self) -> MetaResult<Vec<PbSource>> {
367 self.catalog_controller.list_sources().await
368 }
369
370 pub async fn get_upstream_root_fragments(
375 &self,
376 upstream_table_ids: &HashSet<TableId>,
377 ) -> MetaResult<HashMap<JobId, Fragment>> {
378 let upstream_root_fragments = self
379 .catalog_controller
380 .get_root_fragments(upstream_table_ids.iter().map(|id| id.as_job_id()).collect())
381 .await?;
382
383 Ok(upstream_root_fragments)
384 }
385
386 pub async fn get_streaming_cluster_info(&self) -> MetaResult<StreamingClusterInfo> {
387 self.cluster_controller.get_streaming_cluster_info().await
388 }
389
390 pub async fn get_all_table_options(&self) -> MetaResult<HashMap<TableId, TableOption>> {
391 self.catalog_controller.get_all_table_options().await
392 }
393
394 pub async fn get_table_name_type_mapping(
395 &self,
396 ) -> MetaResult<HashMap<TableId, (String, String)>> {
397 self.catalog_controller.get_table_name_type_mapping().await
398 }
399
400 pub async fn get_created_table_ids(&self) -> MetaResult<Vec<TableId>> {
401 self.catalog_controller.get_created_table_ids().await
402 }
403
404 pub async fn get_table_associated_source_id(
405 &self,
406 table_id: TableId,
407 ) -> MetaResult<Option<SourceId>> {
408 self.catalog_controller
409 .get_table_associated_source_id(table_id)
410 .await
411 }
412
413 pub async fn get_table_catalog_by_ids(&self, ids: &[TableId]) -> MetaResult<Vec<PbTable>> {
414 self.catalog_controller
415 .get_table_by_ids(ids.to_vec(), false)
416 .await
417 }
418
419 pub async fn list_refresh_jobs(&self) -> MetaResult<Vec<refresh_job::Model>> {
420 self.catalog_controller.list_refresh_jobs().await
421 }
422
423 pub async fn list_refreshing_jobs(
424 &self,
425 database_id: Option<DatabaseId>,
426 ) -> MetaResult<Vec<refresh_job::Model>> {
427 self.catalog_controller
428 .list_refreshing_jobs(database_id)
429 .await
430 }
431
432 pub async fn list_refreshable_table_ids(&self) -> MetaResult<Vec<TableId>> {
433 self.catalog_controller.list_refreshable_table_ids().await
434 }
435
436 pub async fn ensure_refresh_job(&self, table_id: TableId) -> MetaResult<()> {
437 self.catalog_controller.ensure_refresh_job(table_id).await
438 }
439
440 pub async fn begin_refresh_job(
441 &self,
442 table_id: TableId,
443 trigger_time: DateTime,
444 ) -> MetaResult<bool> {
445 self.catalog_controller
446 .begin_refresh_job(table_id, trigger_time)
447 .await
448 }
449
450 pub async fn finish_refresh_job(
451 &self,
452 table_id: TableId,
453 trigger_time: DateTime,
454 success: bool,
455 ) -> MetaResult<bool> {
456 self.catalog_controller
457 .finish_refresh_job(table_id, trigger_time, success)
458 .await
459 }
460
461 pub async fn update_refresh_job_interval(
462 &self,
463 table_id: TableId,
464 trigger_interval_secs: Option<i64>,
465 ) -> MetaResult<()> {
466 self.catalog_controller
467 .update_refresh_job_interval(table_id, trigger_interval_secs)
468 .await
469 }
470
471 pub async fn get_sink_state_table_ids(&self, sink_id: SinkId) -> MetaResult<Vec<TableId>> {
472 self.catalog_controller
473 .get_sink_state_table_ids(sink_id)
474 .await
475 }
476
477 pub async fn get_table_catalog_by_cdc_table_id(
478 &self,
479 cdc_table_id: &String,
480 ) -> MetaResult<Vec<PbTable>> {
481 self.catalog_controller
482 .get_table_by_cdc_table_id(cdc_table_id)
483 .await
484 }
485
486 pub async fn get_downstream_fragments(
487 &self,
488 job_id: JobId,
489 ) -> MetaResult<Vec<(PbDispatcherType, Fragment)>> {
490 self.catalog_controller
491 .get_downstream_fragments(job_id)
492 .await
493 }
494
495 pub async fn get_job_id_to_internal_table_ids_mapping(
496 &self,
497 ) -> Option<Vec<(JobId, Vec<TableId>)>> {
498 self.catalog_controller.get_job_internal_table_ids().await
499 }
500
501 pub async fn get_job_fragments_by_id(&self, job_id: JobId) -> MetaResult<StreamJobFragments> {
502 Ok(self
503 .catalog_controller
504 .get_job_fragments_by_id(job_id)
505 .await?
506 .0)
507 }
508
509 pub fn get_running_actors_of_fragment(&self, id: FragmentId) -> MetaResult<HashSet<ActorId>> {
510 self.catalog_controller
511 .get_running_actors_of_fragment(id as _)
512 }
513
514 pub async fn get_running_actors_for_source_backfill(
516 &self,
517 source_backfill_fragment_id: FragmentId,
518 source_fragment_id: FragmentId,
519 ) -> MetaResult<HashSet<(ActorId, ActorId)>> {
520 let actor_ids = self
521 .catalog_controller
522 .get_running_actors_for_source_backfill(
523 source_backfill_fragment_id as _,
524 source_fragment_id as _,
525 )
526 .await?;
527 Ok(actor_ids
528 .into_iter()
529 .map(|(id, upstream)| (id as ActorId, upstream as ActorId))
530 .collect())
531 }
532
533 pub fn worker_actor_count(&self) -> MetaResult<HashMap<WorkerId, usize>> {
534 let actor_cnt = self.catalog_controller.worker_actor_count()?;
535 Ok(actor_cnt
536 .into_iter()
537 .map(|(id, cnt)| (id as WorkerId, cnt))
538 .collect())
539 }
540
541 pub async fn count_streaming_job(&self) -> MetaResult<usize> {
542 self.catalog_controller.count_streaming_jobs().await
543 }
544
545 pub async fn list_stream_job_desc(&self) -> MetaResult<Vec<MetaTelemetryJobDesc>> {
546 self.catalog_controller
547 .list_stream_job_desc_for_telemetry()
548 .await
549 }
550
551 pub async fn update_source_rate_limit_by_source_id(
552 &self,
553 source_id: SourceId,
554 rate_limit: Option<u32>,
555 ) -> MetaResult<HashMap<FragmentId, PbStreamNode>> {
556 self.catalog_controller
557 .update_source_rate_limit_by_source_id(source_id as _, rate_limit)
558 .await
559 }
560
561 pub async fn update_backfill_rate_limit_by_job_id(
562 &self,
563 job_id: JobId,
564 rate_limit: Option<u32>,
565 ) -> MetaResult<HashMap<FragmentId, PbStreamNode>> {
566 self.catalog_controller
567 .update_backfill_rate_limit_by_job_id(job_id, rate_limit)
568 .await
569 }
570
571 pub async fn update_sink_rate_limit_by_sink_id(
572 &self,
573 sink_id: SinkId,
574 rate_limit: Option<u32>,
575 ) -> MetaResult<HashMap<FragmentId, PbStreamNode>> {
576 self.catalog_controller
577 .update_sink_rate_limit_by_job_id(sink_id, rate_limit)
578 .await
579 }
580
581 pub async fn update_dml_rate_limit_by_job_id(
582 &self,
583 job_id: JobId,
584 rate_limit: Option<u32>,
585 ) -> MetaResult<HashMap<FragmentId, PbStreamNode>> {
586 self.catalog_controller
587 .update_dml_rate_limit_by_job_id(job_id, rate_limit)
588 .await
589 }
590
591 pub async fn update_sink_props_by_sink_id(
592 &self,
593 sink_id: SinkId,
594 props: BTreeMap<String, String>,
595 ) -> MetaResult<HashMap<String, String>> {
596 let new_props = self
597 .catalog_controller
598 .update_sink_props_by_sink_id(sink_id, props)
599 .await?;
600 Ok(new_props)
601 }
602
603 pub async fn update_iceberg_table_props_by_table_id(
604 &self,
605 table_id: TableId,
606 props: BTreeMap<String, String>,
607 alter_iceberg_table_props: Option<
608 risingwave_pb::meta::alter_connector_props_request::PbExtraOptions,
609 >,
610 ) -> MetaResult<(HashMap<String, String>, SinkId)> {
611 let (new_props, sink_id) = self
612 .catalog_controller
613 .update_iceberg_table_props_by_table_id(table_id, props, alter_iceberg_table_props)
614 .await?;
615 Ok((new_props, sink_id))
616 }
617
618 pub async fn update_fragment_rate_limit_by_fragment_id(
619 &self,
620 fragment_id: FragmentId,
621 throttle_type: risingwave_pb::common::ThrottleType,
622 rate_limit: Option<u32>,
623 ) -> MetaResult<PbStreamNode> {
624 self.catalog_controller
625 .update_fragment_rate_limit_by_fragment_id(fragment_id as _, throttle_type, rate_limit)
626 .await
627 }
628
629 #[await_tree::instrument]
630 pub async fn update_fragment_splits(
631 &self,
632 split_assignment: &SplitAssignment,
633 ) -> MetaResult<()> {
634 let fragment_splits = split_assignment
635 .iter()
636 .map(|(fragment_id, splits)| {
637 (
638 *fragment_id as _,
639 splits.values().flatten().cloned().collect_vec(),
640 )
641 })
642 .collect();
643
644 let inner = self.catalog_controller.inner.write().await;
645
646 self.catalog_controller
647 .update_fragment_splits(&inner.db, &fragment_splits)
648 .await
649 }
650
651 pub async fn get_mv_depended_subscriptions(
652 &self,
653 database_id: Option<DatabaseId>,
654 ) -> MetaResult<HashMap<TableId, HashMap<SubscriptionId, u64>>> {
655 Ok(self
656 .catalog_controller
657 .get_mv_depended_subscriptions(database_id)
658 .await?
659 .into_iter()
660 .map(|(table_id, subscriptions)| {
661 (
662 table_id,
663 subscriptions
664 .into_iter()
665 .map(|(subscription_id, retention_time)| {
666 (subscription_id as SubscriptionId, retention_time)
667 })
668 .collect(),
669 )
670 })
671 .collect())
672 }
673
674 pub async fn get_job_max_parallelism(&self, job_id: JobId) -> MetaResult<usize> {
675 self.catalog_controller
676 .get_max_parallelism_by_id(job_id)
677 .await
678 }
679
680 pub async fn get_existing_job_resource_group(
681 &self,
682 streaming_job_id: JobId,
683 ) -> MetaResult<String> {
684 self.catalog_controller
685 .get_existing_job_resource_group(streaming_job_id)
686 .await
687 }
688
689 pub async fn get_database_resource_group(&self, database_id: DatabaseId) -> MetaResult<String> {
690 self.catalog_controller
691 .get_database_resource_group(database_id)
692 .await
693 }
694
695 pub fn cluster_id(&self) -> &ClusterId {
696 self.cluster_controller.cluster_id()
697 }
698
699 pub async fn list_rate_limits(&self) -> MetaResult<Vec<RateLimitInfo>> {
700 let rate_limits = self.catalog_controller.list_rate_limits().await?;
701 Ok(rate_limits)
702 }
703
704 pub async fn get_job_backfill_scan_types(
705 &self,
706 job_id: JobId,
707 ) -> MetaResult<HashMap<FragmentId, PbStreamScanType>> {
708 let backfill_types = self
709 .catalog_controller
710 .get_job_fragment_backfill_scan_type(job_id)
711 .await?;
712 Ok(backfill_types)
713 }
714
715 pub async fn collect_online_unreschedulable_backfill_jobs(
717 &self,
718 job_ids: impl IntoIterator<Item = &JobId>,
719 ) -> MetaResult<HashSet<JobId>> {
720 let mut unreschedulable = HashSet::new();
721
722 for job_id in job_ids {
723 let scan_types = self
724 .catalog_controller
725 .get_job_fragment_backfill_scan_type(*job_id)
726 .await?;
727 if scan_types
728 .values()
729 .any(|scan_type| !scan_type.is_reschedulable(true))
730 {
731 unreschedulable.insert(*job_id);
732 }
733 }
734
735 Ok(unreschedulable)
736 }
737
738 pub async fn collect_reschedule_blocked_jobs_for_creating_jobs(
739 &self,
740 creating_job_ids: impl IntoIterator<Item = &JobId>,
741 is_online: bool,
742 ) -> MetaResult<HashSet<JobId>> {
743 let creating_job_ids: HashSet<_> = creating_job_ids.into_iter().copied().collect();
744 if creating_job_ids.is_empty() {
745 return Ok(HashSet::new());
746 }
747
748 let inner = self.catalog_controller.inner.read().await;
749 let txn = inner.db.begin().await?;
750
751 let mut initial_fragment_ids = HashSet::new();
752 for job_id in &creating_job_ids {
753 let scan_types = self
754 .catalog_controller
755 .get_job_fragment_backfill_scan_type_in_txn(&txn, *job_id)
756 .await?;
757 initial_fragment_ids.extend(scan_types.into_iter().filter_map(
758 |(fragment_id, scan_type)| {
759 (!scan_type.is_reschedulable(is_online)).then_some(fragment_id)
760 },
761 ));
762 }
763
764 if !initial_fragment_ids.is_empty() {
765 let upstream_fragments = self
766 .catalog_controller
767 .upstream_fragments_in_txn(&txn, initial_fragment_ids.iter().copied())
768 .await?;
769 initial_fragment_ids.extend(upstream_fragments.into_values().flatten());
770 }
771
772 let mut blocked_fragment_ids = initial_fragment_ids.clone();
773 if !initial_fragment_ids.is_empty() {
774 let initial_fragment_ids = initial_fragment_ids.into_iter().collect_vec();
775 let ensembles =
776 find_fragment_no_shuffle_dags_detailed(&txn, &initial_fragment_ids).await?;
777 for ensemble in ensembles {
778 blocked_fragment_ids.extend(ensemble.fragments());
779 }
780 }
781
782 let mut blocked_job_ids = HashSet::new();
783 if !blocked_fragment_ids.is_empty() {
784 let fragment_ids = blocked_fragment_ids.into_iter().collect_vec();
785 let fragment_job_ids = self
786 .catalog_controller
787 .get_fragment_job_id_in_txn(&txn, fragment_ids)
788 .await?;
789 blocked_job_ids.extend(
790 fragment_job_ids
791 .into_iter()
792 .map(|job_id| job_id.as_job_id()),
793 );
794 }
795
796 txn.commit().await?;
797
798 Ok(blocked_job_ids)
799 }
800}
801
802impl MetadataManager {
803 #[await_tree::instrument]
806 pub async fn wait_streaming_job_finished(
807 &self,
808 database_id: DatabaseId,
809 id: JobId,
810 ) -> MetaResult<NotificationVersion> {
811 tracing::debug!("wait_streaming_job_finished: {id:?}");
812 let mut mgr = self.catalog_controller.get_inner_write_guard().await;
813 if mgr.streaming_job_is_finished(id).await? {
814 return Ok(self.catalog_controller.notify_frontend_trivial().await);
815 }
816 let (tx, rx) = oneshot::channel();
817
818 mgr.register_finish_notifier(database_id, id, tx);
819 drop(mgr);
820 rx.await
821 .map_err(|_| "no received reason".to_owned())
822 .and_then(|result| result)
823 .map_err(|reason| anyhow!("failed to wait streaming job finish: {}", reason).into())
824 }
825
826 pub(crate) async fn notify_finish_failed(&self, database_id: Option<DatabaseId>, err: String) {
827 let mut mgr = self.catalog_controller.get_inner_write_guard().await;
828 mgr.notify_finish_failed(database_id, err);
829 }
830}