1use 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
70pub type CreateStreamingJobResult =
73 Result<NotificationVersion, (MetaError, bool, Option<oneshot::Sender<bool>>)>;
74
75const 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 }
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 pub sink_original_target_columns: Vec<PbColumnCatalog>,
118 pub project_exprs: Vec<PbExprNode>,
119 pub new_sink_downstream: DownstreamFragmentRelation,
120}
121
122pub struct CreateStreamingJobContext {
126 pub upstream_fragment_downstreams: FragmentDownstreamRelation,
128
129 pub database_resource_group: String,
131
132 pub definition: String,
134
135 pub create_type: CreateType,
136
137 pub job_type: StreamingJobType,
138
139 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 pub streaming_job_model: streaming_job::Model,
161
162 pub replace_sink: Option<SinkId>,
164
165 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 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 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
289pub struct ReplaceStreamJobContext {
293 pub old_fragments: StreamJobFragments,
295
296 pub replace_upstream: FragmentReplaceUpstream,
298
299 pub upstream_fragment_downstreams: FragmentDownstreamRelation,
301
302 pub streaming_job: StreamingJob,
303
304 pub database_resource_group: String,
306
307 pub tmp_id: JobId,
308
309 pub drop_table_connector_ctx: Option<DropTableConnectorContext>,
311
312 pub auto_refresh_schema_sinks: Option<Vec<AutoRefreshSchemaSinkContext>>,
313
314 pub streaming_job_model: streaming_job::Model,
316}
317
318pub struct GlobalStreamManager {
320 pub env: MetaSrvEnv,
321
322 pub metadata_manager: MetadataManager,
323
324 pub barrier_scheduler: BarrierScheduler,
326
327 pub hummock_manager: HummockManagerRef,
328
329 pub source_manager: SourceManagerRef,
331
332 pub refresh_manager: GlobalRefreshManagerRef,
333
334 pub iceberg_compaction_manager: IcebergCompactionManagerRef,
335
336 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 #[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 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 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 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 #[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 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 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 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 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 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 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 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 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 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 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}