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