Skip to main content

risingwave_meta/stream/
stream_manager.rs

1// Copyright 2022 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::HashMap;
16use std::sync::Arc;
17
18use anyhow::Context;
19use await_tree::span;
20use futures::future::join_all;
21use itertools::Itertools;
22use risingwave_common::bail;
23use risingwave_common::catalog::{DatabaseId, Field, FragmentTypeFlag, FragmentTypeMask, TableId};
24use risingwave_common::hash::VnodeCountCompat;
25use risingwave_common::id::{JobId, SinkId};
26use risingwave_connector::source::CdcTableSnapshotSplitRaw;
27use risingwave_meta_model::prelude::Fragment as FragmentModel;
28use risingwave_meta_model::{StreamingParallelism, WorkerId, fragment, streaming_job};
29use risingwave_pb::catalog::{CreateType, PbSink, PbTable, Subscription};
30use risingwave_pb::ddl_service::streaming_job_resource_type;
31use risingwave_pb::expr::PbExprNode;
32use risingwave_pb::plan_common::{PbColumnCatalog, PbField};
33use risingwave_pb::serverless_backfill_controller::{
34    ProvisionRequest, node_group_controller_service_client,
35};
36use risingwave_rpc_client::error::TonicStatusWrapper;
37use sea_orm::{ColumnTrait, EntityTrait, QueryFilter, QuerySelect};
38use thiserror_ext::AsReport;
39use tokio::sync::{Mutex, OwnedSemaphorePermit, RwLockReadGuard, oneshot};
40use tokio::time::{Duration, Instant};
41use tracing::Instrument;
42
43use super::{
44    GlobalRefreshManagerRef, ParallelismPolicy, ReschedulePolicy, ScaleControllerRef,
45    StreamFragmentGraph, UserDefinedFragmentBackfillOrder,
46};
47use crate::barrier::{
48    BarrierScheduler, Command, CreateStreamingJobCommandInfo, CreateStreamingJobType,
49    IndependentStreamingJobType, ReplaceStreamJobPlan, SinceEpochInfo, SnapshotBackfillInfo,
50};
51use crate::controller::catalog::DropTableConnectorContext;
52use crate::controller::fragment::{InflightActorInfo, InflightFragmentInfo};
53use crate::error::bail_invalid_parameter;
54use crate::hummock::HummockManagerRef;
55use crate::manager::iceberg_compaction::IcebergCompactionManagerRef;
56use crate::manager::{
57    MetaSrvEnv, MetadataManager, NotificationVersion, StreamingJob, StreamingJobType,
58};
59use crate::model::{
60    ActorId, DownstreamFragmentRelation, Fragment, FragmentDownstreamRelation, FragmentId,
61    FragmentReplaceUpstream, StreamActor, StreamContext, StreamJobFragments,
62    StreamJobFragmentsToCreate, SubscriptionId,
63};
64use crate::stream::cdc::is_parallelized_backfill_enabled_cdc_scan_fragment;
65use crate::stream::{ReplaceJobSplitPlan, SourceManagerRef};
66use crate::{MetaError, MetaResult};
67
68pub type GlobalStreamManagerRef = Arc<GlobalStreamManager>;
69
70/// The error carries whether the caller should explicitly cancel the creating job and an optional
71/// notifier for an awaited cancellation request.
72pub type CreateStreamingJobResult =
73    Result<NotificationVersion, (MetaError, bool, Option<oneshot::Sender<bool>>)>;
74
75/// A user is assumed to stay focused on a streaming-job creation for at most 30 seconds. If an
76/// error occurs during that time, cancel the job so that they can investigate the error. After
77/// that, prioritize eventual completion by continuing to wait through transient errors.
78const FOREGROUND_DDL_EARLY_FAILURE_TIMEOUT: Duration = Duration::from_secs(30);
79
80pub(crate) async fn cleanup_dropped_streaming_jobs(
81    refresh_manager: &GlobalRefreshManagerRef,
82    hummock_manager: &HummockManagerRef,
83    metadata_manager: &MetadataManager,
84    streaming_job_ids: impl IntoIterator<Item = JobId>,
85    state_table_ids: Vec<TableId>,
86    progress_status: &str,
87) -> MetaResult<()> {
88    for job_id in streaming_job_ids {
89        refresh_manager.remove_progress_tracker(job_id.as_mv_table_id(), progress_status);
90    }
91
92    if state_table_ids.is_empty() {
93        return Ok(());
94    }
95
96    hummock_manager
97        .unregister_table_ids(state_table_ids.clone())
98        .await?;
99    metadata_manager
100        .catalog_controller
101        .complete_dropped_tables(state_table_ids)
102        .await;
103    Ok(())
104}
105
106#[derive(Default)]
107pub struct CreateStreamingJobOption {
108    // leave empty as a placeholder for future option if there is any
109}
110
111#[derive(Debug, Clone)]
112pub struct UpstreamSinkInfo {
113    pub sink_id: SinkId,
114    pub sink_fragment_id: FragmentId,
115    pub sink_output_fields: Vec<PbField>,
116    // for backwards compatibility
117    pub sink_original_target_columns: Vec<PbColumnCatalog>,
118    pub project_exprs: Vec<PbExprNode>,
119    pub new_sink_downstream: DownstreamFragmentRelation,
120}
121
122/// [`CreateStreamingJobContext`] carries one-time infos for creating a streaming job.
123///
124/// Note: for better readability, keep this struct complete and immutable once created.
125pub struct CreateStreamingJobContext {
126    /// New fragment relation to add from upstream fragments to downstream fragments.
127    pub upstream_fragment_downstreams: FragmentDownstreamRelation,
128
129    /// The resource group of the database this job belongs to.
130    pub database_resource_group: String,
131
132    /// DDL definition.
133    pub definition: String,
134
135    pub create_type: CreateType,
136
137    pub job_type: StreamingJobType,
138
139    /// Used for sink-into-table.
140    pub new_upstream_sink: Option<UpstreamSinkInfo>,
141
142    pub snapshot_backfill_info: Option<SnapshotBackfillInfo>,
143    pub cross_db_snapshot_backfill_info: SnapshotBackfillInfo,
144
145    pub cdc_table_snapshot_splits: Option<Vec<CdcTableSnapshotSplitRaw>>,
146
147    pub option: CreateStreamingJobOption,
148
149    pub streaming_job: StreamingJob,
150
151    pub fragment_backfill_ordering: UserDefinedFragmentBackfillOrder,
152
153    pub locality_fragment_state_table_mapping: HashMap<FragmentId, Vec<TableId>>,
154
155    pub is_serverless_backfill: bool,
156
157    pub resource_type: streaming_job_resource_type::ResourceType,
158
159    /// The `streaming_job::Model` for this job, loaded from meta store.
160    pub streaming_job_model: streaming_job::Model,
161
162    /// If set, this create command replaces an existing sink while creating the new sink job.
163    pub replace_sink: Option<SinkId>,
164
165    /// Batch refresh interval in seconds. If set, the MV uses batch refresh semantics.
166    pub refresh_interval_sec: Option<u64>,
167
168    pub since_timestamp_epoch: Option<u64>,
169}
170
171struct StreamingJobExecution {
172    id: JobId,
173    shutdown_tx: Option<oneshot::Sender<oneshot::Sender<bool>>>,
174    _permit: OwnedSemaphorePermit,
175}
176
177impl StreamingJobExecution {
178    fn new(
179        id: JobId,
180        shutdown_tx: oneshot::Sender<oneshot::Sender<bool>>,
181        permit: OwnedSemaphorePermit,
182    ) -> Self {
183        Self {
184            id,
185            shutdown_tx: Some(shutdown_tx),
186            _permit: permit,
187        }
188    }
189}
190
191#[derive(Default)]
192struct CreatingStreamingJobInfo {
193    streaming_jobs: Mutex<HashMap<JobId, StreamingJobExecution>>,
194}
195
196impl CreatingStreamingJobInfo {
197    async fn add_job(&self, job: StreamingJobExecution) {
198        let mut jobs = self.streaming_jobs.lock().await;
199        jobs.insert(job.id, job);
200    }
201
202    async fn delete_job(&self, job_id: JobId) {
203        let mut jobs = self.streaming_jobs.lock().await;
204        jobs.remove(&job_id);
205    }
206
207    async fn cancel_jobs(
208        &self,
209        job_ids: Vec<JobId>,
210    ) -> MetaResult<(HashMap<JobId, oneshot::Receiver<bool>>, Vec<JobId>)> {
211        let mut jobs = self.streaming_jobs.lock().await;
212        let mut receivers = HashMap::new();
213        let mut background_job_ids = vec![];
214        for job_id in job_ids {
215            if let Some(job) = jobs.get_mut(&job_id) {
216                if let Some(shutdown_tx) = job.shutdown_tx.take() {
217                    let (tx, rx) = oneshot::channel();
218                    match shutdown_tx.send(tx) {
219                        Ok(()) => {
220                            receivers.insert(job_id, rx);
221                        }
222                        Err(_) => {
223                            return Err(anyhow::anyhow!(
224                                "failed to send shutdown signal for streaming job {}: receiver dropped",
225                                job_id
226                            )
227                            .into());
228                        }
229                    }
230                }
231            } else {
232                // If these job ids do not exist in streaming_jobs, they should be background creating jobs.
233                background_job_ids.push(job_id);
234            }
235        }
236
237        Ok((receivers, background_job_ids))
238    }
239}
240
241type CreatingStreamingJobInfoRef = Arc<CreatingStreamingJobInfo>;
242
243#[derive(Debug, Clone)]
244pub struct AutoRefreshSchemaSinkContext {
245    pub tmp_sink_id: SinkId,
246    pub original_sink: PbSink,
247    pub original_fragment: Fragment,
248    pub new_schema: Vec<PbColumnCatalog>,
249    pub newly_add_fields: Vec<Field>,
250    pub removed_column_names: Vec<String>,
251    pub new_fragment: Fragment,
252    pub new_log_store_table: Option<Box<PbTable>>,
253    /// The sink's own stream context (timezone, `config_override`).
254    pub ctx: StreamContext,
255}
256
257impl AutoRefreshSchemaSinkContext {
258    pub fn new_fragment_info(
259        &self,
260        stream_actors: &HashMap<FragmentId, Vec<StreamActor>>,
261        actor_location: &HashMap<ActorId, WorkerId>,
262    ) -> InflightFragmentInfo {
263        InflightFragmentInfo {
264            fragment_id: self.new_fragment.fragment_id,
265            distribution_type: self.new_fragment.distribution_type.into(),
266            fragment_type_mask: self.new_fragment.fragment_type_mask,
267            vnode_count: self.new_fragment.vnode_count(),
268            nodes: self.new_fragment.nodes.clone(),
269            actors: stream_actors
270                .get(&self.new_fragment.fragment_id)
271                .into_iter()
272                .flatten()
273                .map(|actor| {
274                    (
275                        actor.actor_id,
276                        InflightActorInfo {
277                            worker_id: actor_location[&actor.actor_id],
278                            vnode_bitmap: actor.vnode_bitmap.clone(),
279                            splits: vec![],
280                        },
281                    )
282                })
283                .collect(),
284            state_table_ids: self.new_fragment.state_table_ids.iter().copied().collect(),
285        }
286    }
287}
288
289/// [`ReplaceStreamJobContext`] carries one-time infos for replacing the plan of an existing stream job.
290///
291/// Note: for better readability, keep this struct complete and immutable once created.
292pub struct ReplaceStreamJobContext {
293    /// The old job fragments to be replaced.
294    pub old_fragments: StreamJobFragments,
295
296    /// The updates to be applied to the downstream chain actors. Used for schema change.
297    pub replace_upstream: FragmentReplaceUpstream,
298
299    /// New fragment relation to add from existing upstream fragment to downstream fragment.
300    pub upstream_fragment_downstreams: FragmentDownstreamRelation,
301
302    pub streaming_job: StreamingJob,
303
304    /// The resource group of the database this job belongs to.
305    pub database_resource_group: String,
306
307    pub tmp_id: JobId,
308
309    /// Used for dropping an associated source. Dropping source and related internal tables.
310    pub drop_table_connector_ctx: Option<DropTableConnectorContext>,
311
312    pub auto_refresh_schema_sinks: Option<Vec<AutoRefreshSchemaSinkContext>>,
313
314    /// The `streaming_job::Model` for this job, loaded from meta store.
315    pub streaming_job_model: streaming_job::Model,
316}
317
318/// `GlobalStreamManager` manages all the streams in the system.
319pub struct GlobalStreamManager {
320    pub env: MetaSrvEnv,
321
322    pub metadata_manager: MetadataManager,
323
324    /// Broadcasts and collect barriers
325    pub barrier_scheduler: BarrierScheduler,
326
327    pub hummock_manager: HummockManagerRef,
328
329    /// Maintains streaming sources from external system like kafka
330    pub source_manager: SourceManagerRef,
331
332    pub refresh_manager: GlobalRefreshManagerRef,
333
334    pub iceberg_compaction_manager: IcebergCompactionManagerRef,
335
336    /// Creating streaming job info.
337    creating_job_info: CreatingStreamingJobInfoRef,
338
339    pub scale_controller: ScaleControllerRef,
340}
341
342impl GlobalStreamManager {
343    pub fn new(
344        env: MetaSrvEnv,
345        metadata_manager: MetadataManager,
346        barrier_scheduler: BarrierScheduler,
347        hummock_manager: HummockManagerRef,
348        source_manager: SourceManagerRef,
349        refresh_manager: GlobalRefreshManagerRef,
350        iceberg_compaction_manager: IcebergCompactionManagerRef,
351        scale_controller: ScaleControllerRef,
352    ) -> MetaResult<Self> {
353        Ok(Self {
354            env,
355            metadata_manager,
356            barrier_scheduler,
357            hummock_manager,
358            source_manager,
359            refresh_manager,
360            iceberg_compaction_manager,
361            creating_job_info: Arc::new(CreatingStreamingJobInfo::default()),
362            scale_controller,
363        })
364    }
365
366    /// Create streaming job, it works as follows:
367    ///
368    /// 1. Broadcast the actor info based on the scheduling result in the context, build the hanging
369    ///    channels in upstream worker nodes.
370    /// 2. (optional) Get the split information of the `StreamSource` via source manager and patch
371    ///    actors.
372    /// 3. Notify related worker nodes to update and build the actors.
373    /// 4. Store related meta data.
374    ///
375    /// This function is a wrapper over [`Self::run_create_streaming_job_command`].
376    #[await_tree::instrument]
377    pub async fn create_streaming_job(
378        self: &Arc<Self>,
379        stream_job_fragments: StreamJobFragmentsToCreate,
380        ctx: CreateStreamingJobContext,
381        permit: OwnedSemaphorePermit,
382        reschedule_job_lock: RwLockReadGuard<'_, ()>,
383    ) -> CreateStreamingJobResult {
384        let await_tree_key = format!("Create Streaming Job Worker ({})", ctx.streaming_job.id());
385        let await_tree_span = span!(
386            "{:?}({})",
387            ctx.streaming_job.job_type(),
388            ctx.streaming_job.name()
389        );
390
391        let job_id = stream_job_fragments.stream_job_id();
392        let database_id = ctx.streaming_job.database_id();
393
394        let (cancel_tx, cancel_rx) = oneshot::channel();
395        let execution = StreamingJobExecution::new(job_id, cancel_tx, permit);
396        self.creating_job_info.add_job(execution).await;
397
398        let stream_manager = self.clone();
399        let fut = async move {
400            let create_type = ctx.create_type;
401            let streaming_job = stream_manager
402                .run_create_streaming_job_command(stream_job_fragments, ctx)
403                .await
404                .map_err(|err| (err, false, None))?;
405            // The create command has been collected, so rescheduling no longer conflicts with
406            // planning or scheduling this job. In particular, do not hold this lock while a
407            // foreground job waits through recovery.
408            drop(reschedule_job_lock);
409            let version = match create_type {
410                CreateType::Background => {
411                    stream_manager
412                        .metadata_manager
413                        .catalog_controller
414                        .notify_frontend_trivial()
415                        .await
416                }
417                CreateType::Foreground => {
418                    let job_id = streaming_job.id() as _;
419                    let wait_started_at = Instant::now();
420                    loop {
421                        match stream_manager
422                            .metadata_manager
423                            .wait_streaming_job_finished(database_id, job_id)
424                            .await
425                        {
426                            Ok(version) => break version,
427                            Err(err) if err.is_catalog_id_not_found("streaming job") => {
428                                return Err((err, false, None));
429                            }
430                            Err(err)
431                                if wait_started_at.elapsed()
432                                    < FOREGROUND_DDL_EARLY_FAILURE_TIMEOUT =>
433                            {
434                                tracing::warn!(
435                                    id = %job_id,
436                                    error = %err.as_report(),
437                                    elapsed = ?wait_started_at.elapsed(),
438                                    "foreground streaming job failed shortly after waiting started; cancelling it"
439                                );
440                                return Err((err, true, None));
441                            }
442                            Err(err) => {
443                                tracing::warn!(
444                                    id = %job_id,
445                                    error = %err.as_report(),
446                                    "failed to wait for foreground streaming job; registering another finish notifier"
447                                );
448                            }
449                        }
450                    }
451                }
452                CreateType::Unspecified => unreachable!(),
453            };
454
455            tracing::debug!(?streaming_job, "stream job finish");
456            Ok(version)
457        }
458        .in_current_span();
459
460        let create_fut = (self.env.await_tree_reg())
461            .register(await_tree_key, await_tree_span)
462            .instrument(Box::pin(fut));
463
464        let result = async {
465            tokio::select! {
466                biased;
467
468                res = create_fut => res,
469                notifier = cancel_rx => {
470                    let notifier = notifier.expect("sender should not be dropped");
471                    tracing::debug!(id=%job_id, "cancelling streaming job");
472
473                    enum CancelResult {
474                        Completed(CreateStreamingJobResult),
475                        Cancelled { explicitly_cancel: bool },
476                    }
477
478                    let cancel_res = if let Ok(job_fragments) =
479                        self.metadata_manager.get_job_fragments_by_id(job_id).await
480                    {
481                        // try to cancel buffered creating command.
482                        if self
483                            .barrier_scheduler
484                            .try_cancel_scheduled_create(database_id, job_id)
485                        {
486                            tracing::debug!(
487                                id=%job_id,
488                                "cancelling streaming job in buffer queue."
489                            );
490                            CancelResult::Cancelled {
491                                explicitly_cancel: false,
492                            }
493                        } else if !job_fragments.is_created() {
494                            tracing::debug!(
495                                id=%job_id,
496                                "cancelling streaming job by issue cancel command."
497                            );
498                            CancelResult::Cancelled {
499                                explicitly_cancel: true,
500                            }
501                        } else {
502                            // streaming job is already completed
503                            CancelResult::Completed(
504                                self.metadata_manager
505                                    .wait_streaming_job_finished(database_id, job_id)
506                                    .await
507                                    .map_err(|err| (err, false, None)),
508                            )
509                        }
510                    } else {
511                        CancelResult::Cancelled {
512                            explicitly_cancel: false,
513                        }
514                    };
515
516                    match cancel_res {
517                        CancelResult::Completed(result) => {
518                            let _ = notifier.send(false).inspect_err(|err| {
519                                tracing::warn!("failed to notify cancellation result: {err}")
520                            });
521                            result
522                        }
523                        CancelResult::Cancelled { explicitly_cancel } => {
524                            Err((MetaError::cancelled("create"), explicitly_cancel, Some(notifier)))
525                        }
526                    }
527                }
528            }
529        }
530        .await;
531
532        tracing::debug!("cleaning creating job info: {}", job_id);
533        self.creating_job_info.delete_job(job_id).await;
534        result
535    }
536
537    async fn provision_serverless_backfill_resource_group(&self) -> MetaResult<String> {
538        let sbc_addr = &self.env.opts.serverless_backfill_controller_addr;
539        if sbc_addr.is_empty() {
540            bail_invalid_parameter!(
541                "Serverless Backfill is disabled. Use RisingWave cloud at https://cloud.risingwave.com/auth/signup to try this feature"
542            );
543        }
544
545        let request = tonic::Request::new(ProvisionRequest {});
546        let mut client =
547            node_group_controller_service_client::NodeGroupControllerServiceClient::connect(
548                sbc_addr.clone(),
549            )
550            .await
551            .with_context(|| {
552                format!(
553                    "unable to reach serverless backfill controller at addr {}",
554                    sbc_addr
555                )
556            })?;
557
558        match client.provision(request).await {
559            Ok(resp) => Ok(resp.into_inner().resource_group),
560            Err(e) => Err(anyhow::Error::new(TonicStatusWrapper::new(e))
561                .context("serverless backfill controller returned error")
562                .into()),
563        }
564    }
565
566    async fn finalize_create_streaming_job_resource_group(
567        &self,
568        resource_type: &streaming_job_resource_type::ResourceType,
569        streaming_job_model: &mut streaming_job::Model,
570    ) -> MetaResult<()> {
571        if !matches!(
572            resource_type,
573            streaming_job_resource_type::ResourceType::ServerlessBackfill(true)
574        ) {
575            return Ok(());
576        }
577
578        let group = self.provision_serverless_backfill_resource_group().await?;
579        tracing::info!(
580            resource_group = group,
581            "provisioning serverless backfill resource group"
582        );
583
584        self.metadata_manager
585            .catalog_controller
586            .update_streaming_job_resource_group(streaming_job_model.job_id, group.clone())
587            .await?;
588        streaming_job_model.specific_resource_group = Some(group);
589
590        Ok(())
591    }
592
593    /// The function will return after barrier collected
594    /// ([`crate::manager::MetadataManager::wait_streaming_job_finished`]).
595    #[await_tree::instrument]
596    async fn run_create_streaming_job_command(
597        &self,
598        stream_job_fragments: StreamJobFragmentsToCreate,
599        CreateStreamingJobContext {
600            streaming_job,
601            upstream_fragment_downstreams,
602            database_resource_group,
603            definition,
604            create_type,
605            job_type,
606            new_upstream_sink,
607            snapshot_backfill_info,
608            cross_db_snapshot_backfill_info,
609            fragment_backfill_ordering,
610            locality_fragment_state_table_mapping,
611            cdc_table_snapshot_splits,
612            is_serverless_backfill,
613            resource_type,
614            mut streaming_job_model,
615            replace_sink,
616            refresh_interval_sec,
617            since_timestamp_epoch,
618            ..
619        }: CreateStreamingJobContext,
620    ) -> MetaResult<StreamingJob> {
621        tracing::debug!(
622            table_id = %stream_job_fragments.stream_job_id(),
623            "built actors finished"
624        );
625
626        // Phase 1: Gather fragment-level split information.
627        // - For source fragments: discover splits from the external source.
628        // - For backfill fragments: splits will be aligned in Phase 2 inside the barrier worker
629        //   using the actor-level no-shuffle mapping produced by render_actors.
630        let init_split_assignment = self
631            .source_manager
632            .discover_splits(&stream_job_fragments)
633            .await?;
634
635        let fragment_backfill_ordering =
636            StreamFragmentGraph::extend_fragment_backfill_ordering_with_locality_backfill(
637                fragment_backfill_ordering,
638                &stream_job_fragments.downstreams,
639                || {
640                    stream_job_fragments
641                        .fragments
642                        .iter()
643                        .map(|(fragment_id, fragment)| {
644                            (*fragment_id, fragment.fragment_type_mask, &fragment.nodes)
645                        })
646                },
647            );
648
649        self.finalize_create_streaming_job_resource_group(&resource_type, &mut streaming_job_model)
650            .await?;
651
652        let info = CreateStreamingJobCommandInfo {
653            stream_job_fragments,
654            upstream_fragment_downstreams,
655            init_split_assignment,
656            definition: definition.clone(),
657            streaming_job: streaming_job.clone(),
658            job_type,
659            create_type,
660            database_resource_group,
661            fragment_backfill_ordering,
662            cdc_table_snapshot_splits,
663            locality_fragment_state_table_mapping,
664            is_serverless: is_serverless_backfill,
665            streaming_job_model,
666            replace_sink,
667            refresh_interval_sec,
668        };
669
670        let create_job_type = if let Some(refresh_interval_sec) = refresh_interval_sec {
671            if since_timestamp_epoch.is_some() {
672                bail!("since_timestamp should not be specified when no snapshot backfill");
673            }
674            let snapshot_backfill_info = snapshot_backfill_info.ok_or_else(|| {
675                anyhow::anyhow!(
676                    "batch refresh materialized view must have snapshot backfill upstream"
677                )
678            })?;
679            // Batch refresh jobs must not contain source or source-backfill nodes,
680            // because we skip split assignment resolution for them.
681            for fragment in info.stream_job_fragments.inner.fragments.values() {
682                let mask = fragment.fragment_type_mask;
683                if mask.contains(FragmentTypeFlag::Source)
684                    || mask.contains(FragmentTypeFlag::SourceScan)
685                {
686                    bail!(
687                        "batch refresh materialized views must not depend on sources directly; \
688                         fragment {} has source/source-backfill nodes",
689                        fragment.fragment_id
690                    );
691                }
692            }
693            tracing::debug!(
694                ?snapshot_backfill_info,
695                refresh_interval_sec,
696                "sending Command::CreateBatchRefreshStreamingJob"
697            );
698            CreateStreamingJobType::Independent {
699                snapshot_backfill_info,
700                kind: IndependentStreamingJobType::BatchRefresh {
701                    refresh_interval_sec,
702                },
703            }
704        } else if let Some(snapshot_backfill_info) = snapshot_backfill_info {
705            tracing::debug!(
706                ?snapshot_backfill_info,
707                "sending Command::CreateSnapshotBackfillStreamingJob"
708            );
709            CreateStreamingJobType::Independent {
710                snapshot_backfill_info,
711                kind: IndependentStreamingJobType::SnapshotBackfill {
712                    since_epoch: since_timestamp_epoch.map(|provided_since_epoch| SinceEpochInfo {
713                        provided_since_epoch,
714                        resolved: None,
715                    }),
716                },
717            }
718        } else {
719            if since_timestamp_epoch.is_some() {
720                bail!("since_timestamp should not be specified when no snapshot backfill");
721            }
722            tracing::debug!("sending Command::CreateStreamingJob");
723            if let Some(new_upstream_sink) = new_upstream_sink {
724                CreateStreamingJobType::SinkIntoTable(new_upstream_sink)
725            } else {
726                CreateStreamingJobType::Normal
727            }
728        };
729
730        let command = Command::CreateStreamingJob {
731            info,
732            job_type: create_job_type,
733            cross_db_snapshot_backfill_info,
734        };
735
736        self.barrier_scheduler
737            .run_command(streaming_job.database_id(), command)
738            .await?;
739
740        tracing::debug!(?streaming_job, "first barrier collected for stream job");
741
742        Ok(streaming_job)
743    }
744
745    /// Send replace job command to barrier scheduler.
746    pub async fn replace_stream_job(
747        &self,
748        new_fragments: StreamJobFragmentsToCreate,
749        ReplaceStreamJobContext {
750            old_fragments,
751            replace_upstream,
752            upstream_fragment_downstreams,
753            tmp_id,
754            streaming_job,
755            drop_table_connector_ctx,
756            auto_refresh_schema_sinks,
757            streaming_job_model,
758            database_resource_group,
759        }: ReplaceStreamJobContext,
760    ) -> MetaResult<()> {
761        // Phase 1: Gather fragment-level split information.
762        // For replace source with existing downstream, splits will be aligned
763        // in Phase 2 inside the barrier worker using actor-level no-shuffle produced by render_actors.
764        // For replace source with no downstream (or non-source), discover splits fresh.
765        let split_plan = if streaming_job.is_source() {
766            match self
767                .source_manager
768                .discover_splits_for_replace_source(&new_fragments, &replace_upstream)
769                .await?
770            {
771                Some(discovered) => ReplaceJobSplitPlan::Discovered(discovered),
772                None => ReplaceJobSplitPlan::AlignFromPrevious,
773            }
774        } else {
775            let discovered = self.source_manager.discover_splits(&new_fragments).await?;
776            ReplaceJobSplitPlan::Discovered(discovered)
777        };
778        tracing::info!("replace_stream_job - split plan: {:?}", split_plan);
779
780        self.barrier_scheduler
781            .run_command(
782                streaming_job.database_id(),
783                Command::ReplaceStreamJob(ReplaceStreamJobPlan {
784                    old_fragments,
785                    new_fragments,
786                    database_resource_group,
787                    replace_upstream,
788                    upstream_fragment_downstreams,
789                    split_plan,
790                    streaming_job,
791                    streaming_job_model,
792                    tmp_id,
793                    to_drop_state_table_ids: {
794                        if let Some(drop_table_connector_ctx) = &drop_table_connector_ctx {
795                            vec![drop_table_connector_ctx.to_remove_state_table_id]
796                        } else {
797                            Vec::new()
798                        }
799                    },
800                    auto_refresh_schema_sinks,
801                }),
802            )
803            .await?;
804
805        Ok(())
806    }
807
808    /// Drop streaming jobs by barrier manager, and clean up all related resources. The error will
809    /// be ignored because the recovery process will take over it in cleaning part. Check
810    /// [`Command::DropStreamingJobs`] for details.
811    pub async fn drop_streaming_jobs(
812        &self,
813        database_id: DatabaseId,
814        streaming_job_ids: Vec<JobId>,
815        state_table_ids: Vec<TableId>,
816        dropped_sink_fragment_by_targets: HashMap<FragmentId, Vec<FragmentId>>,
817    ) {
818        if !streaming_job_ids.is_empty() || !state_table_ids.is_empty() {
819            let cleanup_streaming_job_ids = streaming_job_ids.clone();
820            let cleanup_state_table_ids = state_table_ids.clone();
821            let run_result = self
822                .barrier_scheduler
823                .run_command(
824                    database_id,
825                    Command::DropStreamingJobs {
826                        streaming_job_ids: streaming_job_ids.into_iter().collect(),
827                        unregistered_state_table_ids: state_table_ids.iter().copied().collect(),
828                        dropped_sink_fragment_by_targets,
829                    },
830                )
831                .await;
832            let result = match run_result {
833                Ok(()) => {
834                    cleanup_dropped_streaming_jobs(
835                        &self.refresh_manager,
836                        &self.hummock_manager,
837                        &self.metadata_manager,
838                        cleanup_streaming_job_ids,
839                        cleanup_state_table_ids,
840                        "drop_streaming_jobs",
841                    )
842                    .await
843                }
844                Err(err) => Err(err),
845            };
846            let _ = result.inspect_err(|err| {
847                tracing::error!(error = ?err.as_report(), "failed to run drop command");
848            });
849        }
850    }
851
852    /// Cancel streaming jobs and return the canceled table ids.
853    /// 1. Send cancel message to stream jobs (via `cancel_jobs`).
854    /// 2. Send cancel message to recovered stream jobs (via `barrier_scheduler`).
855    ///
856    /// Cleanup of their state is handled by the caller after the drop command is collected.
857    pub async fn cancel_streaming_jobs(&self, job_ids: Vec<JobId>) -> MetaResult<Vec<JobId>> {
858        if job_ids.is_empty() {
859            return Ok(vec![]);
860        }
861
862        let _reschedule_job_lock = self.reschedule_lock_read_guard().await;
863        let (receivers, background_job_ids) = self.creating_job_info.cancel_jobs(job_ids).await?;
864
865        let futures = receivers.into_iter().map(|(id, receiver)| async move {
866            if let Ok(cancelled) = receiver.await
867                && cancelled
868            {
869                tracing::info!("canceled streaming job {id}");
870                Ok(id)
871            } else {
872                Err(MetaError::from(anyhow::anyhow!(
873                    "failed to cancel streaming job {id}"
874                )))
875            }
876        });
877        let mut cancelled_ids = join_all(futures)
878            .await
879            .into_iter()
880            .collect::<MetaResult<Vec<_>>>()?;
881
882        // NOTE(kwannoel): For background_job_ids stream jobs that not tracked in streaming manager,
883        // we can directly cancel them by running the barrier command.
884        let futures = background_job_ids.into_iter().map(|id| async move {
885            let abort_result = self
886                .metadata_manager
887                .catalog_controller
888                .try_abort_creating_streaming_job(id, true)
889                .await?;
890            self.iceberg_compaction_manager
891                .clear_maintenance_for_aborted_job(&abort_result);
892            let Some(cancel_info) = abort_result.cancel_info else {
893                return Ok(None);
894            };
895
896            if let Some(database_id) = abort_result.database_id {
897                self.barrier_scheduler
898                    .run_command(database_id, cancel_info.command)
899                    .await?;
900                cleanup_dropped_streaming_jobs(
901                    &self.refresh_manager,
902                    &self.hummock_manager,
903                    &self.metadata_manager,
904                    cancel_info.streaming_job_ids,
905                    cancel_info.state_table_ids,
906                    "cancel_streaming_job",
907                )
908                .await?;
909            }
910
911            tracing::info!(?id, "cancelled background streaming job");
912            Ok(Some(id))
913        });
914        let cancelled_recovered_ids = join_all(futures)
915            .await
916            .into_iter()
917            .collect::<MetaResult<Vec<_>>>()?;
918
919        cancelled_ids.extend(cancelled_recovered_ids.into_iter().flatten());
920        Ok(cancelled_ids)
921    }
922
923    pub(crate) async fn reschedule_streaming_job(
924        &self,
925        job_id: JobId,
926        policy: ReschedulePolicy,
927        deferred: bool,
928    ) -> MetaResult<()> {
929        let _reschedule_job_lock = self.reschedule_lock_write_guard().await;
930
931        let creating_jobs = self.metadata_manager.list_creating_jobs().await?;
932
933        if !creating_jobs.is_empty() {
934            let blocked_jobs = self
935                .metadata_manager
936                .collect_reschedule_blocked_jobs_for_creating_jobs(&creating_jobs, !deferred)
937                .await?;
938
939            if blocked_jobs.contains(&job_id) {
940                bail!(
941                    "Cannot alter the job {} because it is blocked by creating unreschedulable backfill jobs",
942                    job_id,
943                );
944            }
945        }
946
947        let commands = self
948            .scale_controller
949            .reschedule_inplace(HashMap::from([(job_id, policy)]))
950            .await?;
951
952        if !deferred {
953            let _source_pause_guard = self.source_manager.pause_tick().await;
954
955            for (database_id, command) in commands {
956                self.barrier_scheduler
957                    .run_command(database_id, command)
958                    .await?;
959            }
960        }
961
962        Ok(())
963    }
964
965    pub(crate) async fn reschedule_streaming_job_backfill_parallelism(
966        &self,
967        job_id: JobId,
968        parallelism: Option<ParallelismPolicy>,
969        deferred: bool,
970    ) -> MetaResult<()> {
971        let _reschedule_job_lock = self.reschedule_lock_write_guard().await;
972
973        if !deferred {
974            let creating_jobs = self.metadata_manager.list_creating_jobs().await?;
975
976            if !creating_jobs.is_empty() {
977                let jobs_with_unreschedulable_scan = self
978                    .metadata_manager
979                    .collect_online_unreschedulable_backfill_jobs(&creating_jobs)
980                    .await?;
981
982                if jobs_with_unreschedulable_scan.contains(&job_id) {
983                    bail!(
984                        "Cannot alter the job {} because its creating backfill contains a scan type that does not support online rescheduling",
985                        job_id,
986                    );
987                }
988            }
989        }
990
991        let commands = self
992            .scale_controller
993            .reschedule_backfill_parallelism_inplace(HashMap::from([(job_id, parallelism)]))
994            .await?;
995
996        if !deferred {
997            let _source_pause_guard = self.source_manager.pause_tick().await;
998
999            for (database_id, command) in commands {
1000                self.barrier_scheduler
1001                    .run_command(database_id, command)
1002                    .await?;
1003            }
1004        }
1005
1006        Ok(())
1007    }
1008
1009    /// This method is copied from `GlobalStreamManager::reschedule_streaming_job` and modified to handle reschedule CDC table backfill.
1010    pub(crate) async fn reschedule_cdc_table_backfill(
1011        &self,
1012        job_id: JobId,
1013        target: ReschedulePolicy,
1014    ) -> MetaResult<()> {
1015        let _reschedule_job_lock = self.reschedule_lock_write_guard().await;
1016
1017        let parallelism_policy = match target {
1018            ReschedulePolicy::Parallelism(policy)
1019                if matches!(policy.parallelism, StreamingParallelism::Fixed(_)) =>
1020            {
1021                policy
1022            }
1023            _ => bail_invalid_parameter!(
1024                "CDC backfill reschedule only supports fixed parallelism targets"
1025            ),
1026        };
1027
1028        let cdc_fragment_id = {
1029            let inner = self.metadata_manager.catalog_controller.inner.read().await;
1030            let fragments: Vec<(
1031                risingwave_meta_model::FragmentId,
1032                i32,
1033                risingwave_meta_model::StreamNode,
1034            )> = FragmentModel::find()
1035                .select_only()
1036                .columns([
1037                    fragment::Column::FragmentId,
1038                    fragment::Column::FragmentTypeMask,
1039                    fragment::Column::StreamNode,
1040                ])
1041                .filter(fragment::Column::JobId.eq(job_id))
1042                .into_tuple()
1043                .all(&inner.db)
1044                .await?;
1045
1046            let cdc_fragments = fragments
1047                .into_iter()
1048                .filter_map(|(fragment_id, mask, stream_node)| {
1049                    is_parallelized_backfill_enabled_cdc_scan_fragment(
1050                        FragmentTypeMask::from(mask),
1051                        &stream_node.to_protobuf(),
1052                    )
1053                    .is_some()
1054                    .then_some(fragment_id)
1055                })
1056                .collect_vec();
1057
1058            match cdc_fragments.len() {
1059                0 => bail_invalid_parameter!("no StreamCdcScan fragments found for job {}", job_id),
1060                1 => cdc_fragments[0],
1061                _ => bail_invalid_parameter!(
1062                    "multiple StreamCdcScan fragments found for job {}; expected exactly one",
1063                    job_id
1064                ),
1065            }
1066        };
1067
1068        let fragment_policy = HashMap::from([(
1069            cdc_fragment_id,
1070            Some(parallelism_policy.parallelism.clone()),
1071        )]);
1072
1073        let commands = self
1074            .scale_controller
1075            .reschedule_fragment_inplace(fragment_policy)
1076            .await?;
1077
1078        let _source_pause_guard = self.source_manager.pause_tick().await;
1079
1080        for (database_id, command) in commands {
1081            self.barrier_scheduler
1082                .run_command(database_id, command)
1083                .await?;
1084        }
1085
1086        Ok(())
1087    }
1088
1089    pub(crate) async fn reschedule_fragments(
1090        &self,
1091        fragment_targets: HashMap<FragmentId, Option<StreamingParallelism>>,
1092    ) -> MetaResult<()> {
1093        if fragment_targets.is_empty() {
1094            return Ok(());
1095        }
1096
1097        let _reschedule_job_lock = self.reschedule_lock_write_guard().await;
1098
1099        let fragment_policy = fragment_targets
1100            .into_iter()
1101            .map(|(fragment_id, parallelism)| (fragment_id as _, parallelism))
1102            .collect();
1103
1104        let commands = self
1105            .scale_controller
1106            .reschedule_fragment_inplace(fragment_policy)
1107            .await?;
1108
1109        let _source_pause_guard = self.source_manager.pause_tick().await;
1110
1111        for (database_id, command) in commands {
1112            self.barrier_scheduler
1113                .run_command(database_id, command)
1114                .await?;
1115        }
1116
1117        Ok(())
1118    }
1119
1120    // Don't need to add actor, just send a command
1121    pub async fn create_subscription(
1122        self: &Arc<Self>,
1123        subscription: &Subscription,
1124    ) -> MetaResult<()> {
1125        let command = Command::CreateSubscription {
1126            subscription_id: subscription.id,
1127            upstream_mv_table_id: subscription.dependent_table_id,
1128            retention_second: subscription.retention_seconds,
1129        };
1130
1131        tracing::debug!("sending Command::CreateSubscription");
1132        self.barrier_scheduler
1133            .run_command(subscription.database_id, command)
1134            .await?;
1135        Ok(())
1136    }
1137
1138    // Don't need to add actor, just send a command
1139    pub async fn drop_subscription(
1140        self: &Arc<Self>,
1141        database_id: DatabaseId,
1142        subscription_id: SubscriptionId,
1143        table_id: TableId,
1144    ) {
1145        let command = Command::DropSubscription {
1146            subscription_id,
1147            upstream_mv_table_id: table_id,
1148        };
1149
1150        tracing::debug!("sending Command::DropSubscriptions");
1151        let _ = self
1152            .barrier_scheduler
1153            .run_command(database_id, command)
1154            .await
1155            .inspect_err(|err| {
1156                tracing::error!(error = ?err.as_report(), "failed to run drop command");
1157            });
1158    }
1159
1160    pub async fn alter_subscription_retention(
1161        self: &Arc<Self>,
1162        database_id: DatabaseId,
1163        subscription_id: SubscriptionId,
1164        table_id: TableId,
1165        retention_second: u64,
1166    ) -> MetaResult<()> {
1167        let command = Command::AlterSubscriptionRetention {
1168            subscription_id,
1169            upstream_mv_table_id: table_id,
1170            retention_second,
1171        };
1172
1173        tracing::debug!("sending Command::AlterSubscriptionRetention");
1174        self.barrier_scheduler
1175            .run_command(database_id, command)
1176            .await?;
1177        Ok(())
1178    }
1179}