Skip to main content

risingwave_meta/manager/
metadata.rs

1// Copyright 2024 RisingWave Labs
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use 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    /// Return an uninitialized one as a placeholder for future initialized
95    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    /// Get and filter the "**root**" fragments of the specified relations.
371    /// The root fragment is the bottom-most fragment of its fragment graph, and can be a `MView` or a `Source`.
372    ///
373    /// See also [`crate::controller::catalog::CatalogController::get_root_fragments`].
374    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    // (backfill_actor_id, upstream_source_actor_id)
515    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    /// Returns jobs containing a scan that cannot be rescheduled online.
716    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    /// Wait for job finishing notification in `TrackingJob::finish`.
804    /// The progress is updated per barrier.
805    #[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}