1use std::cmp::Ordering;
16use std::collections::{HashMap, HashSet};
17use std::num::NonZeroUsize;
18use std::sync::Arc;
19use std::sync::atomic::AtomicU64;
20use std::time::Duration;
21
22use anyhow::{Context, anyhow};
23use await_tree::InstrumentAwait;
24use either::Either;
25use itertools::Itertools;
26use risingwave_common::catalog::{
27 AlterDatabaseParam, ColumnCatalog, ColumnId, Field, FragmentTypeFlag,
28};
29use risingwave_common::hash::VnodeCountCompat;
30use risingwave_common::id::{JobId, TableId};
31use risingwave_common::secret::{LocalSecretManager, SecretEncryption};
32use risingwave_common::system_param::adaptive_parallelism_strategy::parse_strategy;
33use risingwave_common::system_param::reader::SystemParamsRead;
34use risingwave_common::util::stream_graph_visitor::visit_stream_node_cont_mut;
35use risingwave_common::{bail, bail_not_implemented};
36use risingwave_connector::WithOptionsSecResolved;
37use risingwave_connector::connector_common::validate_connection;
38use risingwave_connector::sink::SinkParam;
39use risingwave_connector::sink::iceberg::IcebergSink;
40use risingwave_connector::source::cdc::CdcScanOptions;
41use risingwave_connector::source::{
42 ConnectorProperties, SourceEnumeratorContext, UPSTREAM_SOURCE_KEY,
43};
44use risingwave_meta_model::object::ObjectType;
45use risingwave_meta_model::refresh_job::RefreshState;
46use risingwave_meta_model::{
47 ConnectionId, DatabaseId, DispatcherType, FragmentId, FunctionId, IndexId, JobStatus, ObjectId,
48 SchemaId, SecretId, SinkId, SourceId, StreamingParallelism, SubscriptionId, UserId, ViewId,
49 streaming_job,
50};
51use risingwave_pb::catalog::{
52 Comment, Connection, CreateType, Database, Function, PbTable, Schema, Secret, Source,
53 Subscription, Table, View,
54};
55use risingwave_pb::ddl_service::alter_owner_request::Object;
56use risingwave_pb::ddl_service::{
57 DdlProgress, TableJobType, WaitVersion, alter_name_request, alter_set_schema_request,
58 alter_swap_rename_request, streaming_job_resource_type,
59};
60use risingwave_pb::meta::table_fragments::fragment::FragmentDistributionType as PbFragmentDistributionType;
61use risingwave_pb::plan_common::{PbColumnCatalog, PbExternalTableDesc};
62use risingwave_pb::stream_plan::stream_node::NodeBody;
63use risingwave_pb::stream_plan::{
64 PbDispatchOutputMapping, PbStreamFragmentGraph, PbStreamNode, PbUpstreamSinkInfo,
65 StreamFragmentGraph as StreamFragmentGraphProto,
66};
67use risingwave_pb::telemetry::{PbTelemetryDatabaseObject, PbTelemetryEventStage};
68use strum::Display;
69use thiserror_ext::AsReport;
70use tokio::sync::Semaphore;
71use tokio::time::sleep;
72use tracing::Instrument;
73
74use crate::barrier::{BarrierManagerRef, Command};
75use crate::controller::catalog::{DropTableConnectorContext, ReleaseContext};
76use crate::controller::streaming_job::{FinishAutoRefreshSchemaSinkContext, SinkIntoTableContext};
77use crate::controller::utils::build_select_node_list;
78use crate::error::{MetaErrorInner, bail_invalid_parameter};
79use crate::manager::iceberg_compaction::IcebergCompactionManagerRef;
80use crate::manager::iceberg_pk_index_sink::IcebergPkIndexSinkManager;
81use crate::manager::sink_coordination::SinkCoordinatorManager;
82use crate::manager::{
83 IGNORED_NOTIFICATION_VERSION, LocalNotification, MetaSrvEnv, MetadataManager,
84 NotificationVersion, StreamingJob, StreamingJobType,
85};
86use crate::model::{
87 DownstreamFragmentRelation, FragmentDownstreamRelation, FragmentId as CatalogFragmentId,
88 StreamContext, StreamJobFragments, StreamJobFragmentsToCreate,
89};
90use crate::stream::cdc::{
91 parallel_cdc_table_backfill_fragment, try_init_parallel_cdc_table_snapshot_splits,
92};
93use crate::stream::{
94 ActorGraphBuildResult, ActorGraphBuilder, AutoRefreshSchemaSinkContext,
95 CompleteStreamFragmentGraph, CreateStreamingJobContext, CreateStreamingJobOption,
96 FragmentGraphDownstreamContext, FragmentGraphUpstreamContext, GlobalStreamManagerRef,
97 ParallelismPolicy, ReplaceStreamJobContext, ReschedulePolicy, SourceChange, SourceManagerRef,
98 StreamFragmentGraph, UpstreamSinkInfo, check_sink_fragments_support_refresh_schema,
99 cleanup_dropped_streaming_jobs, create_source_worker, first_variant_column,
100 rewrite_refresh_schema_sink_fragment, state_match, validate_sink,
101};
102use crate::telemetry::report_event;
103use crate::{MetaError, MetaResult};
104
105#[derive(PartialEq)]
106pub enum DropMode {
107 Restrict,
108 Cascade,
109}
110
111impl DropMode {
112 pub fn from_request_setting(cascade: bool) -> DropMode {
113 if cascade {
114 DropMode::Cascade
115 } else {
116 DropMode::Restrict
117 }
118 }
119}
120
121#[derive(strum::AsRefStr)]
122pub enum StreamingJobId {
123 MaterializedView(TableId),
124 Sink(SinkId),
125 Table(Option<SourceId>, TableId),
126 Index(IndexId),
127}
128
129impl std::fmt::Display for StreamingJobId {
130 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
131 write!(f, "{}", self.as_ref())?;
132 write!(f, "({})", self.id())
133 }
134}
135
136impl StreamingJobId {
137 fn id(&self) -> JobId {
138 match self {
139 StreamingJobId::MaterializedView(id) | StreamingJobId::Table(_, id) => id.as_job_id(),
140 StreamingJobId::Index(id) => id.as_job_id(),
141 StreamingJobId::Sink(id) => id.as_job_id(),
142 }
143 }
144}
145
146pub struct ReplaceStreamJobInfo {
149 pub streaming_job: StreamingJob,
150 pub fragment_graph: StreamFragmentGraphProto,
151}
152
153#[derive(Display)]
154pub enum DdlCommand {
155 CreateDatabase(Database),
156 DropDatabase(DatabaseId),
157 CreateSchema(Schema),
158 DropSchema(SchemaId, DropMode),
159 CreateNonSharedSource(Source, Option<TableId>),
160 DropSource(SourceId, DropMode),
161 ResetSource(SourceId),
162 CreateFunction(Function),
163 DropFunction(FunctionId, DropMode),
164 CreateView(View, HashSet<ObjectId>),
165 DropView(ViewId, DropMode),
166 CreateStreamingJob {
167 stream_job: StreamingJob,
168 fragment_graph: StreamFragmentGraphProto,
169 dependencies: HashSet<ObjectId>,
170 resource_type: streaming_job_resource_type::ResourceType,
171 if_not_exists: bool,
172 refresh_interval_sec: Option<u64>,
173 replace_sink: Option<SinkId>,
174 since_timestamp_epoch: Option<u64>,
175 },
176 DropStreamingJob {
177 job_id: StreamingJobId,
178 drop_mode: DropMode,
179 },
180 AlterName(alter_name_request::Object, String),
181 AlterSwapRename(alter_swap_rename_request::Object),
182 ReplaceStreamJob(ReplaceStreamJobInfo),
183 AlterNonSharedSource(Source),
184 AlterObjectOwner(Object, UserId),
185 AlterSetSchema(alter_set_schema_request::Object, SchemaId),
186 CreateConnection(Connection),
187 DropConnection(ConnectionId, DropMode),
188 CreateSecret(Secret),
189 AlterSecret(Secret),
190 DropSecret(SecretId, DropMode),
191 CommentOn(Comment),
192 CreateSubscription(Subscription),
193 DropSubscription(SubscriptionId, DropMode),
194 AlterSubscriptionRetention {
195 subscription_id: SubscriptionId,
196 retention_seconds: u64,
197 definition: String,
198 },
199 AlterDatabaseParam(DatabaseId, AlterDatabaseParam),
200 AlterDatabaseResourceGroup(DatabaseId, Option<String>, bool),
201 AlterStreamingJobConfig(JobId, HashMap<String, String>, Vec<String>),
202}
203
204impl DdlCommand {
205 fn object(&self) -> Either<String, ObjectId> {
207 use Either::*;
208 match self {
209 DdlCommand::CreateDatabase(database) => Left(database.name.clone()),
210 DdlCommand::DropDatabase(id) => Right(id.as_object_id()),
211 DdlCommand::CreateSchema(schema) => Left(schema.name.clone()),
212 DdlCommand::DropSchema(id, _) => Right(id.as_object_id()),
213 DdlCommand::CreateNonSharedSource(source, _) => Left(source.name.clone()),
214 DdlCommand::DropSource(id, _) => Right(id.as_object_id()),
215 DdlCommand::ResetSource(id) => Right(id.as_object_id()),
216 DdlCommand::CreateFunction(function) => Left(function.name.clone()),
217 DdlCommand::DropFunction(id, _) => Right(id.as_object_id()),
218 DdlCommand::CreateView(view, _) => Left(view.name.clone()),
219 DdlCommand::DropView(id, _) => Right(id.as_object_id()),
220 DdlCommand::CreateStreamingJob { stream_job, .. } => Left(stream_job.name()),
221 DdlCommand::DropStreamingJob { job_id, .. } => Right(job_id.id().as_object_id()),
222 DdlCommand::AlterName(object, _) => Left(format!("{object:?}")),
223 DdlCommand::AlterSwapRename(object) => Left(format!("{object:?}")),
224 DdlCommand::ReplaceStreamJob(info) => Left(info.streaming_job.name()),
225 DdlCommand::AlterNonSharedSource(source) => Left(source.name.clone()),
226 DdlCommand::AlterObjectOwner(object, _) => Left(format!("{object:?}")),
227 DdlCommand::AlterSetSchema(object, _) => Left(format!("{object:?}")),
228 DdlCommand::CreateConnection(connection) => Left(connection.name.clone()),
229 DdlCommand::DropConnection(id, _) => Right(id.as_object_id()),
230 DdlCommand::CreateSecret(secret) => Left(secret.name.clone()),
231 DdlCommand::AlterSecret(secret) => Left(secret.name.clone()),
232 DdlCommand::DropSecret(id, _) => Right(id.as_object_id()),
233 DdlCommand::CommentOn(comment) => Right(comment.table_id.into()),
234 DdlCommand::CreateSubscription(subscription) => Left(subscription.name.clone()),
235 DdlCommand::DropSubscription(id, _) => Right(id.as_object_id()),
236 DdlCommand::AlterSubscriptionRetention {
237 subscription_id, ..
238 } => Right(subscription_id.as_object_id()),
239 DdlCommand::AlterDatabaseParam(id, _) => Right(id.as_object_id()),
240 DdlCommand::AlterDatabaseResourceGroup(id, _, _) => Right(id.as_object_id()),
241 DdlCommand::AlterStreamingJobConfig(job_id, _, _) => Right(job_id.as_object_id()),
242 }
243 }
244
245 fn allow_in_recovery(&self) -> bool {
246 match self {
247 DdlCommand::DropDatabase(_)
248 | DdlCommand::DropSchema(_, _)
249 | DdlCommand::DropSource(_, _)
250 | DdlCommand::DropFunction(_, _)
251 | DdlCommand::DropView(_, _)
252 | DdlCommand::DropStreamingJob { .. }
253 | DdlCommand::DropConnection(_, _)
254 | DdlCommand::DropSecret(_, _)
255 | DdlCommand::DropSubscription(_, _)
256 | DdlCommand::AlterName(_, _)
257 | DdlCommand::AlterObjectOwner(_, _)
258 | DdlCommand::AlterSetSchema(_, _)
259 | DdlCommand::CreateDatabase(_)
260 | DdlCommand::CreateSchema(_)
261 | DdlCommand::CreateFunction(_)
262 | DdlCommand::CreateView(_, _)
263 | DdlCommand::CreateConnection(_)
264 | DdlCommand::CommentOn(_)
265 | DdlCommand::CreateSecret(_)
266 | DdlCommand::AlterSecret(_)
267 | DdlCommand::AlterSwapRename(_)
268 | DdlCommand::AlterDatabaseParam(_, _)
269 | DdlCommand::AlterDatabaseResourceGroup(_, _, _)
270 | DdlCommand::AlterStreamingJobConfig(_, _, _)
271 | DdlCommand::AlterSubscriptionRetention { .. } => true,
272 DdlCommand::CreateStreamingJob { .. }
273 | DdlCommand::CreateNonSharedSource(_, _)
274 | DdlCommand::ReplaceStreamJob(_)
275 | DdlCommand::AlterNonSharedSource(_)
276 | DdlCommand::ResetSource(_)
277 | DdlCommand::CreateSubscription(_) => false,
278 }
279 }
280}
281
282#[derive(Clone)]
283pub struct DdlController {
284 pub(crate) env: MetaSrvEnv,
285
286 pub(crate) metadata_manager: MetadataManager,
287 pub(crate) stream_manager: GlobalStreamManagerRef,
288 pub(crate) source_manager: SourceManagerRef,
289 barrier_manager: BarrierManagerRef,
290 sink_manager: SinkCoordinatorManager,
291 iceberg_compaction_manager: IcebergCompactionManagerRef,
292 iceberg_pk_index_sink_manager: IcebergPkIndexSinkManager,
293
294 pub(crate) creating_streaming_job_permits: Arc<CreatingStreamingJobPermit>,
296
297 seq: Arc<AtomicU64>,
299}
300
301#[derive(Clone)]
302pub struct CreatingStreamingJobPermit {
303 pub(crate) semaphore: Arc<Semaphore>,
304}
305
306impl CreatingStreamingJobPermit {
307 async fn new(env: &MetaSrvEnv) -> Self {
308 let mut permits = env
309 .system_params_reader()
310 .await
311 .max_concurrent_creating_streaming_jobs() as usize;
312 if permits == 0 {
313 permits = Semaphore::MAX_PERMITS;
315 }
316 let semaphore = Arc::new(Semaphore::new(permits));
317
318 let (local_notification_tx, mut local_notification_rx) =
319 tokio::sync::mpsc::unbounded_channel();
320 env.notification_manager()
321 .insert_local_sender(local_notification_tx);
322 let semaphore_clone = semaphore.clone();
323 tokio::spawn(async move {
324 while let Some(notification) = local_notification_rx.recv().await {
325 let LocalNotification::SystemParamsChange(p) = ¬ification else {
326 continue;
327 };
328 let mut new_permits = p.max_concurrent_creating_streaming_jobs() as usize;
329 if new_permits == 0 {
330 new_permits = Semaphore::MAX_PERMITS;
331 }
332 match permits.cmp(&new_permits) {
333 Ordering::Less => {
334 semaphore_clone.add_permits(new_permits - permits);
335 }
336 Ordering::Equal => continue,
337 Ordering::Greater => {
338 let to_release = permits - new_permits;
339 let reduced = semaphore_clone.forget_permits(to_release);
340 if reduced != to_release {
342 tracing::warn!(
343 "no enough permits to release, expected {}, but reduced {}",
344 to_release,
345 reduced
346 );
347 }
348 }
349 }
350 tracing::info!(
351 "max_concurrent_creating_streaming_jobs changed from {} to {}",
352 permits,
353 new_permits
354 );
355 permits = new_permits;
356 }
357 });
358
359 Self { semaphore }
360 }
361}
362
363impl DdlController {
364 fn validate_specified_parallelism(
365 specified_parallelism: Option<NonZeroUsize>,
366 specified_backfill_parallelism: Option<NonZeroUsize>,
367 max_parallelism: NonZeroUsize,
368 ) -> MetaResult<()> {
369 if let Some(parallelism) = specified_parallelism
370 && parallelism > max_parallelism
371 {
372 bail_invalid_parameter!(
373 "specified parallelism {} should not exceed max parallelism {}",
374 parallelism,
375 max_parallelism,
376 );
377 }
378 if let Some(backfill_parallelism) = specified_backfill_parallelism
379 && backfill_parallelism > max_parallelism
380 {
381 bail_invalid_parameter!(
382 "specified backfill parallelism {} should not exceed max parallelism {}",
383 backfill_parallelism,
384 max_parallelism,
385 );
386 }
387 Ok(())
388 }
389
390 fn validate_serverless_backfill_enabled(
391 &self,
392 resource_type: &streaming_job_resource_type::ResourceType,
393 ) -> MetaResult<()> {
394 if matches!(
395 resource_type,
396 streaming_job_resource_type::ResourceType::ServerlessBackfill(true)
397 ) && self.env.opts.serverless_backfill_controller_addr.is_empty()
398 {
399 bail_invalid_parameter!(
400 "Serverless Backfill is disabled. Use RisingWave cloud at https://cloud.risingwave.com/auth/signup to try this feature"
401 );
402 }
403
404 Ok(())
405 }
406
407 pub async fn new(
408 env: MetaSrvEnv,
409 metadata_manager: MetadataManager,
410 stream_manager: GlobalStreamManagerRef,
411 source_manager: SourceManagerRef,
412 barrier_manager: BarrierManagerRef,
413 sink_manager: SinkCoordinatorManager,
414 iceberg_compaction_manager: IcebergCompactionManagerRef,
415 iceberg_pk_index_sink_manager: IcebergPkIndexSinkManager,
416 ) -> Self {
417 let creating_streaming_job_permits = Arc::new(CreatingStreamingJobPermit::new(&env).await);
418 Self {
419 env,
420 metadata_manager,
421 stream_manager,
422 source_manager,
423 barrier_manager,
424 sink_manager,
425 iceberg_compaction_manager,
426 iceberg_pk_index_sink_manager,
427 creating_streaming_job_permits,
428 seq: Arc::new(AtomicU64::new(0)),
429 }
430 }
431
432 pub fn next_seq(&self) -> u64 {
434 self.seq.fetch_add(1, std::sync::atomic::Ordering::Relaxed)
436 }
437
438 pub async fn run_command(&self, command: DdlCommand) -> MetaResult<Option<WaitVersion>> {
445 if !command.allow_in_recovery() {
446 self.barrier_manager.check_status_running()?;
447 }
448
449 let await_tree_key = format!("DDL Command {}", self.next_seq());
450 let await_tree_span = await_tree::span!("{command}({})", command.object());
451
452 let ctrl = self.clone();
453 let fut = Box::pin(async move {
454 match command {
455 DdlCommand::CreateDatabase(database) => ctrl.create_database(database).await,
456 DdlCommand::DropDatabase(database_id) => ctrl.drop_database(database_id).await,
457 DdlCommand::CreateSchema(schema) => ctrl.create_schema(schema).await,
458 DdlCommand::DropSchema(schema_id, drop_mode) => {
459 ctrl.drop_schema(schema_id, drop_mode).await
460 }
461 DdlCommand::CreateNonSharedSource(source, iceberg_table_id) => {
462 ctrl.create_non_shared_source(source, iceberg_table_id)
463 .await
464 }
465 DdlCommand::DropSource(source_id, drop_mode) => {
466 ctrl.drop_source(source_id, drop_mode).await
467 }
468 DdlCommand::ResetSource(source_id) => ctrl.reset_source(source_id).await,
469 DdlCommand::CreateFunction(function) => ctrl.create_function(function).await,
470 DdlCommand::DropFunction(function_id, drop_mode) => {
471 ctrl.drop_function(function_id, drop_mode).await
472 }
473 DdlCommand::CreateView(view, dependencies) => {
474 ctrl.create_view(view, dependencies).await
475 }
476 DdlCommand::DropView(view_id, drop_mode) => {
477 ctrl.drop_view(view_id, drop_mode).await
478 }
479 DdlCommand::CreateStreamingJob {
480 stream_job,
481 fragment_graph,
482 dependencies,
483 resource_type,
484 if_not_exists,
485 refresh_interval_sec,
486 replace_sink,
487 since_timestamp_epoch,
488 } => {
489 ctrl.create_streaming_job(
490 stream_job,
491 fragment_graph,
492 dependencies,
493 resource_type,
494 if_not_exists,
495 refresh_interval_sec,
496 replace_sink,
497 since_timestamp_epoch,
498 )
499 .await
500 }
501 DdlCommand::DropStreamingJob { job_id, drop_mode } => {
502 ctrl.drop_streaming_job(job_id, drop_mode).await
503 }
504 DdlCommand::ReplaceStreamJob(ReplaceStreamJobInfo {
505 streaming_job,
506 fragment_graph,
507 }) => ctrl.replace_job(streaming_job, fragment_graph).await,
508 DdlCommand::AlterName(relation, name) => ctrl.alter_name(relation, &name).await,
509 DdlCommand::AlterObjectOwner(object, owner_id) => {
510 ctrl.alter_owner(object, owner_id).await
511 }
512 DdlCommand::AlterSetSchema(object, new_schema_id) => {
513 ctrl.alter_set_schema(object, new_schema_id).await
514 }
515 DdlCommand::CreateConnection(connection) => {
516 ctrl.create_connection(connection).await
517 }
518 DdlCommand::DropConnection(connection_id, drop_mode) => {
519 ctrl.drop_connection(connection_id, drop_mode).await
520 }
521 DdlCommand::CreateSecret(secret) => ctrl.create_secret(secret).await,
522 DdlCommand::DropSecret(secret_id, drop_mode) => {
523 ctrl.drop_secret(secret_id, drop_mode).await
524 }
525 DdlCommand::AlterSecret(secret) => ctrl.alter_secret(secret).await,
526 DdlCommand::AlterNonSharedSource(source) => {
527 ctrl.alter_non_shared_source(source).await
528 }
529 DdlCommand::CommentOn(comment) => ctrl.comment_on(comment).await,
530 DdlCommand::CreateSubscription(subscription) => {
531 ctrl.create_subscription(subscription).await
532 }
533 DdlCommand::DropSubscription(subscription_id, drop_mode) => {
534 ctrl.drop_subscription(subscription_id, drop_mode).await
535 }
536 DdlCommand::AlterSubscriptionRetention {
537 subscription_id,
538 retention_seconds,
539 definition,
540 } => {
541 ctrl.alter_subscription_retention(
542 subscription_id,
543 retention_seconds,
544 definition,
545 )
546 .await
547 }
548 DdlCommand::AlterSwapRename(objects) => ctrl.alter_swap_rename(objects).await,
549 DdlCommand::AlterDatabaseParam(database_id, param) => {
550 ctrl.alter_database_param(database_id, param).await
551 }
552 DdlCommand::AlterDatabaseResourceGroup(database_id, resource_group, deferred) => {
553 ctrl.alter_database_resource_group(database_id, resource_group, deferred)
554 .await
555 }
556 DdlCommand::AlterStreamingJobConfig(job_id, entries_to_add, keys_to_remove) => {
557 ctrl.alter_streaming_job_config(job_id, entries_to_add, keys_to_remove)
558 .await
559 }
560 }
561 })
562 .in_current_span();
563 let fut = (self.env.await_tree_reg())
564 .register(await_tree_key, await_tree_span)
565 .instrument(Box::pin(fut));
566 let notification_version = tokio::spawn(fut).await.map_err(|e| anyhow!(e))??;
567 Ok(Some(WaitVersion {
568 catalog_version: notification_version,
569 hummock_version_id: self.barrier_manager.get_hummock_version_id().await,
570 }))
571 }
572
573 pub async fn get_ddl_progress(&self) -> MetaResult<Vec<DdlProgress>> {
574 self.barrier_manager.get_ddl_progress().await
575 }
576
577 async fn create_database(&self, database: Database) -> MetaResult<NotificationVersion> {
578 let (version, updated_db) = self
579 .metadata_manager
580 .catalog_controller
581 .create_database(database)
582 .await?;
583 self.barrier_manager
585 .update_database_barrier(
586 updated_db.database_id,
587 updated_db.barrier_interval_ms.map(|v| v as u32),
588 updated_db.checkpoint_frequency.map(|v| v as u64),
589 )
590 .await?;
591 Ok(version)
592 }
593
594 #[tracing::instrument(skip(self), level = "debug")]
595 pub async fn reschedule_streaming_job(
596 &self,
597 job_id: JobId,
598 target: ReschedulePolicy,
599 mut deferred: bool,
600 ) -> MetaResult<()> {
601 tracing::info!("altering parallelism for job {}", job_id);
602 if self.barrier_manager.check_status_running().is_err() {
603 tracing::info!(
604 "alter parallelism is set to deferred mode because the system is in recovery state"
605 );
606 deferred = true;
607 }
608
609 self.stream_manager
610 .reschedule_streaming_job(job_id, target, deferred)
611 .await
612 }
613
614 pub async fn reschedule_streaming_job_backfill_parallelism(
615 &self,
616 job_id: JobId,
617 parallelism: Option<ParallelismPolicy>,
618 mut deferred: bool,
619 ) -> MetaResult<()> {
620 tracing::info!("altering backfill parallelism for job {}", job_id);
621 if self.barrier_manager.check_status_running().is_err() {
622 tracing::info!(
623 "alter backfill parallelism is set to deferred mode because the system is in recovery state"
624 );
625 deferred = true;
626 }
627
628 self.stream_manager
629 .reschedule_streaming_job_backfill_parallelism(job_id, parallelism, deferred)
630 .await
631 }
632
633 pub async fn reschedule_cdc_table_backfill(
634 &self,
635 job_id: JobId,
636 target: ReschedulePolicy,
637 ) -> MetaResult<()> {
638 tracing::info!("alter CDC table backfill parallelism");
639 if self.barrier_manager.check_status_running().is_err() {
640 return Err(anyhow::anyhow!("CDC table backfill reschedule is unavailable because the system is in recovery state").into());
641 }
642 self.stream_manager
643 .reschedule_cdc_table_backfill(job_id, target)
644 .await
645 }
646
647 pub async fn reschedule_fragments(
648 &self,
649 fragment_targets: HashMap<FragmentId, Option<StreamingParallelism>>,
650 ) -> MetaResult<()> {
651 tracing::info!(
652 "altering parallelism for fragments {:?}",
653 fragment_targets.keys()
654 );
655 let fragment_targets = fragment_targets
656 .into_iter()
657 .map(|(fragment_id, parallelism)| (fragment_id as CatalogFragmentId, parallelism))
658 .collect();
659
660 self.stream_manager
661 .reschedule_fragments(fragment_targets)
662 .await
663 }
664
665 async fn drop_database(&self, database_id: DatabaseId) -> MetaResult<NotificationVersion> {
666 self.drop_object(ObjectType::Database, database_id, DropMode::Cascade)
667 .await
668 }
669
670 async fn create_schema(&self, schema: Schema) -> MetaResult<NotificationVersion> {
671 self.metadata_manager
672 .catalog_controller
673 .create_schema(schema)
674 .await
675 }
676
677 async fn drop_schema(
678 &self,
679 schema_id: SchemaId,
680 drop_mode: DropMode,
681 ) -> MetaResult<NotificationVersion> {
682 self.drop_object(ObjectType::Schema, schema_id, drop_mode)
683 .await
684 }
685
686 async fn create_non_shared_source(
688 &self,
689 source: Source,
690 iceberg_table_id: Option<TableId>,
691 ) -> MetaResult<NotificationVersion> {
692 if let Some(cdc_table_desc) = source
693 .info
694 .as_ref()
695 .and_then(|info| info.external_table.as_ref())
696 {
697 assert!(iceberg_table_id.is_none());
698 self.validate_cdc_table_desc(cdc_table_desc).await?;
702 let (_, version) = self
703 .metadata_manager
704 .catalog_controller
705 .create_source(source, iceberg_table_id)
706 .await?;
707 return Ok(version);
708 }
709
710 let handle = create_source_worker(
711 &source,
712 self.source_manager.metrics.clone(),
713 self.env.await_tree_reg().clone(),
714 )
715 .await
716 .context("failed to create source worker")?;
717
718 let (source_id, version) = self
719 .metadata_manager
720 .catalog_controller
721 .create_source(source, iceberg_table_id)
722 .await?;
723 self.source_manager
724 .register_source_with_handle(source_id, handle)
725 .await;
726 Ok(version)
727 }
728
729 async fn drop_source(
730 &self,
731 source_id: SourceId,
732 drop_mode: DropMode,
733 ) -> MetaResult<NotificationVersion> {
734 self.drop_object(ObjectType::Source, source_id, drop_mode)
735 .await
736 }
737
738 async fn reset_source(&self, source_id: SourceId) -> MetaResult<NotificationVersion> {
739 tracing::info!(source_id = %source_id, "resetting CDC source offset to latest");
740
741 let database_id = self
743 .metadata_manager
744 .catalog_controller
745 .get_object_database_id(source_id)
746 .await?;
747
748 self.stream_manager
749 .barrier_scheduler
750 .run_command(database_id, Command::ResetSource { source_id })
751 .await?;
752
753 let version = self
755 .metadata_manager
756 .catalog_controller
757 .notify_frontend_trivial()
758 .await;
759 Ok(version)
760 }
761
762 async fn alter_non_shared_source(&self, source: Source) -> MetaResult<NotificationVersion> {
765 self.metadata_manager
766 .catalog_controller
767 .alter_non_shared_source(source)
768 .await
769 }
770
771 async fn create_function(&self, function: Function) -> MetaResult<NotificationVersion> {
772 self.metadata_manager
773 .catalog_controller
774 .create_function(function)
775 .await
776 }
777
778 async fn drop_function(
779 &self,
780 function_id: FunctionId,
781 drop_mode: DropMode,
782 ) -> MetaResult<NotificationVersion> {
783 self.drop_object(ObjectType::Function, function_id, drop_mode)
784 .await
785 }
786
787 async fn create_view(
788 &self,
789 view: View,
790 dependencies: HashSet<ObjectId>,
791 ) -> MetaResult<NotificationVersion> {
792 self.metadata_manager
793 .catalog_controller
794 .create_view(view, dependencies)
795 .await
796 }
797
798 async fn drop_view(
799 &self,
800 view_id: ViewId,
801 drop_mode: DropMode,
802 ) -> MetaResult<NotificationVersion> {
803 self.drop_object(ObjectType::View, view_id, drop_mode).await
804 }
805
806 async fn create_connection(&self, connection: Connection) -> MetaResult<NotificationVersion> {
807 validate_connection(&connection).await?;
808 self.metadata_manager
809 .catalog_controller
810 .create_connection(connection)
811 .await
812 }
813
814 async fn drop_connection(
815 &self,
816 connection_id: ConnectionId,
817 drop_mode: DropMode,
818 ) -> MetaResult<NotificationVersion> {
819 self.drop_object(ObjectType::Connection, connection_id, drop_mode)
820 .await
821 }
822
823 async fn alter_database_param(
824 &self,
825 database_id: DatabaseId,
826 param: AlterDatabaseParam,
827 ) -> MetaResult<NotificationVersion> {
828 let (version, updated_db) = self
829 .metadata_manager
830 .catalog_controller
831 .alter_database_param(database_id, param)
832 .await?;
833 self.barrier_manager
835 .update_database_barrier(
836 database_id,
837 updated_db.barrier_interval_ms.map(|v| v as u32),
838 updated_db.checkpoint_frequency.map(|v| v as u64),
839 )
840 .await?;
841 Ok(version)
842 }
843
844 async fn alter_database_resource_group(
845 &self,
846 database_id: DatabaseId,
847 resource_group: Option<String>,
848 _deferred: bool,
849 ) -> MetaResult<NotificationVersion> {
850 let version = self
851 .metadata_manager
852 .catalog_controller
853 .alter_database_resource_group(database_id, resource_group)
854 .await?;
855
856 Ok(version)
857 }
858
859 fn get_encrypted_payload(&self, secret: &Secret) -> MetaResult<Vec<u8>> {
862 let secret_store_private_key = self
863 .env
864 .opts
865 .secret_store_private_key
866 .clone()
867 .ok_or_else(|| anyhow!("secret_store_private_key is not configured"))?;
868
869 let encrypted_payload = SecretEncryption::encrypt(
870 secret_store_private_key.as_slice(),
871 secret.get_value().as_slice(),
872 )
873 .context(format!("failed to encrypt secret {}", secret.name))?;
874 Ok(encrypted_payload
875 .serialize()
876 .context(format!("failed to serialize secret {}", secret.name))?)
877 }
878
879 async fn create_secret(&self, mut secret: Secret) -> MetaResult<NotificationVersion> {
880 let secret_plain_payload = secret.value.clone();
883 let encrypted_payload = self.get_encrypted_payload(&secret)?;
884 secret.value = encrypted_payload;
885
886 self.metadata_manager
887 .catalog_controller
888 .create_secret(secret, secret_plain_payload)
889 .await
890 }
891
892 async fn drop_secret(
893 &self,
894 secret_id: SecretId,
895 drop_mode: DropMode,
896 ) -> MetaResult<NotificationVersion> {
897 self.drop_object(ObjectType::Secret, secret_id, drop_mode)
898 .await
899 }
900
901 async fn alter_secret(&self, mut secret: Secret) -> MetaResult<NotificationVersion> {
902 let secret_plain_payload = secret.value.clone();
903 let encrypted_payload = self.get_encrypted_payload(&secret)?;
904 secret.value = encrypted_payload;
905 self.metadata_manager
906 .catalog_controller
907 .alter_secret(secret, secret_plain_payload)
908 .await
909 }
910
911 async fn create_subscription(
912 &self,
913 mut subscription: Subscription,
914 ) -> MetaResult<NotificationVersion> {
915 tracing::debug!("create subscription");
916 let _permit = self
917 .creating_streaming_job_permits
918 .semaphore
919 .acquire()
920 .await
921 .unwrap();
922 let _reschedule_job_lock = self.stream_manager.reschedule_lock_read_guard().await;
923 self.metadata_manager
924 .catalog_controller
925 .create_subscription_catalog(&mut subscription)
926 .await?;
927 if let Err(err) = self.stream_manager.create_subscription(&subscription).await {
928 tracing::debug!(error = %err.as_report(), "failed to create subscription");
929 let _ = self
930 .metadata_manager
931 .catalog_controller
932 .try_abort_creating_subscription(subscription.id)
933 .await
934 .inspect_err(|e| {
935 tracing::error!(
936 error = %e.as_report(),
937 "failed to abort create subscription after failure"
938 );
939 });
940 return Err(err);
941 }
942
943 let version = self
944 .metadata_manager
945 .catalog_controller
946 .notify_create_subscription(subscription.id)
947 .await?;
948 tracing::debug!("finish create subscription");
949 Ok(version)
950 }
951
952 async fn drop_subscription(
953 &self,
954 subscription_id: SubscriptionId,
955 drop_mode: DropMode,
956 ) -> MetaResult<NotificationVersion> {
957 tracing::debug!("preparing drop subscription");
958 let _reschedule_job_lock = self.stream_manager.reschedule_lock_read_guard().await;
959 let subscription = self
960 .metadata_manager
961 .catalog_controller
962 .get_subscription_by_id(subscription_id)
963 .await?;
964 let table_id = subscription.dependent_table_id;
965 let database_id = subscription.database_id;
966 let (_, version) = self
967 .metadata_manager
968 .catalog_controller
969 .drop_object(ObjectType::Subscription, subscription_id, drop_mode)
970 .await?;
971 self.stream_manager
972 .drop_subscription(database_id, subscription_id, table_id)
973 .await;
974 tracing::debug!("finish drop subscription");
975 Ok(version)
976 }
977
978 async fn alter_subscription_retention(
979 &self,
980 subscription_id: SubscriptionId,
981 retention_seconds: u64,
982 definition: String,
983 ) -> MetaResult<NotificationVersion> {
984 tracing::debug!("alter subscription retention");
985 let _reschedule_job_lock = self.stream_manager.reschedule_lock_read_guard().await;
986 let (version, subscription) = self
987 .metadata_manager
988 .catalog_controller
989 .alter_subscription_retention(subscription_id, retention_seconds, definition)
990 .await?;
991 self.stream_manager
992 .alter_subscription_retention(
993 subscription.database_id,
994 subscription.id,
995 subscription.dependent_table_id,
996 subscription.retention_seconds,
997 )
998 .await?;
999 tracing::debug!("finish alter subscription retention");
1000 Ok(version)
1001 }
1002
1003 #[await_tree::instrument]
1005 pub(crate) async fn validate_cdc_table(
1006 &self,
1007 table: &Table,
1008 table_fragments: &StreamJobFragments,
1009 ) -> MetaResult<()> {
1010 let stream_scan_fragment =
1011 Itertools::exactly_one(table_fragments.fragments.values().filter(|f| {
1012 f.fragment_type_mask.contains(FragmentTypeFlag::StreamScan)
1013 || f.fragment_type_mask
1014 .contains(FragmentTypeFlag::StreamCdcScan)
1015 }))
1016 .ok()
1017 .with_context(|| {
1018 format!(
1019 "expect exactly one stream scan fragment, got: {:?}",
1020 table_fragments.fragments
1021 )
1022 })?;
1023 fn assert_parallelism(
1024 distribution_type: PbFragmentDistributionType,
1025 node_body: &Option<NodeBody>,
1026 ) {
1027 if let Some(NodeBody::StreamCdcScan(node)) = node_body {
1028 if let Some(o) = node.options
1029 && CdcScanOptions::from_proto(&o).is_parallelized_backfill()
1030 {
1031 } else {
1033 assert_eq!(
1034 distribution_type,
1035 PbFragmentDistributionType::Single,
1036 "Non-parallelized CDC scan fragment should have Single distribution"
1037 );
1038 }
1039 }
1040 }
1041 let mut found_cdc_scan = false;
1042 match &stream_scan_fragment.nodes.node_body {
1043 Some(NodeBody::StreamCdcScan(_)) => {
1044 assert_parallelism(
1045 stream_scan_fragment.distribution_type,
1046 &stream_scan_fragment.nodes.node_body,
1047 );
1048 if self
1049 .validate_cdc_table_inner(&stream_scan_fragment.nodes.node_body, table.id)
1050 .await?
1051 {
1052 found_cdc_scan = true;
1053 }
1054 }
1055 Some(NodeBody::Project(_)) => {
1057 for input in &stream_scan_fragment.nodes.input {
1058 assert_parallelism(stream_scan_fragment.distribution_type, &input.node_body);
1059 if self
1060 .validate_cdc_table_inner(&input.node_body, table.id)
1061 .await?
1062 {
1063 found_cdc_scan = true;
1064 }
1065 }
1066 }
1067 _ => {
1068 bail!("Unexpected node body for stream cdc scan");
1069 }
1070 };
1071 if !found_cdc_scan {
1072 bail!("No stream cdc scan node found in stream scan fragment");
1073 }
1074 Ok(())
1075 }
1076
1077 async fn validate_cdc_table_inner(
1078 &self,
1079 node_body: &Option<NodeBody>,
1080 table_id: TableId,
1081 ) -> MetaResult<bool> {
1082 if let Some(NodeBody::StreamCdcScan(stream_cdc_scan)) = node_body
1083 && let Some(ref cdc_table_desc) = stream_cdc_scan.cdc_table_desc
1084 {
1085 self.validate_cdc_table_desc(cdc_table_desc).await?;
1086 tracing::debug!(?table_id, "validate cdc table success");
1087 Ok(true)
1088 } else {
1089 Ok(false)
1090 }
1091 }
1092
1093 pub(crate) async fn validate_cdc_table_desc(
1097 &self,
1098 cdc_table_desc: &PbExternalTableDesc,
1099 ) -> MetaResult<()> {
1100 let options_with_secret = WithOptionsSecResolved::new(
1101 cdc_table_desc.connect_properties.clone(),
1102 cdc_table_desc.secret_refs.clone(),
1103 );
1104
1105 let mut props = ConnectorProperties::extract(options_with_secret, true)?;
1106 props.init_from_pb_cdc_table_desc(cdc_table_desc);
1107
1108 let _enumerator = props
1109 .create_split_enumerator(SourceEnumeratorContext::dummy().into())
1110 .await?;
1111
1112 Ok(())
1113 }
1114
1115 pub async fn validate_table_for_sink(&self, table_id: TableId) -> MetaResult<()> {
1116 let migrated = self
1117 .metadata_manager
1118 .catalog_controller
1119 .has_table_been_migrated(table_id)
1120 .await?;
1121 if !migrated {
1122 Err(anyhow::anyhow!("Creating sink into table is not allowed for unmigrated table {}. Please migrate it first.", table_id).into())
1123 } else {
1124 Ok(())
1125 }
1126 }
1127
1128 #[await_tree::instrument(boxed, "create_streaming_job({streaming_job})")]
1131 pub async fn create_streaming_job(
1132 &self,
1133 mut streaming_job: StreamingJob,
1134 fragment_graph: StreamFragmentGraphProto,
1135 dependencies: HashSet<ObjectId>,
1136 resource_type: streaming_job_resource_type::ResourceType,
1137 if_not_exists: bool,
1138 refresh_interval_sec: Option<u64>,
1139 replace_sink: Option<SinkId>,
1140 since_timestamp_epoch: Option<u64>,
1141 ) -> MetaResult<NotificationVersion> {
1142 let replace_sink_info = if let Some(old_sink_id) = replace_sink {
1143 let StreamingJob::Sink(sink, _) = &streaming_job else {
1144 bail!("replace sink requires a sink job")
1145 };
1146 if sink.target_table.is_some() {
1147 bail_not_implemented!("replace sink into table")
1148 }
1149
1150 Some(old_sink_id)
1151 } else {
1152 if let StreamingJob::Sink(sink, _) = &streaming_job
1153 && let Some(target_table) = sink.target_table
1154 {
1155 self.validate_table_for_sink(target_table).await?;
1156 }
1157 None
1158 };
1159 self.validate_serverless_backfill_enabled(&resource_type)?;
1160 let ctx = StreamContext::from_protobuf(fragment_graph.get_ctx().unwrap());
1161 let adaptive_parallelism_strategy =
1162 (!fragment_graph.adaptive_parallelism_strategy.is_empty()).then(|| {
1163 parse_strategy(&fragment_graph.adaptive_parallelism_strategy)
1164 .expect("adaptive parallelism strategy should be validated in frontend")
1165 });
1166 let backfill_adaptive_parallelism_strategy = (!fragment_graph
1167 .backfill_adaptive_parallelism_strategy
1168 .is_empty())
1169 .then(|| {
1170 parse_strategy(&fragment_graph.backfill_adaptive_parallelism_strategy)
1171 .expect("backfill adaptive parallelism strategy should be validated in frontend")
1172 });
1173
1174 let streaming_job_model = match self
1175 .metadata_manager
1176 .catalog_controller
1177 .create_job_catalog(
1178 &mut streaming_job,
1179 &ctx,
1180 &fragment_graph.parallelism,
1181 fragment_graph.max_parallelism as _,
1182 dependencies,
1183 resource_type.clone(),
1184 &fragment_graph.backfill_parallelism,
1185 adaptive_parallelism_strategy,
1186 backfill_adaptive_parallelism_strategy,
1187 replace_sink_info.as_ref(),
1188 refresh_interval_sec,
1189 )
1190 .await
1191 {
1192 Ok(model) => model,
1193 Err(meta_err) => {
1194 if !if_not_exists {
1195 return Err(meta_err);
1196 }
1197 return if let MetaErrorInner::Duplicated(_, _, Some(job_id)) = meta_err.inner() {
1198 if streaming_job.create_type() == CreateType::Foreground {
1199 let database_id = streaming_job.database_id();
1200 self.metadata_manager
1201 .wait_streaming_job_finished(database_id, *job_id)
1202 .await
1203 } else {
1204 Ok(IGNORED_NOTIFICATION_VERSION)
1205 }
1206 } else {
1207 Err(meta_err)
1208 };
1209 }
1210 };
1211 let job_id = streaming_job.id();
1212 if let Some(old_sink_id) = replace_sink_info.as_ref() {
1213 tracing::debug!(
1214 old_sink_id = %old_sink_id,
1215 new_sink_id = %job_id,
1216 definition = streaming_job.definition(),
1217 create_type = streaming_job.create_type().as_str_name(),
1218 "starting replacement sink",
1219 );
1220 } else {
1221 tracing::debug!(
1222 id = %job_id,
1223 definition = streaming_job.definition(),
1224 create_type = streaming_job.create_type().as_str_name(),
1225 job_type = ?streaming_job.job_type(),
1226 "starting streaming job",
1227 );
1228 }
1229 let permit = self
1231 .creating_streaming_job_permits
1232 .semaphore
1233 .clone()
1234 .acquire_owned()
1235 .instrument_await("acquire_creating_streaming_job_permit")
1236 .await
1237 .unwrap();
1238 let reschedule_job_lock = self.stream_manager.reschedule_lock_read_guard().await;
1239
1240 let name = streaming_job.name();
1241 let definition = streaming_job.definition();
1242 let database_id = streaming_job.database_id();
1243 let source_id = match &streaming_job {
1244 StreamingJob::Table(Some(src), _, _) | StreamingJob::Source(src) => Some(src.id),
1245 _ => None,
1246 };
1247 let create_result = match self
1248 .generate_streaming_job(
1249 ctx,
1250 streaming_job,
1251 fragment_graph,
1252 resource_type.clone(),
1253 streaming_job_model,
1254 replace_sink_info,
1255 since_timestamp_epoch,
1256 )
1257 .await
1258 {
1259 Ok((stream_job_fragments, ctx)) => {
1260 self.stream_manager
1261 .create_streaming_job(stream_job_fragments, ctx, permit, reschedule_job_lock)
1262 .await
1263 }
1264 Err(err) => Err((err, false, None)),
1265 };
1266
1267 match create_result {
1268 Ok(version) => Ok(version),
1269 Err((err, is_cancelled, cancel_notifier)) => {
1270 tracing::error!(id = %job_id, error = %err.as_report(), "failed to create streaming job");
1271 let event = risingwave_pb::meta::event_log::EventCreateStreamJobFail {
1272 id: job_id,
1273 name,
1274 definition,
1275 error: err.as_report().to_string(),
1276 };
1277 self.env.event_log_manager_ref().add_event_logs(vec![
1278 risingwave_pb::meta::event_log::Event::CreateStreamJobFail(event),
1279 ]);
1280 let abort_result = self
1281 .metadata_manager
1282 .catalog_controller
1283 .try_abort_creating_streaming_job(job_id, is_cancelled)
1284 .await?;
1285 self.iceberg_compaction_manager
1286 .clear_maintenance_for_aborted_job(&abort_result);
1287 if let Some(cancel_info) = abort_result.cancel_info {
1288 self.stream_manager
1289 .barrier_scheduler
1290 .run_command(database_id, cancel_info.command)
1291 .await?;
1292 cleanup_dropped_streaming_jobs(
1293 &self.stream_manager.refresh_manager,
1294 &self.stream_manager.hummock_manager,
1295 &self.stream_manager.metadata_manager,
1296 cancel_info.streaming_job_ids,
1297 cancel_info.state_table_ids,
1298 "cancel_streaming_job",
1299 )
1300 .await?;
1301 }
1302 if let Some(cancel_notifier) = cancel_notifier {
1303 let _ = cancel_notifier.send(true).inspect_err(|err| {
1304 tracing::warn!("failed to notify cancellation result: {err}")
1305 });
1306 }
1307 if abort_result.aborted {
1308 tracing::warn!(id = %job_id, is_cancelled, "aborted streaming job");
1309 if let Some(source_id) = source_id {
1311 self.source_manager
1312 .apply_source_change(SourceChange::DropSource {
1313 dropped_source_ids: vec![source_id],
1314 })
1315 .await;
1316 }
1317 }
1318 Err(err)
1319 }
1320 }
1321 }
1322
1323 #[await_tree::instrument(boxed)]
1324 async fn generate_streaming_job(
1325 &self,
1326 ctx: StreamContext,
1327 mut streaming_job: StreamingJob,
1328 fragment_graph: StreamFragmentGraphProto,
1329 resource_type: streaming_job_resource_type::ResourceType,
1330 streaming_job_model: streaming_job::Model,
1331 replace_sink: Option<SinkId>,
1332 since_timestamp_epoch: Option<u64>,
1333 ) -> MetaResult<(StreamJobFragmentsToCreate, CreateStreamingJobContext)> {
1334 let mut fragment_graph =
1335 StreamFragmentGraph::new(&self.env, fragment_graph, &streaming_job)?;
1336 streaming_job.set_info_from_graph(&fragment_graph);
1337
1338 let incomplete_internal_tables = fragment_graph
1340 .incomplete_internal_tables()
1341 .into_values()
1342 .collect_vec();
1343 let table_id_map = self
1344 .metadata_manager
1345 .catalog_controller
1346 .create_internal_table_catalog(&streaming_job, incomplete_internal_tables)
1347 .await?;
1348 fragment_graph.refill_internal_table_ids(table_id_map);
1349
1350 tracing::debug!(id = %streaming_job.id(), "building streaming job");
1352 let (mut ctx, stream_job_fragments) = self
1353 .build_stream_job(
1354 ctx,
1355 streaming_job,
1356 fragment_graph,
1357 resource_type,
1358 streaming_job_model,
1359 since_timestamp_epoch,
1360 )
1361 .await?;
1362 ctx.replace_sink = replace_sink;
1363
1364 let streaming_job = &ctx.streaming_job;
1365
1366 match streaming_job {
1367 StreamingJob::Table(None, table, TableJobType::SharedCdcSource) => {
1368 self.validate_cdc_table(table, &stream_job_fragments)
1369 .await?;
1370 }
1371 StreamingJob::Table(Some(source), ..) => {
1372 self.source_manager.register_source(source).await?;
1374 let connector_name = source
1375 .get_with_properties()
1376 .get(UPSTREAM_SOURCE_KEY)
1377 .cloned();
1378 let attr = source.info.as_ref().map(|source_info| {
1379 jsonbb::json!({
1380 "format": source_info.format().as_str_name(),
1381 "encode": source_info.row_encode().as_str_name(),
1382 })
1383 });
1384 report_create_object(
1385 streaming_job.id(),
1386 "source",
1387 PbTelemetryDatabaseObject::Source,
1388 connector_name,
1389 attr,
1390 );
1391 }
1392 StreamingJob::Sink(sink, _) => {
1393 if sink.auto_refresh_schema_from_table.is_some() {
1394 check_sink_fragments_support_refresh_schema(&stream_job_fragments.fragments)?;
1395 }
1396 validate_sink(sink).await?;
1398 if crate::manager::iceberg_pk_index_sink::is_iceberg_pk_index_sink(&sink.properties)
1403 {
1404 let iceberg_config =
1405 crate::manager::iceberg_pk_index_sink::build_iceberg_config(sink)?;
1406 self.iceberg_pk_index_sink_manager
1407 .register_sink(
1408 sink.id,
1409 crate::barrier::to_partial_graph_id(sink.database_id, None),
1410 iceberg_config,
1411 )
1412 .await
1413 .map_err(|e| anyhow!(e).context("register v3 sink worker"))?;
1414 }
1415 let connector_name = sink.get_properties().get(UPSTREAM_SOURCE_KEY).cloned();
1416 let attr = sink.format_desc.as_ref().map(|sink_info| {
1417 jsonbb::json!({
1418 "format": sink_info.format().as_str_name(),
1419 "encode": sink_info.encode().as_str_name(),
1420 })
1421 });
1422 report_create_object(
1423 streaming_job.id(),
1424 "sink",
1425 PbTelemetryDatabaseObject::Sink,
1426 connector_name,
1427 attr,
1428 );
1429 }
1430 StreamingJob::Source(source) => {
1431 self.source_manager.register_source(source).await?;
1433 let connector_name = source
1434 .get_with_properties()
1435 .get(UPSTREAM_SOURCE_KEY)
1436 .cloned();
1437 let attr = source.info.as_ref().map(|source_info| {
1438 jsonbb::json!({
1439 "format": source_info.format().as_str_name(),
1440 "encode": source_info.row_encode().as_str_name(),
1441 })
1442 });
1443 report_create_object(
1444 streaming_job.id(),
1445 "source",
1446 PbTelemetryDatabaseObject::Source,
1447 connector_name,
1448 attr,
1449 );
1450 }
1451 _ => {}
1452 }
1453
1454 let backfill_orders = ctx.fragment_backfill_ordering.to_meta_model();
1455 self.metadata_manager
1456 .catalog_controller
1457 .prepare_stream_job_fragments(
1458 &stream_job_fragments,
1459 streaming_job,
1460 false,
1461 Some(backfill_orders),
1462 )
1463 .await?;
1464
1465 Ok((stream_job_fragments, ctx))
1466 }
1467
1468 pub async fn drop_object(
1470 &self,
1471 object_type: ObjectType,
1472 object_id: impl Into<ObjectId>,
1473 drop_mode: DropMode,
1474 ) -> MetaResult<NotificationVersion> {
1475 let object_id = object_id.into();
1476 let _reschedule_job_lock = self.stream_manager.reschedule_lock_read_guard().await;
1479 let _source_tick_pause_guard = self.source_manager.pause_tick().await;
1480
1481 let (release_ctx, version) = self
1482 .metadata_manager
1483 .catalog_controller
1484 .drop_object(object_type, object_id, drop_mode)
1485 .await?;
1486
1487 if object_type == ObjectType::Source {
1488 self.env
1489 .notification_manager_ref()
1490 .notify_local_subscribers(LocalNotification::SourceDropped(object_id));
1491 }
1492
1493 let ReleaseContext {
1494 database_id,
1495 removed_streaming_job_ids,
1496 removed_state_table_ids,
1497 removed_source_ids,
1498 removed_secret_ids: secret_ids,
1499 removed_source_fragments,
1500 removed_fragments,
1501 removed_sink_fragment_by_targets,
1502 removed_iceberg_table_sinks,
1503 removed_iceberg_sink_ids,
1504 removed_iceberg_pk_index_sink_ids,
1505 } = release_ctx;
1506 let removed_job_ids_for_sink_coordinators = removed_streaming_job_ids.clone();
1507
1508 self.env
1512 .notification_manager_ref()
1513 .notify_serving_fragment_mapping_delete(
1514 removed_fragments.iter().map(|id| *id as _).collect(),
1515 );
1516
1517 self.stream_manager
1518 .drop_streaming_jobs(
1519 database_id,
1520 removed_streaming_job_ids,
1521 removed_state_table_ids,
1522 removed_sink_fragment_by_targets
1523 .into_iter()
1524 .map(|(target, sinks)| {
1525 (target as _, sinks.into_iter().map(|id| id as _).collect())
1526 })
1527 .collect(),
1528 )
1529 .await;
1530
1531 self.source_manager
1534 .apply_source_change(SourceChange::DropSource {
1535 dropped_source_ids: removed_source_ids.into_iter().map(|id| id as _).collect(),
1536 })
1537 .await;
1538
1539 let dropped_source_fragments = removed_source_fragments;
1542 self.source_manager
1543 .apply_source_change(SourceChange::DropMv {
1544 dropped_source_fragments,
1545 })
1546 .await;
1547
1548 for sink in removed_iceberg_table_sinks {
1550 let sink_param = SinkParam::try_from_sink_catalog(sink.into())
1551 .expect("Iceberg sink should be valid");
1552 let iceberg_sink =
1553 IcebergSink::try_from(sink_param).expect("Iceberg sink should be valid");
1554 if let Ok(iceberg_catalog) = iceberg_sink.config.create_catalog().await {
1555 let table_identifier = iceberg_sink.config.full_table_name().unwrap();
1556 tracing::info!(
1557 "dropping iceberg table {} for dropped sink",
1558 table_identifier
1559 );
1560
1561 let _ = iceberg_catalog
1562 .drop_table(&table_identifier)
1563 .await
1564 .inspect_err(|err| {
1565 tracing::error!(
1566 "failed to drop iceberg table {} during cleanup: {}",
1567 table_identifier,
1568 err.as_report()
1569 );
1570 });
1571 }
1572 }
1573
1574 if !removed_job_ids_for_sink_coordinators.is_empty() {
1576 self.sink_manager
1577 .stop_sink_coordinators_for_jobs(removed_job_ids_for_sink_coordinators)
1578 .await;
1579 }
1580
1581 for sink_id in removed_iceberg_sink_ids {
1584 self.iceberg_compaction_manager
1585 .clear_iceberg_maintenance_by_sink_id(sink_id);
1586 }
1587
1588 if !removed_iceberg_pk_index_sink_ids.is_empty() {
1592 self.iceberg_pk_index_sink_manager.unregister_jobs(
1593 removed_iceberg_pk_index_sink_ids
1594 .into_iter()
1595 .map(|sink_id| sink_id.as_job_id()),
1596 );
1597 }
1598
1599 for secret in secret_ids {
1601 LocalSecretManager::global().remove_secret(secret);
1602 }
1603 Ok(version)
1604 }
1605
1606 #[await_tree::instrument(boxed, "replace_streaming_job({streaming_job})")]
1608 pub async fn replace_job(
1609 &self,
1610 mut streaming_job: StreamingJob,
1611 fragment_graph: StreamFragmentGraphProto,
1612 ) -> MetaResult<NotificationVersion> {
1613 match &streaming_job {
1614 StreamingJob::Table(..)
1615 | StreamingJob::Source(..)
1616 | StreamingJob::MaterializedView(..) => {}
1617 StreamingJob::Sink(..) | StreamingJob::Index(..) => {
1618 bail_not_implemented!("schema change for {}", streaming_job.job_type_str())
1619 }
1620 }
1621
1622 let job_id = streaming_job.id();
1623
1624 let _reschedule_job_lock = self.stream_manager.reschedule_lock_write_guard().await;
1625 if let StreamingJob::Table(_, table, _) = &streaming_job
1626 && self
1627 .metadata_manager
1628 .catalog_controller
1629 .get_refresh_job_state(table.id)
1630 .await?
1631 .is_some_and(|state| state != RefreshState::Idle)
1632 {
1633 bail!(
1634 "Cannot alter table {} because it is being refreshed",
1635 table.name
1636 );
1637 }
1638 let ctx = StreamContext::from_protobuf(fragment_graph.get_ctx().unwrap());
1639
1640 let original_max_parallelism = self
1642 .metadata_manager
1643 .get_job_max_parallelism(streaming_job.id())
1644 .await?;
1645 let fragment_graph = PbStreamFragmentGraph {
1646 max_parallelism: original_max_parallelism as _,
1647 ..fragment_graph
1648 };
1649
1650 let fragment_graph = StreamFragmentGraph::new(&self.env, fragment_graph, &streaming_job)?;
1652 streaming_job.set_info_from_graph(&fragment_graph);
1653
1654 let streaming_job = streaming_job;
1656
1657 let auto_refresh_schema_sinks = if let StreamingJob::Table(_, table, _) = &streaming_job {
1658 let auto_refresh_schema_sinks = self
1659 .metadata_manager
1660 .catalog_controller
1661 .get_sink_auto_refresh_schema_from(table.id)
1662 .await?;
1663 if !auto_refresh_schema_sinks.is_empty() {
1664 let original_table_columns = self
1665 .metadata_manager
1666 .catalog_controller
1667 .get_table_columns(table.id)
1668 .await?;
1669 let original_table_column_ids: HashSet<_> = original_table_columns
1671 .iter()
1672 .map(|col| col.column_id())
1673 .collect();
1674 let new_table_column_ids: HashSet<_> = table
1675 .columns
1676 .iter()
1677 .map(|col| ColumnId::new(col.column_desc.as_ref().unwrap().column_id as _))
1678 .collect();
1679 let newly_added_columns = table
1680 .columns
1681 .iter()
1682 .filter(|col| {
1683 !original_table_column_ids.contains(&ColumnId::new(
1684 col.column_desc.as_ref().unwrap().column_id as _,
1685 ))
1686 })
1687 .map(|col| ColumnCatalog::from(col.clone()))
1688 .collect_vec();
1689 let removed_columns = original_table_columns
1690 .iter()
1691 .filter(|col| !new_table_column_ids.contains(&col.column_id()))
1692 .cloned()
1693 .collect_vec();
1694 if let Some(variant_column) = first_variant_column(&newly_added_columns) {
1697 let sink_names = auto_refresh_schema_sinks
1698 .iter()
1699 .map(|sink| format!("`{}`", sink.name))
1700 .join(", ");
1701 return Err(MetaError::invalid_parameter(format!(
1702 "cannot add VARIANT column `{}` because sink(s) {} with auto schema refresh do not support VARIANT",
1703 variant_column.name_with_hidden(),
1704 sink_names,
1705 )));
1706 }
1707 let mut sinks = Vec::with_capacity(auto_refresh_schema_sinks.len());
1708 for sink in auto_refresh_schema_sinks {
1709 let sink_job_fragments = self
1710 .metadata_manager
1711 .get_job_fragments_by_id(sink.id.as_job_id())
1712 .await?;
1713 if sink_job_fragments.fragments.len() != 1 {
1714 return Err(anyhow!(
1715 "auto schema refresh sink must have only one fragment, but got {}",
1716 sink_job_fragments.fragments.len()
1717 )
1718 .into());
1719 }
1720 let sink_ctx = sink_job_fragments.ctx;
1721 let original_sink_fragment =
1722 sink_job_fragments.fragments.into_values().next().unwrap();
1723 let (new_sink_fragment, new_schema, new_log_store_table) =
1724 rewrite_refresh_schema_sink_fragment(
1725 &original_sink_fragment,
1726 &sink,
1727 &newly_added_columns,
1728 &removed_columns,
1729 table,
1730 fragment_graph.table_fragment_id(),
1731 self.env.id_gen_manager(),
1732 )?;
1733
1734 let streaming_job = StreamingJob::Sink(sink, None);
1735
1736 let tmp_sink_model = self
1737 .metadata_manager
1738 .catalog_controller
1739 .create_job_catalog_for_replace(&streaming_job, None, None, None)
1740 .await?;
1741 let tmp_sink_id = tmp_sink_model.job_id.as_sink_id();
1742 let StreamingJob::Sink(sink, _) = streaming_job else {
1743 unreachable!()
1744 };
1745
1746 sinks.push(AutoRefreshSchemaSinkContext {
1747 tmp_sink_id,
1748 original_sink: sink,
1749 original_fragment: original_sink_fragment,
1750 new_schema,
1751 newly_add_fields: newly_added_columns
1752 .iter()
1753 .map(|col| Field::from(&col.column_desc))
1754 .collect(),
1755 removed_column_names: removed_columns
1756 .iter()
1757 .map(|col| col.name.clone())
1758 .collect(),
1759 new_fragment: new_sink_fragment,
1760 new_log_store_table: new_log_store_table.map(Box::new),
1761 ctx: sink_ctx,
1762 });
1763 }
1764 Some(sinks)
1765 } else {
1766 None
1767 }
1768 } else {
1769 None
1770 };
1771
1772 let streaming_job_model = self
1773 .metadata_manager
1774 .catalog_controller
1775 .create_job_catalog_for_replace(
1776 &streaming_job,
1777 Some(&ctx),
1778 fragment_graph.specified_parallelism().as_ref(),
1779 Some(fragment_graph.max_parallelism()),
1780 )
1781 .await?;
1782 let tmp_id = streaming_job_model.job_id;
1783
1784 let tmp_sink_ids = auto_refresh_schema_sinks.as_ref().map(|sinks| {
1785 sinks
1786 .iter()
1787 .map(|sink| sink.tmp_sink_id.as_object_id())
1788 .collect_vec()
1789 });
1790
1791 tracing::debug!(id = %job_id, "building replace streaming job");
1792 let mut updated_sink_catalogs = vec![];
1793
1794 let mut drop_table_connector_ctx = None;
1795 let result: MetaResult<_> = try {
1796 let (mut ctx, mut stream_job_fragments) = self
1797 .build_replace_job(
1798 ctx,
1799 &streaming_job,
1800 fragment_graph,
1801 tmp_id,
1802 auto_refresh_schema_sinks,
1803 streaming_job_model,
1804 )
1805 .await?;
1806 drop_table_connector_ctx = ctx.drop_table_connector_ctx.clone();
1807 let auto_refresh_schema_sink_finish_ctx =
1808 ctx.auto_refresh_schema_sinks.as_ref().map(|sinks| {
1809 sinks
1810 .iter()
1811 .map(|sink| FinishAutoRefreshSchemaSinkContext {
1812 tmp_sink_id: sink.tmp_sink_id,
1813 original_sink_id: sink.original_sink.id,
1814 columns: sink.new_schema.clone(),
1815 new_log_store_table: sink.new_log_store_table.clone(),
1816 })
1817 .collect()
1818 });
1819
1820 if let StreamingJob::Table(_, table, ..) = &streaming_job {
1822 let union_fragment = stream_job_fragments.inner.union_fragment_for_table();
1823 let upstream_infos = self
1824 .metadata_manager
1825 .catalog_controller
1826 .get_all_upstream_sink_infos(table, union_fragment.fragment_id as _)
1827 .await?;
1828 refill_upstream_sink_union_in_table(&mut union_fragment.nodes, &upstream_infos);
1829
1830 for upstream_info in &upstream_infos {
1831 let upstream_fragment_id = upstream_info.sink_fragment_id;
1832 ctx.upstream_fragment_downstreams
1833 .entry(upstream_fragment_id)
1834 .or_default()
1835 .push(upstream_info.new_sink_downstream.clone());
1836 if upstream_info.sink_original_target_columns.is_empty() {
1837 updated_sink_catalogs.push(upstream_info.sink_id);
1838 }
1839 }
1840 }
1841
1842 let replace_upstream = ctx.replace_upstream.clone();
1843
1844 if let Some(sinks) = &ctx.auto_refresh_schema_sinks {
1845 let empty_downstreams = FragmentDownstreamRelation::default();
1846 for sink in sinks {
1847 self.metadata_manager
1848 .catalog_controller
1849 .prepare_streaming_job(
1850 sink.tmp_sink_id.as_job_id(),
1851 || [&sink.new_fragment].into_iter(),
1852 &empty_downstreams,
1853 true,
1854 None,
1855 None,
1856 )
1857 .await?;
1858 }
1859 }
1860
1861 self.metadata_manager
1862 .catalog_controller
1863 .prepare_stream_job_fragments(&stream_job_fragments, &streaming_job, true, None)
1864 .await?;
1865
1866 self.stream_manager
1867 .replace_stream_job(stream_job_fragments, ctx)
1868 .await?;
1869 (replace_upstream, auto_refresh_schema_sink_finish_ctx)
1870 };
1871
1872 match result {
1873 Ok((replace_upstream, auto_refresh_schema_sink_finish_ctx)) => {
1874 let version = self
1875 .metadata_manager
1876 .catalog_controller
1877 .finish_replace_streaming_job(
1878 tmp_id,
1879 streaming_job,
1880 replace_upstream,
1881 SinkIntoTableContext {
1882 updated_sink_catalogs,
1883 },
1884 drop_table_connector_ctx.as_ref(),
1885 auto_refresh_schema_sink_finish_ctx,
1886 )
1887 .await?;
1888 if let Some(drop_table_connector_ctx) = &drop_table_connector_ctx {
1889 self.source_manager
1890 .apply_source_change(SourceChange::DropSource {
1891 dropped_source_ids: vec![drop_table_connector_ctx.to_remove_source_id],
1892 })
1893 .await;
1894 }
1895 Ok(version)
1896 }
1897 Err(err) => {
1898 tracing::error!(id = %job_id, error = ?err.as_report(), "failed to replace job");
1899 let _ = self.metadata_manager
1900 .catalog_controller
1901 .try_abort_replacing_streaming_job(tmp_id, tmp_sink_ids)
1902 .await.inspect_err(|err| {
1903 tracing::error!(id = %job_id, error = ?err.as_report(), "failed to abort replacing job");
1904 });
1905 Err(err)
1906 }
1907 }
1908 }
1909
1910 #[await_tree::instrument(boxed, "drop_streaming_job{}({job_id})", if let DropMode::Cascade = drop_mode { "_cascade" } else { "" }
1911 )]
1912 async fn drop_streaming_job(
1913 &self,
1914 job_id: StreamingJobId,
1915 drop_mode: DropMode,
1916 ) -> MetaResult<NotificationVersion> {
1917 let (object_id, object_type) = match job_id {
1918 StreamingJobId::MaterializedView(id) => (id.as_object_id(), ObjectType::Table),
1919 StreamingJobId::Sink(id) => (id.as_object_id(), ObjectType::Sink),
1920 StreamingJobId::Table(_, id) => (id.as_object_id(), ObjectType::Table),
1921 StreamingJobId::Index(idx) => (idx.as_object_id(), ObjectType::Index),
1922 };
1923
1924 let job_status = self
1925 .metadata_manager
1926 .catalog_controller
1927 .get_streaming_job_status(job_id.id())
1928 .await?;
1929 let version = match job_status {
1930 JobStatus::Initial => {
1931 let abort_result = self
1932 .metadata_manager
1933 .catalog_controller
1934 .try_abort_creating_streaming_job(job_id.id(), true)
1935 .await?;
1936 self.iceberg_compaction_manager
1937 .clear_maintenance_for_aborted_job(&abort_result);
1938 IGNORED_NOTIFICATION_VERSION
1939 }
1940 JobStatus::Creating => {
1941 self.stream_manager
1942 .cancel_streaming_jobs(vec![job_id.id()])
1943 .await?;
1944 IGNORED_NOTIFICATION_VERSION
1945 }
1946 JobStatus::Created => self.drop_object(object_type, object_id, drop_mode).await?,
1947 };
1948
1949 Ok(version)
1950 }
1951
1952 #[await_tree::instrument]
1958 pub(crate) async fn build_stream_job(
1959 &self,
1960 stream_ctx: StreamContext,
1961 mut stream_job: StreamingJob,
1962 fragment_graph: StreamFragmentGraph,
1963 resource_type: streaming_job_resource_type::ResourceType,
1964 streaming_job_model: streaming_job::Model,
1965 since_timestamp_epoch: Option<u64>,
1966 ) -> MetaResult<(CreateStreamingJobContext, StreamJobFragmentsToCreate)> {
1967 let id = stream_job.id();
1968 let max_parallelism = NonZeroUsize::new(fragment_graph.max_parallelism()).unwrap();
1969 Self::validate_specified_parallelism(
1970 fragment_graph.specified_parallelism(),
1971 fragment_graph.specified_backfill_parallelism(),
1972 max_parallelism,
1973 )?;
1974
1975 let fragment_backfill_ordering = fragment_graph.create_fragment_backfill_ordering();
1977
1978 let (snapshot_backfill_info, cross_db_snapshot_backfill_info) =
1982 fragment_graph.collect_snapshot_backfill_info()?;
1983 assert!(
1984 snapshot_backfill_info
1985 .iter()
1986 .chain([&cross_db_snapshot_backfill_info])
1987 .flat_map(|info| info.upstream_mv_table_id_to_backfill_epoch.values())
1988 .all(|backfill_epoch| backfill_epoch.is_none()),
1989 "should not set backfill epoch when initially build the job: {:?} {:?}",
1990 snapshot_backfill_info,
1991 cross_db_snapshot_backfill_info
1992 );
1993
1994 let locality_fragment_state_table_mapping =
1995 fragment_graph.find_locality_provider_fragment_state_table_mapping();
1996
1997 self.metadata_manager
1999 .catalog_controller
2000 .validate_cross_db_snapshot_backfill(&cross_db_snapshot_backfill_info)
2001 .await?;
2002
2003 let upstream_table_ids = fragment_graph
2004 .dependent_table_ids()
2005 .iter()
2006 .filter(|id| {
2007 !cross_db_snapshot_backfill_info
2008 .upstream_mv_table_id_to_backfill_epoch
2009 .contains_key(*id)
2010 })
2011 .cloned()
2012 .collect();
2013
2014 let upstream_root_fragments = self
2015 .metadata_manager
2016 .get_upstream_root_fragments(&upstream_table_ids)
2017 .await?;
2018
2019 if snapshot_backfill_info.is_some() {
2020 match stream_job {
2021 StreamingJob::MaterializedView(_)
2022 | StreamingJob::Sink(..)
2023 | StreamingJob::Index(_, _) => {}
2024 StreamingJob::Table(_, _, _) | StreamingJob::Source(_) => {
2025 return Err(
2026 anyhow!("snapshot_backfill not enabled for table and source").into(),
2027 );
2028 }
2029 }
2030 }
2031
2032 let complete_graph = CompleteStreamFragmentGraph::with_upstreams(
2033 fragment_graph,
2034 FragmentGraphUpstreamContext {
2035 upstream_root_fragments,
2036 },
2037 (&stream_job).into(),
2038 )?;
2039 let database_resource_group = self
2040 .metadata_manager
2041 .get_database_resource_group(stream_job.database_id())
2042 .await?;
2043 let is_serverless_backfill = matches!(
2044 &resource_type,
2045 streaming_job_resource_type::ResourceType::ServerlessBackfill(true)
2046 );
2047
2048 let actor_graph_builder = ActorGraphBuilder::new(complete_graph)?;
2050
2051 let ActorGraphBuildResult {
2052 graph,
2053 downstream_fragment_relations,
2054 upstream_fragment_downstreams,
2055 replace_upstream,
2056 } = actor_graph_builder.generate_graph()?;
2057 assert!(replace_upstream.is_empty());
2058
2059 let stream_job_fragments =
2064 StreamJobFragments::new(id, graph, stream_ctx.clone(), max_parallelism.get());
2065
2066 if let Some(mview_fragment) = stream_job_fragments.mview_fragment() {
2067 stream_job.set_table_vnode_count(mview_fragment.vnode_count());
2068 }
2069
2070 let new_upstream_sink = if let StreamingJob::Sink(sink, _) = &stream_job
2071 && let Ok(table_id) = sink.get_target_table()
2072 {
2073 let tables = self
2074 .metadata_manager
2075 .get_table_catalog_by_ids(&[*table_id])
2076 .await?;
2077 let target_table = tables
2078 .first()
2079 .ok_or_else(|| MetaError::catalog_id_not_found("table", *table_id))?;
2080 let sink_fragment = stream_job_fragments
2081 .sink_fragment()
2082 .ok_or_else(|| anyhow::anyhow!("sink fragment not found for sink {}", sink.id))?;
2083 let mview_fragment_id = self
2084 .metadata_manager
2085 .catalog_controller
2086 .get_mview_fragment_by_id(table_id.as_job_id())
2087 .await?;
2088 let upstream_sink_info = build_upstream_sink_info(
2089 sink.id,
2090 sink.original_target_columns.clone(),
2091 sink_fragment.fragment_id as _,
2092 target_table,
2093 mview_fragment_id,
2094 )?;
2095 Some(upstream_sink_info)
2096 } else {
2097 None
2098 };
2099
2100 let mut cdc_table_snapshot_splits = None;
2101 if let StreamingJob::Table(None, table, TableJobType::SharedCdcSource) = &stream_job
2102 && let Some((_, stream_cdc_scan)) =
2103 parallel_cdc_table_backfill_fragment(stream_job_fragments.fragments.values())
2104 {
2105 {
2106 let splits = try_init_parallel_cdc_table_snapshot_splits(
2108 table.id,
2109 stream_cdc_scan.cdc_table_desc.as_ref().unwrap(),
2110 self.env.meta_store_ref(),
2111 stream_cdc_scan.options.as_ref().unwrap(),
2112 self.env.opts.cdc_table_split_init_insert_batch_size,
2113 self.env.opts.cdc_table_split_init_sleep_interval_splits,
2114 self.env.opts.cdc_table_split_init_sleep_duration_millis,
2115 )
2116 .await?;
2117 cdc_table_snapshot_splits = Some(splits);
2118 }
2119 }
2120
2121 let ctx = CreateStreamingJobContext {
2122 upstream_fragment_downstreams,
2123 database_resource_group,
2124 definition: stream_job.definition(),
2125 create_type: stream_job.create_type(),
2126 job_type: (&stream_job).into(),
2127 streaming_job: stream_job,
2128 new_upstream_sink,
2129 option: CreateStreamingJobOption {},
2130 snapshot_backfill_info,
2131 cross_db_snapshot_backfill_info,
2132 fragment_backfill_ordering,
2133 locality_fragment_state_table_mapping,
2134 cdc_table_snapshot_splits,
2135 is_serverless_backfill,
2136 resource_type,
2137 streaming_job_model: streaming_job_model.clone(),
2138 replace_sink: None,
2139 refresh_interval_sec: streaming_job_model.refresh_interval_sec.map(|s| s as u64),
2140 since_timestamp_epoch,
2141 };
2142
2143 Ok((
2144 ctx,
2145 StreamJobFragmentsToCreate {
2146 inner: stream_job_fragments,
2147 downstreams: downstream_fragment_relations,
2148 },
2149 ))
2150 }
2151
2152 pub(crate) async fn build_replace_job(
2158 &self,
2159 stream_ctx: StreamContext,
2160 stream_job: &StreamingJob,
2161 mut fragment_graph: StreamFragmentGraph,
2162 tmp_job_id: JobId,
2163 auto_refresh_schema_sinks: Option<Vec<AutoRefreshSchemaSinkContext>>,
2164 streaming_job_model: streaming_job::Model,
2165 ) -> MetaResult<(ReplaceStreamJobContext, StreamJobFragmentsToCreate)> {
2166 match &stream_job {
2167 StreamingJob::Table(..)
2168 | StreamingJob::Source(..)
2169 | StreamingJob::MaterializedView(..) => {}
2170 StreamingJob::Sink(..) | StreamingJob::Index(..) => {
2171 bail_not_implemented!("schema change for {}", stream_job.job_type_str())
2172 }
2173 }
2174
2175 let id = stream_job.id();
2176
2177 let mut drop_table_associated_source_id = None;
2179 if let StreamingJob::Table(None, _, _) = &stream_job {
2180 drop_table_associated_source_id = self
2181 .metadata_manager
2182 .get_table_associated_source_id(id.as_mv_table_id())
2183 .await?;
2184 }
2185
2186 let old_fragments = self.metadata_manager.get_job_fragments_by_id(id).await?;
2187 let old_internal_table_ids = old_fragments.internal_table_ids();
2188
2189 let mut drop_table_connector_ctx = None;
2191 if let Some(to_remove_source_id) = drop_table_associated_source_id {
2192 debug_assert!(old_internal_table_ids.len() == 1);
2194
2195 drop_table_connector_ctx = Some(DropTableConnectorContext {
2196 to_change_streaming_job_id: id,
2199 to_remove_state_table_id: old_internal_table_ids[0], to_remove_source_id,
2201 });
2202 } else if stream_job.is_materialized_view() {
2203 let old_fragments_upstreams = self
2206 .metadata_manager
2207 .catalog_controller
2208 .upstream_fragments(old_fragments.fragment_ids())
2209 .await?;
2210
2211 let old_state_graph =
2212 state_match::Graph::from_existing(&old_fragments, &old_fragments_upstreams);
2213 let new_state_graph = state_match::Graph::from_building(&fragment_graph);
2214 let result = state_match::match_graph(&new_state_graph, &old_state_graph)
2215 .context("incompatible altering on the streaming job states")?;
2216
2217 fragment_graph.fit_internal_table_ids_with_mapping(result.table_matches);
2218 fragment_graph.fit_snapshot_backfill_epochs(result.snapshot_backfill_epochs);
2219 } else {
2220 let old_internal_tables = self
2223 .metadata_manager
2224 .get_table_catalog_by_ids(&old_internal_table_ids)
2225 .await?;
2226 fragment_graph.fit_internal_tables_trivial(old_internal_tables)?;
2227 }
2228
2229 let original_root_fragment = old_fragments
2232 .root_fragment()
2233 .expect("root fragment not found");
2234
2235 let job_type = StreamingJobType::from(stream_job);
2236
2237 let mut downstream_fragments = self.metadata_manager.get_downstream_fragments(id).await?;
2239
2240 if let Some(auto_refresh_schema_sinks) = &auto_refresh_schema_sinks {
2241 let mut remaining_fragment: HashSet<_> = auto_refresh_schema_sinks
2242 .iter()
2243 .map(|sink| sink.original_fragment.fragment_id)
2244 .collect();
2245 for (_, downstream_fragment) in &mut downstream_fragments {
2246 if let Some(sink) = auto_refresh_schema_sinks.iter().find(|sink| {
2247 sink.original_fragment.fragment_id == downstream_fragment.fragment_id
2248 }) {
2249 assert!(remaining_fragment.remove(&downstream_fragment.fragment_id));
2250 *downstream_fragment = sink.new_fragment.clone();
2253 }
2254 }
2255 assert!(remaining_fragment.is_empty());
2256 }
2257
2258 let complete_graph = match &job_type {
2260 StreamingJobType::Table(TableJobType::General) | StreamingJobType::Source => {
2261 CompleteStreamFragmentGraph::with_downstreams(
2262 fragment_graph,
2263 FragmentGraphDownstreamContext {
2264 original_root_fragment_id: original_root_fragment.fragment_id,
2265 downstream_fragments,
2266 },
2267 job_type,
2268 )?
2269 }
2270 StreamingJobType::Table(TableJobType::SharedCdcSource)
2271 | StreamingJobType::MaterializedView => {
2272 let upstream_root_fragments = self
2274 .metadata_manager
2275 .get_upstream_root_fragments(fragment_graph.dependent_table_ids())
2276 .await?;
2277
2278 CompleteStreamFragmentGraph::with_upstreams_and_downstreams(
2279 fragment_graph,
2280 FragmentGraphUpstreamContext {
2281 upstream_root_fragments,
2282 },
2283 FragmentGraphDownstreamContext {
2284 original_root_fragment_id: original_root_fragment.fragment_id,
2285 downstream_fragments,
2286 },
2287 job_type,
2288 )?
2289 }
2290 _ => unreachable!(),
2291 };
2292
2293 let resource_group = self
2294 .metadata_manager
2295 .get_database_resource_group(stream_job.database_id())
2296 .await?;
2297
2298 let actor_graph_builder = ActorGraphBuilder::new(complete_graph)?;
2299
2300 let ActorGraphBuildResult {
2301 graph,
2302 downstream_fragment_relations,
2303 upstream_fragment_downstreams,
2304 mut replace_upstream,
2305 } = actor_graph_builder.generate_graph()?;
2306
2307 if matches!(
2309 job_type,
2310 StreamingJobType::Source | StreamingJobType::Table(TableJobType::General)
2311 ) {
2312 assert!(upstream_fragment_downstreams.is_empty());
2313 }
2314
2315 let stream_job_fragments =
2319 StreamJobFragments::new(tmp_job_id, graph, stream_ctx, old_fragments.max_parallelism);
2320
2321 if let Some(sinks) = &auto_refresh_schema_sinks {
2322 for sink in sinks {
2323 replace_upstream
2324 .remove(&sink.new_fragment.fragment_id)
2325 .expect("should exist");
2326 }
2327 }
2328
2329 let ctx = ReplaceStreamJobContext {
2333 old_fragments,
2334 replace_upstream,
2335 upstream_fragment_downstreams,
2336 streaming_job: stream_job.clone(),
2337 database_resource_group: resource_group,
2338 tmp_id: tmp_job_id,
2339 drop_table_connector_ctx,
2340 auto_refresh_schema_sinks,
2341 streaming_job_model,
2342 };
2343
2344 Ok((
2345 ctx,
2346 StreamJobFragmentsToCreate {
2347 inner: stream_job_fragments,
2348 downstreams: downstream_fragment_relations,
2349 },
2350 ))
2351 }
2352
2353 async fn alter_name(
2354 &self,
2355 relation: alter_name_request::Object,
2356 new_name: &str,
2357 ) -> MetaResult<NotificationVersion> {
2358 let (obj_type, id): (ObjectType, ObjectId) = match relation {
2359 alter_name_request::Object::TableId(id) => (ObjectType::Table, id.into()),
2360 alter_name_request::Object::ViewId(id) => (ObjectType::View, id.into()),
2361 alter_name_request::Object::IndexId(id) => (ObjectType::Index, id.into()),
2362 alter_name_request::Object::SinkId(id) => (ObjectType::Sink, id.into()),
2363 alter_name_request::Object::SourceId(id) => (ObjectType::Source, id.into()),
2364 alter_name_request::Object::SchemaId(id) => (ObjectType::Schema, id.into()),
2365 alter_name_request::Object::DatabaseId(id) => (ObjectType::Database, id.into()),
2366 alter_name_request::Object::SubscriptionId(id) => (ObjectType::Subscription, id.into()),
2367 };
2368 self.metadata_manager
2369 .catalog_controller
2370 .alter_name(obj_type, id, new_name)
2371 .await
2372 }
2373
2374 async fn alter_swap_rename(
2375 &self,
2376 object: alter_swap_rename_request::Object,
2377 ) -> MetaResult<NotificationVersion> {
2378 let (obj_type, src_id, dst_id) = match object {
2379 alter_swap_rename_request::Object::Schema(_) => unimplemented!("schema swap"),
2380 alter_swap_rename_request::Object::Table(objs) => {
2381 let (src_id, dst_id) = (objs.src_object_id, objs.dst_object_id);
2382 (ObjectType::Table, src_id, dst_id)
2383 }
2384 alter_swap_rename_request::Object::View(objs) => {
2385 let (src_id, dst_id) = (objs.src_object_id, objs.dst_object_id);
2386 (ObjectType::View, src_id, dst_id)
2387 }
2388 alter_swap_rename_request::Object::Source(objs) => {
2389 let (src_id, dst_id) = (objs.src_object_id, objs.dst_object_id);
2390 (ObjectType::Source, src_id, dst_id)
2391 }
2392 alter_swap_rename_request::Object::Sink(objs) => {
2393 let (src_id, dst_id) = (objs.src_object_id, objs.dst_object_id);
2394 (ObjectType::Sink, src_id, dst_id)
2395 }
2396 alter_swap_rename_request::Object::Subscription(objs) => {
2397 let (src_id, dst_id) = (objs.src_object_id, objs.dst_object_id);
2398 (ObjectType::Subscription, src_id, dst_id)
2399 }
2400 };
2401
2402 self.metadata_manager
2403 .catalog_controller
2404 .alter_swap_rename(obj_type, src_id, dst_id)
2405 .await
2406 }
2407
2408 async fn alter_owner(
2409 &self,
2410 object: Object,
2411 owner_id: UserId,
2412 ) -> MetaResult<NotificationVersion> {
2413 let (obj_type, id): (ObjectType, ObjectId) = match object {
2414 Object::TableId(id) => (ObjectType::Table, id.into()),
2415 Object::ViewId(id) => (ObjectType::View, id.into()),
2416 Object::SourceId(id) => (ObjectType::Source, id.into()),
2417 Object::SinkId(id) => (ObjectType::Sink, id.into()),
2418 Object::SchemaId(id) => (ObjectType::Schema, id.into()),
2419 Object::DatabaseId(id) => (ObjectType::Database, id.into()),
2420 Object::SubscriptionId(id) => (ObjectType::Subscription, id.into()),
2421 Object::ConnectionId(id) => (ObjectType::Connection, id.into()),
2422 Object::FunctionId(id) => (ObjectType::Function, id.into()),
2423 Object::SecretId(id) => (ObjectType::Secret, id.into()),
2424 };
2425 self.metadata_manager
2426 .catalog_controller
2427 .alter_owner(obj_type, id, owner_id as _)
2428 .await
2429 }
2430
2431 async fn alter_set_schema(
2432 &self,
2433 object: alter_set_schema_request::Object,
2434 new_schema_id: SchemaId,
2435 ) -> MetaResult<NotificationVersion> {
2436 let (obj_type, id): (ObjectType, ObjectId) = match object {
2437 alter_set_schema_request::Object::TableId(id) => (ObjectType::Table, id.into()),
2438 alter_set_schema_request::Object::ViewId(id) => (ObjectType::View, id.into()),
2439 alter_set_schema_request::Object::SourceId(id) => (ObjectType::Source, id.into()),
2440 alter_set_schema_request::Object::SinkId(id) => (ObjectType::Sink, id.into()),
2441 alter_set_schema_request::Object::FunctionId(id) => (ObjectType::Function, id.into()),
2442 alter_set_schema_request::Object::ConnectionId(id) => {
2443 (ObjectType::Connection, id.into())
2444 }
2445 alter_set_schema_request::Object::SubscriptionId(id) => {
2446 (ObjectType::Subscription, id.into())
2447 }
2448 };
2449 self.metadata_manager
2450 .catalog_controller
2451 .alter_schema(obj_type, id, new_schema_id)
2452 .await
2453 }
2454
2455 pub async fn wait(&self, job_id: Option<JobId>) -> MetaResult<WaitVersion> {
2456 if let Some(job_id) = job_id {
2457 let database_id = self
2458 .metadata_manager
2459 .catalog_controller
2460 .get_object_database_id(job_id)
2461 .await?;
2462 let catalog_version = self
2463 .metadata_manager
2464 .wait_streaming_job_finished(database_id, job_id)
2465 .await?;
2466 let hummock_version_id = self.barrier_manager.get_hummock_version_id().await;
2467 return Ok(WaitVersion {
2468 catalog_version,
2469 hummock_version_id,
2470 });
2471 }
2472
2473 let timeout_ms = 2 * 60 * 60 * 1000;
2474 let poll_interval = Duration::from_millis(100);
2475 for _ in 0..(timeout_ms / poll_interval.as_millis() as usize) {
2476 let creating_jobs = self
2477 .metadata_manager
2478 .catalog_controller
2479 .list_creating_jobs(true, None)
2480 .await?;
2481 if creating_jobs.is_empty() {
2482 let catalog_version = self
2483 .metadata_manager
2484 .catalog_controller
2485 .notify_frontend_trivial()
2486 .await;
2487 let hummock_version_id = self.barrier_manager.get_hummock_version_id().await;
2488 return Ok(WaitVersion {
2489 catalog_version,
2490 hummock_version_id,
2491 });
2492 }
2493
2494 sleep(poll_interval).await;
2495 }
2496 Err(MetaError::cancelled(format!(
2497 "timeout after {timeout_ms}ms"
2498 )))
2499 }
2500
2501 async fn comment_on(&self, comment: Comment) -> MetaResult<NotificationVersion> {
2502 self.metadata_manager
2503 .catalog_controller
2504 .comment_on(comment)
2505 .await
2506 }
2507
2508 async fn alter_streaming_job_config(
2509 &self,
2510 job_id: JobId,
2511 entries_to_add: HashMap<String, String>,
2512 keys_to_remove: Vec<String>,
2513 ) -> MetaResult<NotificationVersion> {
2514 self.metadata_manager
2515 .catalog_controller
2516 .alter_streaming_job_config(job_id, entries_to_add, keys_to_remove)
2517 .await
2518 }
2519}
2520
2521fn report_create_object(
2522 job_id: JobId,
2523 event_name: &str,
2524 obj_type: PbTelemetryDatabaseObject,
2525 connector_name: Option<String>,
2526 attr_info: Option<jsonbb::Value>,
2527) {
2528 report_event(
2529 PbTelemetryEventStage::CreateStreamJob,
2530 event_name,
2531 job_id.as_raw_id() as _,
2532 connector_name,
2533 Some(obj_type),
2534 attr_info,
2535 );
2536}
2537
2538pub fn build_upstream_sink_info(
2539 sink_id: SinkId,
2540 original_target_columns: Vec<PbColumnCatalog>,
2541 sink_fragment_id: FragmentId,
2542 target_table: &PbTable,
2543 target_fragment_id: FragmentId,
2544) -> MetaResult<UpstreamSinkInfo> {
2545 let sink_columns = if !original_target_columns.is_empty() {
2546 original_target_columns.clone()
2547 } else {
2548 target_table.columns.clone()
2553 };
2554
2555 let sink_output_fields = sink_columns
2556 .iter()
2557 .map(|col| Field::from(col.column_desc.as_ref().unwrap()).to_prost())
2558 .collect_vec();
2559 let output_indices = (0..sink_output_fields.len())
2560 .map(|i| i as u32)
2561 .collect_vec();
2562
2563 let dist_key_indices: anyhow::Result<Vec<u32>> = try {
2564 let sink_idx_by_col_id = sink_columns
2565 .iter()
2566 .enumerate()
2567 .map(|(idx, col)| {
2568 let column_id = col.column_desc.as_ref().unwrap().column_id;
2569 (column_id, idx as u32)
2570 })
2571 .collect::<HashMap<_, _>>();
2572 target_table
2573 .distribution_key
2574 .iter()
2575 .map(|dist_idx| {
2576 let column_id = target_table.columns[*dist_idx as usize]
2577 .column_desc
2578 .as_ref()
2579 .unwrap()
2580 .column_id;
2581 let sink_idx = sink_idx_by_col_id
2582 .get(&column_id)
2583 .ok_or_else(|| anyhow::anyhow!("column id {} not found in sink", column_id))?;
2584 Ok(*sink_idx)
2585 })
2586 .collect::<anyhow::Result<Vec<_>>>()?
2587 };
2588 let dist_key_indices =
2589 dist_key_indices.map_err(|e| e.context("failed to get distribution key indices"))?;
2590 let downstream_fragment_id = target_fragment_id as _;
2591 let new_downstream_relation = DownstreamFragmentRelation {
2592 downstream_fragment_id,
2593 dispatcher_type: DispatcherType::Hash,
2594 dist_key_indices,
2595 output_mapping: PbDispatchOutputMapping::simple(output_indices),
2596 };
2597 let current_target_columns = target_table.get_columns();
2598 let project_exprs = build_select_node_list(&sink_columns, current_target_columns)?;
2599 Ok(UpstreamSinkInfo {
2600 sink_id,
2601 sink_fragment_id: sink_fragment_id as _,
2602 sink_output_fields,
2603 sink_original_target_columns: original_target_columns,
2604 project_exprs,
2605 new_sink_downstream: new_downstream_relation,
2606 })
2607}
2608
2609pub fn refill_upstream_sink_union_in_table(
2610 union_fragment_root: &mut PbStreamNode,
2611 upstream_sink_infos: &Vec<UpstreamSinkInfo>,
2612) {
2613 visit_stream_node_cont_mut(union_fragment_root, |node| {
2614 if let Some(NodeBody::UpstreamSinkUnion(upstream_sink_union)) = &mut node.node_body {
2615 let init_upstreams = upstream_sink_infos
2616 .iter()
2617 .map(|info| PbUpstreamSinkInfo {
2618 upstream_fragment_id: info.sink_fragment_id,
2619 sink_output_schema: info.sink_output_fields.clone(),
2620 project_exprs: info.project_exprs.clone(),
2621 })
2622 .collect();
2623 upstream_sink_union.init_upstreams = init_upstreams;
2624 false
2625 } else {
2626 true
2627 }
2628 });
2629}
2630
2631#[cfg(test)]
2632mod tests {
2633 use std::num::NonZeroUsize;
2634
2635 use super::*;
2636
2637 #[test]
2638 fn test_validate_specified_parallelism_accepts_within_max() {
2639 DdlController::validate_specified_parallelism(
2640 Some(NonZeroUsize::new(4).unwrap()),
2641 Some(NonZeroUsize::new(8).unwrap()),
2642 NonZeroUsize::new(8).unwrap(),
2643 )
2644 .unwrap();
2645 }
2646
2647 #[test]
2648 fn test_validate_specified_parallelism_rejects_parallelism_over_max() {
2649 let result = DdlController::validate_specified_parallelism(
2650 Some(NonZeroUsize::new(9).unwrap()),
2651 None,
2652 NonZeroUsize::new(8).unwrap(),
2653 );
2654 assert!(matches!(
2655 result,
2656 Err(ref e) if matches!(e.inner(), MetaErrorInner::InvalidParameter(_))
2657 ));
2658 }
2659
2660 #[test]
2661 fn test_validate_specified_parallelism_rejects_backfill_parallelism_over_max() {
2662 let result = DdlController::validate_specified_parallelism(
2663 None,
2664 Some(NonZeroUsize::new(9).unwrap()),
2665 NonZeroUsize::new(8).unwrap(),
2666 );
2667 assert!(matches!(
2668 result,
2669 Err(ref e) if matches!(e.inner(), MetaErrorInner::InvalidParameter(_))
2670 ));
2671 }
2672}