1use std::collections::{BTreeMap, HashMap, HashSet};
16use std::num::NonZeroUsize;
17
18use anyhow::anyhow;
19use indexmap::IndexMap;
20use itertools::Itertools;
21use risingwave_common::catalog::{
22 ColumnCatalog, FragmentTypeFlag, FragmentTypeMask, ICEBERG_SINK_PREFIX, ICEBERG_SOURCE_PREFIX,
23 RISINGWAVE_ICEBERG_ROW_ID, ROW_ID_COLUMN_NAME, max_column_id,
24};
25use risingwave_common::config::DefaultParallelism;
26use risingwave_common::hash::VnodeCountCompat;
27use risingwave_common::id::JobId;
28use risingwave_common::secret::LocalSecretManager;
29use risingwave_common::system_param::AdaptiveParallelismStrategy;
30use risingwave_common::system_param::adaptive_parallelism_strategy::parse_strategy;
31use risingwave_common::util::iter_util::ZipEqDebug;
32use risingwave_common::util::stream_graph_visitor::{
33 visit_stream_node_body, visit_stream_node_mut, visit_stream_node_stream_scan,
34};
35use risingwave_common::{bail, current_cluster_version};
36use risingwave_connector::allow_alter_on_fly_fields::check_sink_allow_alter_on_fly_fields;
37use risingwave_connector::connector_common::validate_connection;
38use risingwave_connector::error::ConnectorError;
39use risingwave_connector::sink::file_sink::fs::FsSink;
40use risingwave_connector::sink::{CONNECTOR_TYPE_KEY, SinkError};
41use risingwave_connector::source::{
42 ConnectorProperties, UPSTREAM_SOURCE_KEY, pb_connection_type_to_connection_type,
43};
44use risingwave_connector::{WithOptionsSecResolved, WithPropertiesExt, match_sink_name_str};
45use risingwave_meta_model::object::ObjectType;
46use risingwave_meta_model::prelude::{StreamingJob as StreamingJobModel, *};
47use risingwave_meta_model::refresh_job::RefreshState;
48use risingwave_meta_model::streaming_job::BackfillOrders;
49use risingwave_meta_model::user_privilege::Action;
50use risingwave_meta_model::*;
51use risingwave_pb::catalog::table::PbEngine;
52use risingwave_pb::catalog::{PbConnection, PbCreateType, PbTable};
53use risingwave_pb::common::ThrottleType;
54use risingwave_pb::ddl_service::streaming_job_resource_type;
55use risingwave_pb::meta::alter_connector_props_request::AlterIcebergTableIds;
56use risingwave_pb::meta::list_rate_limits_response::RateLimitInfo;
57use risingwave_pb::meta::object::PbObjectInfo;
58use risingwave_pb::meta::subscribe_response::{
59 Info as NotificationInfo, Operation as NotificationOperation, Operation,
60};
61use risingwave_pb::meta::{ObjectDependency as PbObjectDependency, PbObject, PbObjectGroup};
62use risingwave_pb::plan_common::PbColumnCatalog;
63use risingwave_pb::plan_common::source_refresh_mode::{
64 RefreshMode, SourceRefreshModeFullReload, SourceRefreshModeStreaming,
65};
66use risingwave_pb::secret::PbSecretRef;
67use risingwave_pb::stream_plan::stream_fragment_graph::Parallelism;
68use risingwave_pb::stream_plan::stream_node::PbNodeBody;
69use risingwave_pb::stream_plan::{PbSinkLogStoreType, PbStreamNode, StreamScanType};
70use risingwave_pb::user::PbUserInfo;
71use risingwave_sqlparser::ast::{Engine, SqlOption, Statement};
72use risingwave_sqlparser::parser::{Parser, ParserError};
73use sea_orm::ActiveValue::Set;
74use sea_orm::sea_query::{Expr, Query, SimpleExpr};
75use sea_orm::{
76 ActiveModelTrait, ColumnTrait, Condition, DatabaseTransaction, EntityTrait, IntoActiveModel,
77 JoinType, NotSet, PaginatorTrait, QueryFilter, QuerySelect, RelationTrait, TransactionTrait,
78};
79use thiserror_ext::AsReport;
80
81use super::rename::IndexItemRewriter;
82use crate::barrier::Command;
83use crate::controller::ObjectModel;
84use crate::controller::catalog::{
85 CatalogController, DropTableConnectorContext, load_object_models,
86};
87use crate::controller::fragment::FragmentTypeMaskExt;
88use crate::controller::utils::{
89 PartialObject, build_object_group_for_delete, check_if_belongs_to_iceberg_table,
90 check_relation_name_duplicate, check_sink_into_table_cycle, ensure_job_not_canceled,
91 ensure_object_id, ensure_user_id, fetch_target_fragments, format_with_option_secret_resolved,
92 get_belong_objects, get_belong_objects_by_ids, get_referring_objects, get_table_columns,
93 grant_default_privileges_automatically, insert_fragment_relations,
94 list_object_dependencies_by_object_id, list_user_info_by_ids, upsert_user_privileges,
95};
96use crate::error::MetaErrorInner;
97use crate::manager::{NotificationVersion, StreamingJob, StreamingJobType};
98use crate::model::{
99 FragmentDownstreamRelation, FragmentReplaceUpstream, StreamContext, StreamJobFragments,
100 StreamJobFragmentsToCreate,
101};
102use crate::stream::SplitAssignment;
103use crate::{MetaError, MetaResult};
104
105pub struct AbortCreatingJobResult {
107 pub aborted: bool,
109 pub database_id: Option<DatabaseId>,
111 pub aborted_sink_ids: Vec<SinkId>,
114 pub cancel_info: Option<CancelStreamingJobInfo>,
117}
118
119pub struct CancelStreamingJobInfo {
121 pub command: Command,
123 pub streaming_job_ids: Vec<JobId>,
125 pub state_table_ids: Vec<TableId>,
127}
128
129#[derive(Debug)]
130pub struct IndependentJobChangeLogInfo {
131 pub job_id: JobId,
132 pub state_table_ids: HashSet<TableId>,
133 pub upstream_table_snapshot_epochs: HashMap<TableId, Option<u64>>,
134}
135
136#[derive(Debug)]
137pub struct TableChangeLogTruncateInfo {
138 pub subscription_retention_seconds: HashMap<TableId, u64>,
139 pub independent_jobs: Vec<IndependentJobChangeLogInfo>,
140}
141
142fn serverless_backfill_resource_group_placeholder(job_id: JobId) -> String {
143 format!("SERVERLESS_BACKFILL_RESOURCE_GROUP_TBD_FOR_{job_id}")
144}
145
146fn job_initial_catalog_resource_group(
147 resource_type: &streaming_job_resource_type::ResourceType,
148 job_id: JobId,
149) -> Option<String> {
150 match resource_type {
151 streaming_job_resource_type::ResourceType::Regular(_)
152 | streaming_job_resource_type::ResourceType::ServerlessBackfill(false) => None,
153 streaming_job_resource_type::ResourceType::SpecificResourceGroup(group) => {
154 Some(group.clone())
155 }
156 streaming_job_resource_type::ResourceType::ServerlessBackfill(true) => {
157 Some(serverless_backfill_resource_group_placeholder(job_id))
158 }
159 }
160}
161
162#[derive(Debug)]
167struct DependentSourceFragmentUpdate {
168 job_ids: Vec<JobId>,
169 with_properties: BTreeMap<String, String>,
170 secret_refs: BTreeMap<String, PbSecretRef>,
171 is_shared_source: bool,
172}
173
174#[derive(Debug)]
175struct ReplaceOriginalJobInfo {
176 max_parallelism: i32,
177 timezone: Option<String>,
178 config_override: Option<String>,
179 adaptive_parallelism_strategy: Option<String>,
180 parallelism: StreamingParallelism,
181 specific_resource_group: Option<String>,
182}
183
184impl ReplaceOriginalJobInfo {
185 fn resolved_parallelism(
186 &self,
187 specified_parallelism: Option<&NonZeroUsize>,
188 ) -> StreamingParallelism {
189 specified_parallelism
190 .map(|n| StreamingParallelism::Fixed(n.get() as _))
191 .unwrap_or_else(|| self.parallelism.clone())
192 }
193
194 fn stream_context(&self, ctx_override: Option<&StreamContext>) -> StreamContext {
195 StreamContext {
196 timezone: ctx_override
197 .and_then(|ctx| ctx.timezone.clone())
198 .or_else(|| self.timezone.clone()),
199 config_override: self.config_override.clone().unwrap_or_default().into(),
202 }
203 }
204
205 fn resource_type(&self) -> streaming_job_resource_type::ResourceType {
206 match &self.specific_resource_group {
207 Some(group) => {
208 streaming_job_resource_type::ResourceType::SpecificResourceGroup(group.clone())
209 }
210 None => streaming_job_resource_type::ResourceType::Regular(true),
211 }
212 }
213}
214
215impl From<streaming_job::Model> for ReplaceOriginalJobInfo {
216 fn from(model: streaming_job::Model) -> Self {
217 Self {
218 max_parallelism: model.max_parallelism,
219 timezone: model.timezone,
220 config_override: model.config_override,
221 adaptive_parallelism_strategy: model.adaptive_parallelism_strategy,
222 parallelism: model.parallelism,
223 specific_resource_group: model.specific_resource_group,
224 }
225 }
226}
227
228fn update_sink_node_rate_limit(node: &mut PbNodeBody, rate_limit: Option<u32>) -> MetaResult<bool> {
229 let PbNodeBody::Sink(node) = node else {
230 return Ok(false);
231 };
232 if node.log_store_type != PbSinkLogStoreType::KvLogStore as i32 {
233 return Err(MetaError::invalid_parameter(
234 "sink rate limit is only supported for kv log store, please SET sink_decouple = TRUE before CREATE SINK",
235 ));
236 }
237 node.rate_limit = rate_limit;
238 Ok(true)
239}
240
241impl CatalogController {
242 pub async fn get_table_change_log_truncate_info(
243 &self,
244 ) -> MetaResult<TableChangeLogTruncateInfo> {
245 let inner = self.inner.read().await;
246
247 let subscriptions: Vec<(TableId, i64)> = Subscription::find()
248 .select_only()
249 .columns([
250 subscription::Column::DependentTableId,
251 subscription::Column::RetentionSeconds,
252 ])
253 .into_tuple()
254 .all(&inner.db)
255 .await?;
256 let mut subscription_retention_seconds = HashMap::new();
257 for (table_id, retention_seconds) in subscriptions {
258 let retention_seconds = u64::try_from(retention_seconds).map_err(|_| {
259 anyhow!(
260 "subscription on table {} has invalid retention seconds {}",
261 table_id,
262 retention_seconds
263 )
264 })?;
265 subscription_retention_seconds
266 .entry(table_id)
267 .and_modify(|retention: &mut u64| *retention = (*retention).max(retention_seconds))
268 .or_insert(retention_seconds);
269 }
270
271 let jobs: Vec<JobId> = StreamingJobModel::find()
272 .select_only()
273 .column(streaming_job::Column::JobId)
274 .filter(
275 Condition::any()
276 .add(streaming_job::Column::JobStatus.eq(JobStatus::Creating))
277 .add(streaming_job::Column::RefreshIntervalSec.is_not_null()),
278 )
279 .into_tuple::<JobId>()
280 .all(&inner.db)
281 .await?;
282 let mut job_info: HashMap<_, _> = jobs
283 .into_iter()
284 .map(|job_id| {
285 (
286 job_id,
287 IndependentJobChangeLogInfo {
288 job_id,
289 state_table_ids: HashSet::new(),
290 upstream_table_snapshot_epochs: HashMap::new(),
291 },
292 )
293 })
294 .collect();
295 if !job_info.is_empty() {
296 let fragments = Fragment::find()
297 .filter(fragment::Column::JobId.is_in(job_info.keys().copied()))
298 .all(&inner.db)
299 .await?;
300 for fragment in fragments {
301 let info = job_info
302 .get_mut(&fragment.job_id)
303 .expect("job should exist");
304 info.state_table_ids
305 .extend(fragment.state_table_ids.inner_ref().iter().copied());
306 let mut collection_error = None;
307 visit_stream_node_stream_scan(&fragment.stream_node.to_protobuf(), |stream_scan| {
308 let scan_type = match StreamScanType::try_from(stream_scan.stream_scan_type) {
309 Ok(scan_type) => scan_type,
310 Err(err) => {
311 collection_error = Some(anyhow::Error::new(err).context(format!(
312 "invalid persisted stream scan type {} in job {} fragment {}",
313 stream_scan.stream_scan_type, fragment.job_id, fragment.fragment_id
314 )));
315 return;
316 }
317 };
318 if scan_type != StreamScanType::SnapshotBackfill {
319 return;
320 }
321 match info
322 .upstream_table_snapshot_epochs
323 .entry(stream_scan.table_id)
324 {
325 std::collections::hash_map::Entry::Occupied(entry) => {
326 if entry.get() != &stream_scan.snapshot_backfill_epoch {
327 collection_error = Some(anyhow!(
328 "job {} has inconsistent snapshot epochs for upstream table {}",
329 fragment.job_id,
330 stream_scan.table_id
331 ));
332 }
333 }
334 std::collections::hash_map::Entry::Vacant(entry) => {
335 entry.insert(stream_scan.snapshot_backfill_epoch);
336 }
337 }
338 });
339 if let Some(err) = collection_error {
340 return Err(err.into());
341 }
342 }
343 }
344 let independent_jobs = job_info
345 .into_values()
346 .filter(|info| !info.upstream_table_snapshot_epochs.is_empty())
347 .collect();
348 Ok(TableChangeLogTruncateInfo {
349 subscription_retention_seconds,
350 independent_jobs,
351 })
352 }
353
354 pub async fn get_pinned_snapshot_epochs(&self) -> MetaResult<HashMap<TableId, HashSet<u64>>> {
355 let inner = self.inner.read().await;
358 let creating_job_ids = StreamingJobModel::find()
359 .select_only()
360 .column(streaming_job::Column::JobId)
361 .filter(streaming_job::Column::JobStatus.eq(JobStatus::Creating))
362 .into_tuple::<JobId>()
363 .all(&inner.db)
364 .await?;
365 if creating_job_ids.is_empty() {
366 return Ok(HashMap::new());
367 }
368 let fragments = Fragment::find()
369 .filter(fragment::Column::JobId.is_in(creating_job_ids))
370 .all(&inner.db)
371 .await?;
372 let mut pinned_snapshot_epochs: HashMap<TableId, HashSet<u64>> = HashMap::new();
373 for fragment in fragments {
374 visit_stream_node_stream_scan(&fragment.stream_node.to_protobuf(), |stream_scan| {
375 let scan_type = match StreamScanType::try_from(stream_scan.stream_scan_type) {
376 Ok(scan_type) => scan_type,
377 Err(err) => {
378 tracing::warn!(
379 job_id = %fragment.job_id,
380 fragment_id = %fragment.fragment_id,
381 stream_scan_type = stream_scan.stream_scan_type,
382 error = %err.as_report(),
383 "invalid persisted stream scan type, skip collecting snapshot pin"
384 );
385 return;
386 }
387 };
388 if !matches!(
389 scan_type,
390 StreamScanType::SnapshotBackfill | StreamScanType::CrossDbSnapshotBackfill
391 ) {
392 return;
393 }
394 let Some(epoch) = stream_scan.snapshot_backfill_epoch else {
395 tracing::warn!(
396 job_id = %fragment.job_id,
397 fragment_id = %fragment.fragment_id,
398 table_id = %stream_scan.table_id,
399 "persisted snapshot backfill epoch is not set, skip collecting snapshot pin"
400 );
401 return;
402 };
403 pinned_snapshot_epochs
404 .entry(stream_scan.table_id)
405 .or_default()
406 .insert(epoch);
407 });
408 }
409 Ok(pinned_snapshot_epochs)
410 }
411
412 #[expect(clippy::too_many_arguments)]
413 pub async fn create_streaming_job_obj(
414 txn: &DatabaseTransaction,
415 obj_type: ObjectType,
416 owner_id: UserId,
417 belong_to_oid: Option<ObjectId>,
418 create_type: PbCreateType,
419 ctx: StreamContext,
420 adaptive_parallelism_strategy: Option<AdaptiveParallelismStrategy>,
421 streaming_parallelism: StreamingParallelism,
422 max_parallelism: usize,
423 resource_type: streaming_job_resource_type::ResourceType,
424 backfill_parallelism: Option<StreamingParallelism>,
425 backfill_adaptive_parallelism_strategy: Option<AdaptiveParallelismStrategy>,
426 refresh_interval_sec: Option<u64>,
427 ) -> MetaResult<streaming_job::Model> {
428 let obj = Self::create_object(txn, obj_type, owner_id, belong_to_oid).await?;
429 let job_id = obj.oid.as_job_id();
430 let is_serverless_backfill = matches!(
431 &resource_type,
432 streaming_job_resource_type::ResourceType::ServerlessBackfill(true)
433 );
434 let model = streaming_job::Model {
435 job_id,
436 job_status: JobStatus::Initial,
437 create_type: create_type.into(),
438 timezone: ctx.timezone,
439 config_override: Some(ctx.config_override.to_string()),
440 adaptive_parallelism_strategy: adaptive_parallelism_strategy
441 .as_ref()
442 .map(ToString::to_string),
443 parallelism: streaming_parallelism,
444 backfill_parallelism,
445 backfill_adaptive_parallelism_strategy: backfill_adaptive_parallelism_strategy
446 .as_ref()
447 .map(ToString::to_string),
448 backfill_orders: None,
449 max_parallelism: max_parallelism as _,
450 specific_resource_group: job_initial_catalog_resource_group(&resource_type, job_id),
451 is_serverless_backfill,
452 refresh_interval_sec: refresh_interval_sec.map(|s| s as i64),
453 };
454 let job = model.clone().into_active_model();
455 StreamingJobModel::insert(job).exec(txn).await?;
456
457 Ok(model)
458 }
459
460 #[expect(clippy::too_many_arguments)]
466 #[await_tree::instrument]
467 pub async fn create_job_catalog(
468 &self,
469 streaming_job: &mut StreamingJob,
470 ctx: &StreamContext,
471 parallelism: &Option<Parallelism>,
472 max_parallelism: usize,
473 mut dependencies: HashSet<ObjectId>,
474 resource_type: streaming_job_resource_type::ResourceType,
475 backfill_parallelism: &Option<Parallelism>,
476 adaptive_parallelism_strategy: Option<AdaptiveParallelismStrategy>,
477 backfill_adaptive_parallelism_strategy: Option<AdaptiveParallelismStrategy>,
478 replace_sink: Option<&SinkId>,
479 refresh_interval_sec: Option<u64>,
480 ) -> MetaResult<streaming_job::Model> {
481 let inner = self.inner.write().await;
482 let txn = inner.db.begin().await?;
483 let create_type = streaming_job.create_type();
484
485 let streaming_parallelism = match (parallelism, self.env.opts.default_parallelism) {
486 (None, DefaultParallelism::Full) => StreamingParallelism::Adaptive,
487 (None, DefaultParallelism::Default(n)) => StreamingParallelism::Fixed(n.get()),
488 (Some(n), _) => StreamingParallelism::Fixed(n.parallelism as _),
489 };
490 let backfill_parallelism = backfill_parallelism
491 .as_ref()
492 .map(|p| StreamingParallelism::Fixed(p.parallelism as _))
493 .or_else(|| {
494 backfill_adaptive_parallelism_strategy
495 .as_ref()
496 .map(|_| StreamingParallelism::Adaptive)
497 });
498
499 ensure_user_id(streaming_job.owner() as _, &txn).await?;
500 ensure_object_id(ObjectType::Database, streaming_job.database_id(), &txn).await?;
501 ensure_object_id(ObjectType::Schema, streaming_job.schema_id(), &txn).await?;
502 if let Some(old_sink_id) = replace_sink {
503 let StreamingJob::Sink(sink, _) = streaming_job else {
504 bail!("replacement sink catalog requires a sink job")
505 };
506 let (old_sink, old_object) = Sink::find_by_id(*old_sink_id)
507 .find_also_related(Object)
508 .one(&txn)
509 .await?
510 .and_then(|(sink, object)| object.map(|object| (sink, object)))
511 .ok_or_else(|| MetaError::catalog_id_not_found("sink", *old_sink_id))?;
512 let old_streaming_job = StreamingJobModel::find_by_id(old_sink_id.as_job_id())
513 .one(&txn)
514 .await?
515 .ok_or_else(|| MetaError::catalog_id_not_found("sink", *old_sink_id))?;
516 if old_object.obj_type != ObjectType::Sink
517 || old_object.database_id != Some(sink.database_id)
518 || old_object.schema_id != Some(sink.schema_id)
519 || old_sink.name != sink.name
520 {
521 bail!(
522 "old sink {} does not match replacement sink {}",
523 old_sink_id,
524 sink.name
525 );
526 }
527 if old_sink.target_table.is_some() || sink.target_table.is_some() {
528 bail!("replace sink into table is not supported");
529 }
530 if old_streaming_job.job_status != JobStatus::Created {
531 bail!("sink {} is not ready to be replaced", old_sink_id);
532 }
533 } else {
534 check_relation_name_duplicate(
535 &streaming_job.name(),
536 streaming_job.database_id(),
537 streaming_job.schema_id(),
538 &txn,
539 )
540 .await?;
541 }
542
543 if !dependencies.is_empty() {
545 let altering_cnt = ObjectDependency::find()
546 .join(
547 JoinType::InnerJoin,
548 object_dependency::Relation::Object1.def(),
549 )
550 .join(JoinType::InnerJoin, object::Relation::StreamingJob.def())
551 .filter(
552 object_dependency::Column::Oid
553 .is_in(dependencies.clone())
554 .and(object::Column::ObjType.eq(ObjectType::Table))
555 .and(streaming_job::Column::JobStatus.ne(JobStatus::Created))
556 .and(
557 object::Column::Oid.not_in_subquery(
559 Query::select()
560 .column(table::Column::TableId)
561 .from(Table)
562 .to_owned(),
563 ),
564 ),
565 )
566 .count(&txn)
567 .await?;
568 if altering_cnt != 0 {
569 return Err(MetaError::permission_denied(
570 "some dependent relations are being altered",
571 ));
572 }
573
574 let batch_refresh_cnt = StreamingJobModel::find()
577 .filter(
578 streaming_job::Column::JobId
579 .is_in(dependencies.iter().map(|id| JobId::new(id.as_raw_id())))
580 .and(streaming_job::Column::RefreshIntervalSec.is_not_null()),
581 )
582 .count(&txn)
583 .await?;
584 if batch_refresh_cnt != 0 {
585 return Err(MetaError::permission_denied(
586 "creating streaming jobs on batch refresh materialized views is not supported",
587 ));
588 }
589 }
590
591 let streaming_job_model = match streaming_job {
592 StreamingJob::MaterializedView(table) => {
593 let streaming_job_model = Self::create_streaming_job_obj(
594 &txn,
595 ObjectType::Table,
596 table.owner as _,
597 Some(table.schema_id.as_object_id()),
598 create_type,
599 ctx.clone(),
600 adaptive_parallelism_strategy,
601 streaming_parallelism,
602 max_parallelism,
603 resource_type.clone(),
604 backfill_parallelism.clone(),
605 backfill_adaptive_parallelism_strategy,
606 refresh_interval_sec,
607 )
608 .await?;
609 table.id = streaming_job_model.job_id.as_mv_table_id();
610 let table_model: table::ActiveModel = table.clone().into();
611 Table::insert(table_model).exec(&txn).await?;
612 streaming_job_model
613 }
614 StreamingJob::Sink(sink, belong_to_table_id) => {
615 if let Some(target_table_id) = sink.target_table
616 && check_sink_into_table_cycle(
617 target_table_id.into(),
618 dependencies.iter().cloned().collect(),
619 &txn,
620 )
621 .await?
622 {
623 bail!("Creating such a sink will result in circular dependency.");
624 }
625
626 let streaming_job_model = Self::create_streaming_job_obj(
627 &txn,
628 ObjectType::Sink,
629 sink.owner as _,
630 Some(
631 belong_to_table_id
632 .map(|table_id| table_id.as_object_id())
633 .unwrap_or(sink.schema_id.as_object_id()),
634 ),
635 create_type,
636 ctx.clone(),
637 adaptive_parallelism_strategy,
638 streaming_parallelism,
639 max_parallelism,
640 resource_type.clone(),
641 backfill_parallelism.clone(),
642 backfill_adaptive_parallelism_strategy,
643 None, )
645 .await?;
646 sink.id = streaming_job_model.job_id.as_sink_id();
647 if let Some(old_sink_id) = replace_sink {
648 let final_sink_name = sink.name.clone();
649 sink.name = format!("__rw_replacing_sink_{}_{}", old_sink_id, sink.id);
653 tracing::debug!(
654 old_sink_id = %old_sink_id,
655 new_sink_id = %sink.id,
656 final_name = %final_sink_name,
657 temp_name = %sink.name,
658 "created replacement sink catalog with temporary name"
659 );
660 }
661 let sink_model: sink::ActiveModel = sink.clone().into();
662 Sink::insert(sink_model).exec(&txn).await?;
663 streaming_job_model
664 }
665 StreamingJob::Table(src, table, _) => {
666 let streaming_job_model = Self::create_streaming_job_obj(
667 &txn,
668 ObjectType::Table,
669 table.owner as _,
670 Some(table.schema_id.as_object_id()),
671 create_type,
672 ctx.clone(),
673 adaptive_parallelism_strategy,
674 streaming_parallelism,
675 max_parallelism,
676 resource_type.clone(),
677 backfill_parallelism.clone(),
678 backfill_adaptive_parallelism_strategy,
679 None, )
681 .await?;
682 let job_id = streaming_job_model.job_id;
683 table.id = job_id.as_mv_table_id();
684 if let Some(src) = src {
685 let src_obj = Self::create_object(
686 &txn,
687 ObjectType::Source,
688 src.owner as _,
689 Some(job_id.as_object_id()),
690 )
691 .await?;
692 src.id = src_obj.oid.as_source_id();
693 src.optional_associated_table_id = Some(job_id.as_mv_table_id().into());
694 table.optional_associated_source_id = Some(src_obj.oid.as_source_id().into());
695 let source: source::ActiveModel = src.clone().into();
696 Source::insert(source).exec(&txn).await?;
697 }
698 let table_model: table::ActiveModel = table.clone().into();
699 Table::insert(table_model).exec(&txn).await?;
700
701 if table.refreshable {
702 let trigger_interval_secs = src
703 .as_ref()
704 .and_then(|source_catalog| source_catalog.refresh_mode)
705 .and_then(
706 |source_refresh_mode| match source_refresh_mode.refresh_mode {
707 Some(RefreshMode::FullReload(SourceRefreshModeFullReload {
708 refresh_interval_sec,
709 })) => refresh_interval_sec,
710 Some(RefreshMode::Streaming(SourceRefreshModeStreaming {})) => None,
711 None => None,
712 },
713 );
714
715 RefreshJob::insert(refresh_job::ActiveModel {
716 table_id: Set(table.id),
717 last_trigger_time: Set(None),
718 trigger_interval_secs: Set(trigger_interval_secs),
719 current_status: Set(RefreshState::Idle),
720 last_success_time: Set(None),
721 })
722 .exec(&txn)
723 .await?;
724 }
725 streaming_job_model
726 }
727 StreamingJob::Index(index, table) => {
728 ensure_object_id(ObjectType::Table, index.primary_table_id, &txn).await?;
729 let streaming_job_model = Self::create_streaming_job_obj(
730 &txn,
731 ObjectType::Index,
732 index.owner as _,
733 Some(index.schema_id.as_object_id()),
734 create_type,
735 ctx.clone(),
736 adaptive_parallelism_strategy,
737 streaming_parallelism,
738 max_parallelism,
739 resource_type.clone(),
740 backfill_parallelism.clone(),
741 backfill_adaptive_parallelism_strategy,
742 None, )
744 .await?;
745 let job_id = streaming_job_model.job_id;
747 index.id = job_id.as_index_id();
748 index.index_table_id = job_id.as_mv_table_id();
749 table.id = job_id.as_mv_table_id();
750
751 ObjectDependency::insert(object_dependency::ActiveModel {
752 oid: Set(index.primary_table_id.into()),
753 used_by: Set(table.id.into()),
754 ..Default::default()
755 })
756 .exec(&txn)
757 .await?;
758
759 let table_model: table::ActiveModel = table.clone().into();
760 Table::insert(table_model).exec(&txn).await?;
761 let index_model: index::ActiveModel = index.clone().into();
762 Index::insert(index_model).exec(&txn).await?;
763 streaming_job_model
764 }
765 StreamingJob::Source(src) => {
766 let streaming_job_model = Self::create_streaming_job_obj(
767 &txn,
768 ObjectType::Source,
769 src.owner as _,
770 Some(src.schema_id.as_object_id()),
771 create_type,
772 ctx.clone(),
773 adaptive_parallelism_strategy,
774 streaming_parallelism,
775 max_parallelism,
776 resource_type.clone(),
777 backfill_parallelism.clone(),
778 backfill_adaptive_parallelism_strategy,
779 None, )
781 .await?;
782 src.id = streaming_job_model.job_id.as_shared_source_id();
783 let source_model: source::ActiveModel = src.clone().into();
784 Source::insert(source_model).exec(&txn).await?;
785 streaming_job_model
786 }
787 };
788
789 dependencies.extend(
791 streaming_job
792 .dependent_secret_ids()?
793 .into_iter()
794 .map(|id| id.as_object_id()),
795 );
796 dependencies.extend(
798 streaming_job
799 .dependent_connection_ids()?
800 .into_iter()
801 .map(|id| id.as_object_id()),
802 );
803
804 if !dependencies.is_empty() {
806 ObjectDependency::insert_many(dependencies.into_iter().map(|oid| {
807 object_dependency::ActiveModel {
808 oid: Set(oid),
809 used_by: Set(streaming_job.id().as_object_id()),
810 ..Default::default()
811 }
812 }))
813 .exec(&txn)
814 .await?;
815 }
816
817 txn.commit().await?;
818
819 Ok(streaming_job_model)
820 }
821
822 pub async fn create_internal_table_catalog(
830 &self,
831 job: &StreamingJob,
832 mut incomplete_internal_tables: Vec<PbTable>,
833 ) -> MetaResult<HashMap<TableId, TableId>> {
834 let job_id = job.id();
835 let inner = self.inner.write().await;
836 let txn = inner.db.begin().await?;
837
838 ensure_job_not_canceled(job_id, &txn).await?;
840
841 let mut table_id_map = HashMap::new();
842 for table in &mut incomplete_internal_tables {
843 let table_id = Self::create_object(
844 &txn,
845 ObjectType::Table,
846 table.owner as _,
847 Some(job_id.as_object_id()),
848 )
849 .await?
850 .oid
851 .as_table_id();
852 table_id_map.insert(table.id, table_id);
853 table.id = table_id;
854 table.job_id = Some(job_id);
855
856 let table_model = table::ActiveModel {
857 table_id: Set(table_id),
858 belongs_to_job_id: Set(Some(job_id)),
859 fragment_id: NotSet,
860 ..table.clone().into()
861 };
862 Table::insert(table_model).exec(&txn).await?;
863 }
864 txn.commit().await?;
865
866 Ok(table_id_map)
867 }
868
869 pub async fn update_streaming_job_resource_group(
870 &self,
871 job_id: JobId,
872 resource_group: String,
873 ) -> MetaResult<()> {
874 let inner = self.inner.write().await;
875 let txn = inner.db.begin().await?;
876
877 ensure_job_not_canceled(job_id, &txn).await?;
878
879 let job = streaming_job::ActiveModel {
880 job_id: Set(job_id),
881 specific_resource_group: Set(Some(resource_group)),
882 ..Default::default()
883 };
884 StreamingJobModel::update(job).exec(&txn).await?;
885
886 txn.commit().await?;
887
888 Ok(())
889 }
890
891 pub async fn prepare_stream_job_fragments(
892 &self,
893 stream_job_fragments: &StreamJobFragmentsToCreate,
894 streaming_job: &StreamingJob,
895 for_replace: bool,
896 backfill_orders: Option<BackfillOrders>,
897 ) -> MetaResult<()> {
898 self.prepare_streaming_job(
899 stream_job_fragments.stream_job_id(),
900 || stream_job_fragments.fragments.values(),
901 &stream_job_fragments.downstreams,
902 for_replace,
903 Some(streaming_job),
904 backfill_orders,
905 )
906 .await
907 }
908
909 #[await_tree::instrument("prepare_streaming_job_for_{}", if for_replace { "replace" } else { "create" }
915 )]
916 pub async fn prepare_streaming_job<'a, I: Iterator<Item = &'a crate::model::Fragment> + 'a>(
917 &self,
918 job_id: JobId,
919 get_fragments: impl Fn() -> I + 'a,
920 downstreams: &FragmentDownstreamRelation,
921 for_replace: bool,
922 creating_streaming_job: Option<&'a StreamingJob>,
923 backfill_orders: Option<BackfillOrders>,
924 ) -> MetaResult<()> {
925 let fragments = Self::prepare_fragment_models_from_fragments(job_id, get_fragments())?;
926
927 let inner = self.inner.write().await;
928
929 let txn = inner.db.begin().await?;
930
931 ensure_job_not_canceled(job_id, &txn).await?;
933
934 if let Some(backfill_orders) = backfill_orders {
935 let job = streaming_job::ActiveModel {
936 job_id: Set(job_id),
937 backfill_orders: Set(Some(backfill_orders)),
938 ..Default::default()
939 };
940 StreamingJobModel::update(job).exec(&txn).await?;
941 }
942
943 let state_table_ids = fragments
944 .iter()
945 .flat_map(|fragment| fragment.state_table_ids.inner_ref().clone())
946 .collect_vec();
947
948 let inserted_fragment_ids: Vec<crate::model::FragmentId> = fragments
950 .iter()
951 .map(|f| f.fragment_id as crate::model::FragmentId)
952 .collect();
953
954 if !fragments.is_empty() {
955 let fragment_models = fragments
956 .into_iter()
957 .map(|fragment| fragment.into_active_model())
958 .collect_vec();
959 Fragment::insert_many(fragment_models).exec(&txn).await?;
960 }
961
962 if !for_replace {
965 let all_tables = StreamJobFragments::collect_tables(get_fragments());
966 for state_table_id in state_table_ids {
967 let table = all_tables
971 .get(&state_table_id)
972 .unwrap_or_else(|| panic!("table {} not found", state_table_id));
973 assert_eq!(table.id, state_table_id);
974 let vnode_count = table.vnode_count();
975
976 Table::update(table::ActiveModel {
977 table_id: Set(state_table_id as _),
978 fragment_id: Set(Some(table.fragment_id)),
979 vnode_count: Set(vnode_count as _),
980 ..Default::default()
981 })
982 .exec(&txn)
983 .await?;
984 }
985 }
986
987 insert_fragment_relations(&txn, downstreams).await?;
988
989 if !for_replace {
990 if let Some(StreamingJob::Table(_, table, _)) = creating_streaming_job {
992 Table::update(table::ActiveModel {
993 table_id: Set(table.id),
994 dml_fragment_id: Set(table.dml_fragment_id),
995 ..Default::default()
996 })
997 .exec(&txn)
998 .await?;
999 }
1000 }
1001
1002 txn.commit().await?;
1003
1004 self.env
1008 .notification_manager()
1009 .notify_serving_fragment_mapping_update(inserted_fragment_ids);
1010
1011 self.notify_hummock_table_cache_refill_policy_if_explicit(&inner, job_id)
1012 .await?;
1013
1014 Ok(())
1015 }
1016
1017 #[await_tree::instrument]
1022 pub async fn try_abort_creating_streaming_job(
1023 &self,
1024 job_id: JobId,
1025 is_cancelled: bool,
1026 ) -> MetaResult<AbortCreatingJobResult> {
1027 let mut inner = self.inner.write().await;
1028 let txn = inner.db.begin().await?;
1029
1030 let obj = Object::find_by_id(job_id).one(&txn).await?;
1031 let Some(obj) = obj else {
1032 tracing::warn!(
1033 id = %job_id,
1034 "streaming job not found when aborting creating, might be cancelled already or cleaned by recovery"
1035 );
1036 return Ok(AbortCreatingJobResult {
1037 aborted: true,
1038 database_id: None,
1039 aborted_sink_ids: vec![],
1040 cancel_info: None,
1041 });
1042 };
1043 let database_id = obj
1044 .database_id
1045 .ok_or_else(|| anyhow!("obj has no database id: {:?}", obj))?;
1046 let streaming_job = streaming_job::Entity::find_by_id(job_id).one(&txn).await?;
1047
1048 if let Some(streaming_job) = &streaming_job {
1049 if streaming_job.job_status == JobStatus::Created {
1050 tracing::warn!(%job_id, "streaming job is already created, ignore abort request");
1051 return Ok(AbortCreatingJobResult {
1052 aborted: false,
1053 database_id: Some(database_id),
1054 aborted_sink_ids: vec![],
1055 cancel_info: None,
1056 });
1057 }
1058
1059 if !is_cancelled && streaming_job.job_status == JobStatus::Creating {
1060 if (obj.obj_type == ObjectType::Table || obj.obj_type == ObjectType::Sink)
1061 && check_if_belongs_to_iceberg_table(&txn, job_id).await?
1062 {
1063 } else {
1065 tracing::warn!(
1068 id = %job_id,
1069 "streaming job is still in creating status"
1070 );
1071 return Ok(AbortCreatingJobResult {
1072 aborted: false,
1073 database_id: Some(database_id),
1074 aborted_sink_ids: vec![],
1075 cancel_info: None,
1076 });
1077 }
1078 }
1079 }
1080
1081 let mut objects_to_abort = vec![obj];
1082 objects_to_abort.extend(get_belong_objects(&txn, job_id.as_object_id()).await?);
1083 let object_ids_to_abort = objects_to_abort
1084 .iter()
1085 .map(|object| object.oid)
1086 .collect_vec();
1087 let streaming_job_ids_to_abort: Vec<JobId> = StreamingJobModel::find()
1088 .select_only()
1089 .column(streaming_job::Column::JobId)
1090 .filter(streaming_job::Column::JobId.is_in(object_ids_to_abort.iter().copied()))
1091 .into_tuple()
1092 .all(&txn)
1093 .await?;
1094 let object_models = load_object_models(&txn, &objects_to_abort).await?;
1095
1096 let fragments: Vec<(FragmentId, i32, TableIdArray)> = Fragment::find()
1097 .select_only()
1098 .columns([
1099 fragment::Column::FragmentId,
1100 fragment::Column::FragmentTypeMask,
1101 fragment::Column::StateTableIds,
1102 ])
1103 .filter(fragment::Column::JobId.is_in(streaming_job_ids_to_abort.iter().copied()))
1104 .into_tuple()
1105 .all(&txn)
1106 .await?;
1107 let cancel_info = if is_cancelled {
1108 let state_table_ids: HashSet<_> = fragments
1109 .iter()
1110 .flat_map(|(_, _, state_table_ids)| state_table_ids.inner_ref().iter().copied())
1111 .collect();
1112 let sink_fragment_ids = fragments
1113 .iter()
1114 .filter_map(|(fragment_id, fragment_type_mask, _)| {
1115 FragmentTypeMask::from(*fragment_type_mask)
1116 .contains(FragmentTypeFlag::Sink)
1117 .then_some(*fragment_id)
1118 })
1119 .collect_vec();
1120 let sink_target_fragments = fetch_target_fragments(&txn, sink_fragment_ids).await?;
1121 let mut dropped_sink_fragment_by_targets = HashMap::new();
1122 for (sink_fragment, target_fragments) in sink_target_fragments {
1123 if let Some(target_fragment) = target_fragments.first() {
1124 dropped_sink_fragment_by_targets
1125 .entry(*target_fragment)
1126 .or_insert_with(Vec::new)
1127 .push(sink_fragment);
1128 }
1129 }
1130
1131 Some(CancelStreamingJobInfo {
1132 command: Command::DropStreamingJobs {
1133 streaming_job_ids: streaming_job_ids_to_abort.iter().copied().collect(),
1134 unregistered_state_table_ids: state_table_ids.clone(),
1135 dropped_sink_fragment_by_targets,
1136 },
1137 streaming_job_ids: streaming_job_ids_to_abort.clone(),
1138 state_table_ids: state_table_ids.into_iter().collect(),
1139 })
1140 } else {
1141 None
1142 };
1143
1144 let need_notify = streaming_job
1146 .as_ref()
1147 .is_some_and(|job| job.job_status != JobStatus::Initial);
1148
1149 if is_cancelled {
1150 let dropped_tables = object_models.iter().filter_map(|object_info| {
1151 if let PbObjectInfo::Table(table) = object_info {
1152 Some((table.id, table.clone()))
1153 } else {
1154 None
1155 }
1156 });
1157 inner.dropped_tables.extend(dropped_tables);
1158 }
1159
1160 let objs = if need_notify {
1161 objects_to_abort
1162 .iter()
1163 .map(|object| PartialObject {
1164 oid: object.oid,
1165 obj_type: object.obj_type,
1166 schema_id: object.schema_id,
1167 database_id: object.database_id,
1168 })
1169 .collect_vec()
1170 } else {
1171 vec![]
1172 };
1173
1174 let abort_fragment_ids = fragments
1176 .iter()
1177 .map(|(fragment_id, _, _)| *fragment_id)
1178 .collect_vec();
1179
1180 Object::delete_by_id(job_id).exec(&txn).await?;
1184
1185 let err = if is_cancelled {
1186 MetaError::cancelled(format!("streaming job {job_id} is cancelled"))
1187 } else {
1188 MetaError::catalog_id_not_found("stream job", format!("streaming job {job_id} failed"))
1189 };
1190 let abort_reason = format!("streaming job aborted {}", err.as_report());
1191 for aborted_job_id in streaming_job_ids_to_abort {
1192 for tx in inner
1193 .creating_table_finish_notifier
1194 .get_mut(&database_id)
1195 .map(|creating_tables| creating_tables.remove(&aborted_job_id).into_iter())
1196 .into_iter()
1197 .flatten()
1198 .flatten()
1199 {
1200 let _ = tx.send(Err(abort_reason.clone()));
1201 }
1202 }
1203 txn.commit().await?;
1204
1205 self.env
1207 .notification_manager()
1208 .notify_serving_fragment_mapping_delete(
1209 abort_fragment_ids.iter().map(|id| *id as _).collect(),
1210 );
1211
1212 if !objs.is_empty() {
1213 self.notify_frontend(Operation::Delete, build_object_group_for_delete(objs))
1216 .await;
1217 }
1218 let aborted_sink_ids = objects_to_abort
1219 .iter()
1220 .filter(|object| object.obj_type == ObjectType::Sink)
1221 .map(|object| object.oid.as_sink_id())
1222 .collect();
1223 Ok(AbortCreatingJobResult {
1224 aborted: true,
1225 database_id: Some(database_id),
1226 aborted_sink_ids,
1227 cancel_info,
1228 })
1229 }
1230
1231 async fn build_creating_streaming_job_objects(
1232 txn: &DatabaseTransaction,
1233 job_id: JobId,
1234 ) -> MetaResult<Vec<PbObject>> {
1235 let job_type = Object::find_by_id(job_id)
1236 .select_only()
1237 .column(object::Column::ObjType)
1238 .into_tuple()
1239 .one(txn)
1240 .await?
1241 .ok_or_else(|| MetaError::catalog_id_not_found("streaming job", job_id))?;
1242 let streaming_job = StreamingJobModel::find_by_id(job_id)
1243 .one(txn)
1244 .await?
1245 .ok_or_else(|| MetaError::catalog_id_not_found("streaming job", job_id))?;
1246
1247 let table_objs = Table::find()
1248 .find_also_related(Object)
1249 .filter(
1250 table::Column::BelongsToJobId
1251 .eq(job_id)
1252 .or(table::Column::TableId.eq(job_id.as_mv_table_id())),
1253 )
1254 .all(txn)
1255 .await?;
1256 let associated_source_id = table_objs
1257 .iter()
1258 .find(|(table, _)| table.table_id == job_id.as_mv_table_id())
1259 .and_then(|(table, _)| table.optional_associated_source_id);
1260 let mut objects = table_objs
1261 .into_iter()
1262 .map(|(table, obj)| PbObject {
1263 object_info: Some(PbObjectInfo::Table(
1264 ObjectModel(table, obj.unwrap(), Some(streaming_job.clone())).into(),
1265 )),
1266 })
1267 .collect_vec();
1268
1269 match job_type {
1270 ObjectType::Table => {
1271 if let Some(source_id) = associated_source_id {
1272 let (source, obj) = Source::find_by_id(source_id)
1273 .find_also_related(Object)
1274 .one(txn)
1275 .await?
1276 .ok_or_else(|| MetaError::catalog_id_not_found("source", source_id))?;
1277 objects.push(PbObject {
1278 object_info: Some(PbObjectInfo::Source(
1279 ObjectModel(source, obj.unwrap(), None).into(),
1280 )),
1281 });
1282 }
1283 }
1284 ObjectType::Sink => {
1285 let (sink, obj) = Sink::find_by_id(job_id.as_sink_id())
1286 .find_also_related(Object)
1287 .one(txn)
1288 .await?
1289 .ok_or_else(|| MetaError::catalog_id_not_found("sink", job_id))?;
1290 objects.push(PbObject {
1291 object_info: Some(PbObjectInfo::Sink(
1292 ObjectModel(sink, obj.unwrap(), Some(streaming_job)).into(),
1293 )),
1294 });
1295 }
1296 ObjectType::Index => {
1297 let (index, obj) = Index::find_by_id(job_id.as_index_id())
1298 .find_also_related(Object)
1299 .one(txn)
1300 .await?
1301 .ok_or_else(|| MetaError::catalog_id_not_found("index", job_id))?;
1302 objects.push(PbObject {
1303 object_info: Some(PbObjectInfo::Index(
1304 ObjectModel(index, obj.unwrap(), Some(streaming_job)).into(),
1305 )),
1306 });
1307 }
1308 ObjectType::Source => {
1309 let (source, obj) = Source::find_by_id(job_id.as_shared_source_id())
1310 .find_also_related(Object)
1311 .one(txn)
1312 .await?
1313 .ok_or_else(|| MetaError::catalog_id_not_found("source", job_id))?;
1314 objects.push(PbObject {
1315 object_info: Some(PbObjectInfo::Source(
1316 ObjectModel(source, obj.unwrap(), None).into(),
1317 )),
1318 });
1319 }
1320 _ => unreachable!("invalid streaming job type: {job_type:?}"),
1321 }
1322
1323 Ok(objects)
1324 }
1325
1326 #[await_tree::instrument]
1329 pub async fn post_collect_job_fragments(
1330 &self,
1331 job_id: JobId,
1332 upstream_fragment_new_downstreams: &FragmentDownstreamRelation,
1333 new_sink_downstream: Option<FragmentDownstreamRelation>,
1334 split_assignment: Option<&SplitAssignment>,
1335 replace_sink: Option<&SinkId>,
1336 notify_creating: bool,
1337 ) -> MetaResult<Option<Vec<TableId>>> {
1338 let mut inner = self.inner.write().await;
1339 let txn = inner.db.begin().await?;
1340 let mut replace_sink_post_collect = None;
1341
1342 insert_fragment_relations(&txn, upstream_fragment_new_downstreams).await?;
1343
1344 if let Some(new_downstream) = new_sink_downstream {
1345 insert_fragment_relations(&txn, &new_downstream).await?;
1346 }
1347
1348 let mut streaming_job_model = streaming_job::ActiveModel {
1350 job_id: Set(job_id),
1351 job_status: Set(JobStatus::Creating),
1352 ..Default::default()
1353 };
1354 if replace_sink.is_some() {
1355 streaming_job_model.create_type = Set(CreateType::Background);
1359 }
1360 StreamingJobModel::update(streaming_job_model)
1361 .exec(&txn)
1362 .await?;
1363
1364 if let Some(split_assignment) = split_assignment {
1365 let fragment_splits = split_assignment
1366 .iter()
1367 .map(|(fragment_id, splits)| {
1368 (
1369 *fragment_id as _,
1370 splits.values().flatten().cloned().collect_vec(),
1371 )
1372 })
1373 .collect();
1374
1375 self.update_fragment_splits(&txn, &fragment_splits).await?;
1376 }
1377
1378 if let Some(old_sink_id) = replace_sink {
1379 let old_sink_id = *old_sink_id;
1380
1381 let (old_sink, old_sink_object) = Sink::find_by_id(old_sink_id)
1382 .find_also_related(Object)
1383 .one(&txn)
1384 .await?
1385 .and_then(|(sink, object)| object.map(|object| (sink, object)))
1386 .ok_or_else(|| MetaError::catalog_id_not_found("sink", old_sink_id))?;
1387 if old_sink_object.obj_type != ObjectType::Sink {
1388 bail!("object {} is not a sink", old_sink_id);
1389 }
1390 let final_sink_name = old_sink.name.clone();
1391
1392 let mut old_objects_to_delete = vec![old_sink_object];
1393 old_objects_to_delete
1394 .extend(get_belong_objects(&txn, old_sink_id.as_object_id()).await?);
1395 let old_object_models = load_object_models(&txn, &old_objects_to_delete).await?;
1396 let old_state_table_ids = old_object_models
1397 .iter()
1398 .filter_map(|object_info| match object_info {
1399 PbObjectInfo::Table(table) => Some(table.id),
1400 _ => None,
1401 })
1402 .collect_vec();
1403 let old_fragment_ids: Vec<FragmentId> = Fragment::find()
1404 .select_only()
1405 .column(fragment::Column::FragmentId)
1406 .filter(
1407 fragment::Column::JobId.is_in(
1408 old_objects_to_delete
1409 .iter()
1410 .map(|object| object.oid.as_job_id()),
1411 ),
1412 )
1413 .into_tuple()
1414 .all(&txn)
1415 .await?;
1416 let dropped_tables = old_object_models
1417 .iter()
1418 .filter_map(|object_info| match object_info {
1419 PbObjectInfo::Table(table) => Some(table.clone()),
1420 _ => None,
1421 })
1422 .collect_vec();
1423
1424 Sink::update(sink::ActiveModel {
1425 sink_id: Set(job_id.as_sink_id()),
1426 name: Set(final_sink_name),
1427 ..Default::default()
1428 })
1429 .exec(&txn)
1430 .await?;
1431
1432 let old_objects = old_objects_to_delete
1433 .into_iter()
1434 .map(PartialObject::from)
1435 .collect_vec();
1436 let old_object_ids = old_objects.iter().map(|obj| obj.oid).collect_vec();
1437 let updated_user_ids: Vec<UserId> = UserPrivilege::find()
1438 .select_only()
1439 .distinct()
1440 .column(user_privilege::Column::UserId)
1441 .filter(user_privilege::Column::Oid.is_in(old_object_ids.clone()))
1442 .into_tuple()
1443 .all(&txn)
1444 .await?;
1445
1446 let res = Object::delete_by_id(old_sink_id).exec(&txn).await?;
1448 if res.rows_affected == 0 {
1449 return Err(MetaError::catalog_id_not_found("sink", old_sink_id));
1450 }
1451
1452 let new_root_object = Object::find_by_id(job_id)
1453 .one(&txn)
1454 .await?
1455 .ok_or_else(|| MetaError::catalog_id_not_found("sink", job_id))?;
1456 let mut new_objects_to_notify = vec![new_root_object];
1457 new_objects_to_notify.extend(get_belong_objects(&txn, job_id.as_object_id()).await?);
1458 let new_objects = load_object_models(&txn, &new_objects_to_notify)
1459 .await?
1460 .into_iter()
1461 .map(|object_info| PbObject {
1462 object_info: Some(object_info),
1463 })
1464 .collect_vec();
1465 let dependencies =
1466 list_object_dependencies_by_object_id(&txn, job_id.as_object_id()).await?;
1467 let updated_user_info = list_user_info_by_ids(updated_user_ids, &txn).await?;
1468
1469 replace_sink_post_collect = Some((
1470 old_state_table_ids,
1471 old_fragment_ids,
1472 old_objects,
1473 dropped_tables,
1474 new_objects,
1475 dependencies,
1476 updated_user_info,
1477 ));
1478 }
1479
1480 let creating_objects = if notify_creating {
1481 Some(Self::build_creating_streaming_job_objects(&txn, job_id).await?)
1482 } else {
1483 None
1484 };
1485
1486 txn.commit().await?;
1487
1488 if let Some(objects) = creating_objects {
1489 self.notify_frontend(
1490 NotificationOperation::Add,
1491 NotificationInfo::ObjectGroup(PbObjectGroup {
1492 objects,
1493 dependencies: vec![],
1494 }),
1495 )
1496 .await;
1497 }
1498
1499 let Some((
1500 old_state_table_ids,
1501 old_fragment_ids,
1502 old_objects,
1503 dropped_tables,
1504 new_objects,
1505 dependencies,
1506 updated_user_info,
1507 )) = replace_sink_post_collect
1508 else {
1509 return Ok(None);
1510 };
1511
1512 inner
1513 .dropped_tables
1514 .extend(dropped_tables.into_iter().map(|table| (table.id, table)));
1515 drop(inner);
1516
1517 self.env
1518 .notification_manager()
1519 .notify_serving_fragment_mapping_delete(
1520 old_fragment_ids.iter().map(|id| *id as _).collect(),
1521 );
1522
1523 let _ = self
1524 .notify_frontend(
1525 NotificationOperation::Delete,
1526 build_object_group_for_delete(old_objects),
1527 )
1528 .await;
1529
1530 let _ = self
1531 .notify_frontend(
1532 NotificationOperation::Add,
1535 NotificationInfo::ObjectGroup(PbObjectGroup {
1536 objects: new_objects,
1537 dependencies,
1538 }),
1539 )
1540 .await;
1541
1542 if !updated_user_info.is_empty() {
1543 let _ = self.notify_users_update(updated_user_info).await;
1544 }
1545
1546 Ok(Some(old_state_table_ids))
1547 }
1548
1549 pub async fn create_job_catalog_for_replace(
1550 &self,
1551 streaming_job: &StreamingJob,
1552 ctx: Option<&StreamContext>,
1553 specified_parallelism: Option<&NonZeroUsize>,
1554 expected_original_max_parallelism: Option<usize>,
1555 ) -> MetaResult<streaming_job::Model> {
1556 let id = streaming_job.id();
1557 let inner = self.inner.write().await;
1558 let txn = inner.db.begin().await?;
1559
1560 streaming_job.verify_version_for_replace(&txn).await?;
1562 let referring_objects =
1564 get_referring_objects(id.as_object_id(), streaming_job.object_type(), &txn).await?;
1565 let referring_job_ids = referring_objects
1566 .iter()
1567 .map(|object| object.oid.as_job_id())
1568 .collect_vec();
1569 let non_created_referring_job_count = StreamingJobModel::find()
1570 .filter(
1571 streaming_job::Column::JobId
1572 .is_in(referring_job_ids)
1573 .and(streaming_job::Column::JobStatus.ne(JobStatus::Created)),
1574 )
1575 .count(&txn)
1576 .await?;
1577 if non_created_referring_job_count != 0 {
1578 return Err(MetaError::permission_denied(
1579 "job is being altered or referenced by some creating jobs",
1580 ));
1581 }
1582
1583 let batch_refresh_dep_cnt = ObjectDependency::find()
1586 .join(
1587 JoinType::InnerJoin,
1588 object_dependency::Relation::Object1.def(),
1589 )
1590 .join(JoinType::InnerJoin, object::Relation::StreamingJob.def())
1591 .filter(
1592 object_dependency::Column::Oid
1593 .eq(id)
1594 .and(streaming_job::Column::RefreshIntervalSec.is_not_null()),
1595 )
1596 .count(&txn)
1597 .await?;
1598 if batch_refresh_dep_cnt != 0 {
1599 return Err(MetaError::permission_denied(
1600 "replacing a table with dependent batch refresh materialized views is not supported",
1601 ));
1602 }
1603
1604 let original_job = StreamingJobModel::find_by_id(id)
1606 .one(&txn)
1607 .await?
1608 .map(ReplaceOriginalJobInfo::from)
1609 .ok_or_else(|| MetaError::catalog_id_not_found(streaming_job.job_type_str(), id))?;
1610
1611 if let Some(max_parallelism) = expected_original_max_parallelism
1612 && original_job.max_parallelism != max_parallelism as i32
1613 {
1614 bail!(
1617 "cannot use a different max parallelism \
1618 when replacing streaming job, \
1619 original: {}, new: {}",
1620 original_job.max_parallelism,
1621 max_parallelism
1622 );
1623 }
1624
1625 let parallelism = original_job.resolved_parallelism(specified_parallelism);
1626 let ctx = original_job.stream_context(ctx);
1627 let adaptive_parallelism_strategy = original_job
1628 .adaptive_parallelism_strategy
1629 .as_deref()
1630 .map(|s| parse_strategy(s).expect("strategy should be validated before persisting"));
1631 let resource_type = original_job.resource_type();
1632 let belong_to_oid = Object::find_by_id(id)
1633 .select_only()
1634 .column(object::Column::BelongToOid)
1635 .into_tuple::<Option<ObjectId>>()
1636 .one(&txn)
1637 .await?
1638 .ok_or_else(|| MetaError::catalog_id_not_found(streaming_job.job_type_str(), id))?;
1639
1640 let tmp_model = Self::create_streaming_job_obj(
1642 &txn,
1643 streaming_job.object_type(),
1644 streaming_job.owner() as _,
1645 belong_to_oid,
1646 streaming_job.create_type(),
1647 ctx,
1648 adaptive_parallelism_strategy,
1649 parallelism,
1650 original_job.max_parallelism as _,
1651 resource_type,
1652 None,
1656 None,
1657 None, )
1659 .await?;
1660
1661 ObjectDependency::insert(object_dependency::ActiveModel {
1663 oid: Set(id.as_object_id()),
1664 used_by: Set(tmp_model.job_id.as_object_id()),
1665 ..Default::default()
1666 })
1667 .exec(&txn)
1668 .await?;
1669
1670 txn.commit().await?;
1671
1672 Ok(tmp_model)
1673 }
1674
1675 pub async fn finish_streaming_job(&self, job_id: JobId) -> MetaResult<()> {
1677 let mut inner = self.inner.write().await;
1678 let txn = inner.db.begin().await?;
1679
1680 if check_if_belongs_to_iceberg_table(&txn, job_id).await? {
1682 tracing::info!(
1683 "streaming job {} is for iceberg table, wait for manual finish operation",
1684 job_id
1685 );
1686 return Ok(());
1687 }
1688
1689 let (notification_op, objects, updated_user_info, dependencies) =
1690 self.finish_streaming_job_inner(&txn, job_id).await?;
1691
1692 txn.commit().await?;
1693
1694 let mut version = self
1695 .notify_frontend(
1696 notification_op,
1697 NotificationInfo::ObjectGroup(PbObjectGroup {
1698 objects,
1699 dependencies,
1700 }),
1701 )
1702 .await;
1703
1704 if !updated_user_info.is_empty() {
1706 version = self.notify_users_update(updated_user_info).await;
1707 }
1708
1709 inner
1710 .creating_table_finish_notifier
1711 .values_mut()
1712 .for_each(|creating_tables| {
1713 if let Some(txs) = creating_tables.remove(&job_id) {
1714 for tx in txs {
1715 let _ = tx.send(Ok(version));
1716 }
1717 }
1718 });
1719
1720 Ok(())
1721 }
1722
1723 pub async fn finish_streaming_job_inner(
1725 &self,
1726 txn: &DatabaseTransaction,
1727 job_id: JobId,
1728 ) -> MetaResult<(
1729 Operation,
1730 Vec<risingwave_pb::meta::Object>,
1731 Vec<PbUserInfo>,
1732 Vec<PbObjectDependency>,
1733 )> {
1734 let job_type = Object::find_by_id(job_id)
1735 .select_only()
1736 .column(object::Column::ObjType)
1737 .into_tuple()
1738 .one(txn)
1739 .await?
1740 .ok_or_else(|| MetaError::catalog_id_not_found("streaming job", job_id))?;
1741
1742 let res = Object::update_many()
1744 .col_expr(object::Column::CreatedAt, Expr::current_timestamp().into())
1745 .col_expr(
1746 object::Column::CreatedAtClusterVersion,
1747 current_cluster_version().into(),
1748 )
1749 .filter(object::Column::Oid.eq(job_id))
1750 .exec(txn)
1751 .await?;
1752 if res.rows_affected == 0 {
1753 return Err(MetaError::catalog_id_not_found("streaming job", job_id));
1754 }
1755
1756 let job = streaming_job::ActiveModel {
1758 job_id: Set(job_id),
1759 job_status: Set(JobStatus::Created),
1760 ..Default::default()
1761 };
1762 let streaming_job = Some(job.update(txn).await?);
1763
1764 let internal_table_objs = Table::find()
1766 .find_also_related(Object)
1767 .filter(table::Column::BelongsToJobId.eq(job_id))
1768 .all(txn)
1769 .await?;
1770 let mut objects = internal_table_objs
1771 .iter()
1772 .map(|(table, obj)| PbObject {
1773 object_info: Some(PbObjectInfo::Table(
1774 ObjectModel(table.clone(), obj.clone().unwrap(), streaming_job.clone()).into(),
1775 )),
1776 })
1777 .collect_vec();
1778 let notification_op = NotificationOperation::Update;
1779 let mut updated_user_info = vec![];
1780 let mut need_grant_default_privileges = true;
1781
1782 match job_type {
1783 ObjectType::Table => {
1784 let (table, obj) = Table::find_by_id(job_id.as_mv_table_id())
1785 .find_also_related(Object)
1786 .one(txn)
1787 .await?
1788 .ok_or_else(|| MetaError::catalog_id_not_found("table", job_id))?;
1789 if let Some(source_id) = table.optional_associated_source_id {
1790 let (src, obj) = Source::find_by_id(source_id)
1791 .find_also_related(Object)
1792 .one(txn)
1793 .await?
1794 .ok_or_else(|| MetaError::catalog_id_not_found("source", source_id))?;
1795 objects.push(PbObject {
1796 object_info: Some(PbObjectInfo::Source(
1797 ObjectModel(src, obj.unwrap(), None).into(),
1798 )),
1799 });
1800 }
1801 objects.push(PbObject {
1802 object_info: Some(PbObjectInfo::Table(
1803 ObjectModel(table, obj.unwrap(), streaming_job).into(),
1804 )),
1805 });
1806 }
1807 ObjectType::Sink => {
1808 let (sink, obj) = Sink::find_by_id(job_id.as_sink_id())
1809 .find_also_related(Object)
1810 .one(txn)
1811 .await?
1812 .ok_or_else(|| MetaError::catalog_id_not_found("sink", job_id))?;
1813 if sink.name.starts_with(ICEBERG_SINK_PREFIX) {
1814 need_grant_default_privileges = false;
1815 }
1816 objects.push(PbObject {
1817 object_info: Some(PbObjectInfo::Sink(
1818 ObjectModel(sink, obj.unwrap(), streaming_job).into(),
1819 )),
1820 });
1821 }
1822 ObjectType::Index => {
1823 need_grant_default_privileges = false;
1824 let (index, obj) = Index::find_by_id(job_id.as_index_id())
1825 .find_also_related(Object)
1826 .one(txn)
1827 .await?
1828 .ok_or_else(|| MetaError::catalog_id_not_found("index", job_id))?;
1829 {
1830 let (table, obj) = Table::find_by_id(index.index_table_id)
1831 .find_also_related(Object)
1832 .one(txn)
1833 .await?
1834 .ok_or_else(|| {
1835 MetaError::catalog_id_not_found("table", index.index_table_id)
1836 })?;
1837 objects.push(PbObject {
1838 object_info: Some(PbObjectInfo::Table(
1839 ObjectModel(table, obj.unwrap(), streaming_job.clone()).into(),
1840 )),
1841 });
1842 }
1843
1844 let primary_table_privileges = UserPrivilege::find()
1847 .filter(
1848 user_privilege::Column::Oid
1849 .eq(index.primary_table_id)
1850 .and(user_privilege::Column::Action.eq(Action::Select)),
1851 )
1852 .all(txn)
1853 .await?;
1854 if !primary_table_privileges.is_empty() {
1855 let index_state_table_ids: Vec<TableId> = Table::find()
1856 .select_only()
1857 .column(table::Column::TableId)
1858 .filter(
1859 table::Column::BelongsToJobId
1860 .eq(job_id)
1861 .or(table::Column::TableId.eq(index.index_table_id)),
1862 )
1863 .into_tuple()
1864 .all(txn)
1865 .await?;
1866 let mut new_privileges = vec![];
1867 for privilege in &primary_table_privileges {
1868 for state_table_id in &index_state_table_ids {
1869 new_privileges.push(user_privilege::ActiveModel {
1870 id: Default::default(),
1871 oid: Set(state_table_id.as_object_id()),
1872 user_id: Set(privilege.user_id),
1873 action: Set(Action::Select),
1874 dependent_id: Set(privilege.dependent_id),
1875 granted_by: Set(privilege.granted_by),
1876 with_grant_option: Set(privilege.with_grant_option),
1877 });
1878 }
1879 }
1880 upsert_user_privileges(txn, new_privileges).await?;
1881
1882 updated_user_info = list_user_info_by_ids(
1883 primary_table_privileges.into_iter().map(|p| p.user_id),
1884 txn,
1885 )
1886 .await?;
1887 }
1888
1889 objects.push(PbObject {
1890 object_info: Some(PbObjectInfo::Index(
1891 ObjectModel(index, obj.unwrap(), streaming_job).into(),
1892 )),
1893 });
1894 }
1895 ObjectType::Source => {
1896 let (source, obj) = Source::find_by_id(job_id.as_shared_source_id())
1897 .find_also_related(Object)
1898 .one(txn)
1899 .await?
1900 .ok_or_else(|| MetaError::catalog_id_not_found("source", job_id))?;
1901 objects.push(PbObject {
1902 object_info: Some(PbObjectInfo::Source(
1903 ObjectModel(source, obj.unwrap(), None).into(),
1904 )),
1905 });
1906 }
1907 _ => unreachable!("invalid job type: {:?}", job_type),
1908 }
1909
1910 if need_grant_default_privileges {
1911 updated_user_info = grant_default_privileges_automatically(txn, job_id).await?;
1912 }
1913
1914 let dependencies =
1915 list_object_dependencies_by_object_id(txn, job_id.as_object_id()).await?;
1916
1917 Ok((notification_op, objects, updated_user_info, dependencies))
1918 }
1919
1920 pub async fn finish_replace_streaming_job(
1921 &self,
1922 tmp_id: JobId,
1923 streaming_job: StreamingJob,
1924 replace_upstream: FragmentReplaceUpstream,
1925 sink_into_table_context: SinkIntoTableContext,
1926 drop_table_connector_ctx: Option<&DropTableConnectorContext>,
1927 auto_refresh_schema_sinks: Option<Vec<FinishAutoRefreshSchemaSinkContext>>,
1928 ) -> MetaResult<NotificationVersion> {
1929 let inner = self.inner.write().await;
1930 let txn = inner.db.begin().await?;
1931
1932 let (objects, delete_notification_objs, old_fragment_ids, new_fragment_ids) =
1933 Self::finish_replace_streaming_job_inner(
1934 tmp_id,
1935 replace_upstream,
1936 sink_into_table_context,
1937 &txn,
1938 streaming_job,
1939 drop_table_connector_ctx,
1940 auto_refresh_schema_sinks,
1941 )
1942 .await?;
1943
1944 txn.commit().await?;
1945
1946 let notification_manager = self.env.notification_manager();
1948 notification_manager.notify_serving_fragment_mapping_delete(
1949 old_fragment_ids.iter().map(|id| *id as _).collect(),
1950 );
1951 notification_manager.notify_serving_fragment_mapping_update(
1952 new_fragment_ids.iter().map(|id| *id as _).collect(),
1953 );
1954
1955 let mut version = self
1956 .notify_frontend(
1957 NotificationOperation::Update,
1958 NotificationInfo::ObjectGroup(PbObjectGroup {
1959 objects,
1960 dependencies: vec![],
1961 }),
1962 )
1963 .await;
1964
1965 if let Some((user_infos, to_drop_objects)) = delete_notification_objs {
1966 self.notify_users_update(user_infos).await;
1967 version = self
1968 .notify_frontend(
1969 NotificationOperation::Delete,
1970 build_object_group_for_delete(to_drop_objects),
1971 )
1972 .await;
1973 }
1974
1975 Ok(version)
1976 }
1977
1978 fn update_iceberg_source_columns(
1979 original_source_columns: &[ColumnCatalog],
1980 original_row_id_index: Option<usize>,
1981 new_table_columns: &[ColumnCatalog],
1982 ) -> (Vec<ColumnCatalog>, Option<usize>) {
1983 let row_id_column_name = original_row_id_index
1984 .and_then(|idx| original_source_columns.get(idx))
1985 .map(|col| col.name().to_owned());
1986 let mut next_column_id = max_column_id(original_source_columns).next();
1987 let existing_columns: HashMap<String, ColumnCatalog> = original_source_columns
1988 .iter()
1989 .cloned()
1990 .map(|col| (col.name().to_owned(), col))
1991 .collect();
1992
1993 let mut new_columns = Vec::new();
1994 for table_col in new_table_columns
1995 .iter()
1996 .filter(|col| !col.is_rw_sys_column())
1997 {
1998 let mut source_col_name = table_col.name().to_owned();
1999 if source_col_name == ROW_ID_COLUMN_NAME {
2000 source_col_name = RISINGWAVE_ICEBERG_ROW_ID.to_owned();
2001 }
2002
2003 if let Some(existing) = existing_columns.get(&source_col_name) {
2004 new_columns.push(existing.clone());
2005 } else {
2006 let mut new_col = table_col.clone();
2007 new_col.column_desc.name = source_col_name;
2008 new_col.column_desc.column_id = next_column_id;
2009 next_column_id = next_column_id.next();
2010 new_columns.push(new_col);
2011 }
2012 }
2013
2014 let mut seen_names: HashSet<String> = new_columns
2015 .iter()
2016 .map(|col| col.name().to_owned())
2017 .collect();
2018 for col in original_source_columns
2019 .iter()
2020 .filter(|col| col.is_iceberg_hidden_column())
2021 {
2022 if seen_names.insert(col.name().to_owned()) {
2023 new_columns.push(col.clone());
2024 }
2025 }
2026
2027 let new_row_id_index = row_id_column_name
2028 .as_ref()
2029 .and_then(|name| new_columns.iter().position(|col| col.name() == name));
2030
2031 (new_columns, new_row_id_index)
2032 }
2033
2034 pub async fn finish_replace_streaming_job_inner(
2035 tmp_id: JobId,
2036 replace_upstream: FragmentReplaceUpstream,
2037 SinkIntoTableContext {
2038 updated_sink_catalogs,
2039 }: SinkIntoTableContext,
2040 txn: &DatabaseTransaction,
2041 streaming_job: StreamingJob,
2042 drop_table_connector_ctx: Option<&DropTableConnectorContext>,
2043 auto_refresh_schema_sinks: Option<Vec<FinishAutoRefreshSchemaSinkContext>>,
2044 ) -> MetaResult<(
2045 Vec<PbObject>,
2046 Option<(Vec<PbUserInfo>, Vec<PartialObject>)>,
2047 Vec<FragmentId>,
2048 Vec<FragmentId>,
2049 )> {
2050 let original_job_id = streaming_job.id();
2051 let job_type = streaming_job.job_type();
2052
2053 let old_fragment_ids: Vec<FragmentId> = Fragment::find()
2055 .select_only()
2056 .column(fragment::Column::FragmentId)
2057 .filter(fragment::Column::JobId.eq(original_job_id))
2058 .into_tuple()
2059 .all(txn)
2060 .await?;
2061 let new_fragment_ids: Vec<FragmentId> = Fragment::find()
2062 .select_only()
2063 .column(fragment::Column::FragmentId)
2064 .filter(fragment::Column::JobId.eq(tmp_id))
2065 .into_tuple()
2066 .all(txn)
2067 .await?;
2068
2069 let mut index_item_rewriter = None;
2070 let mut updated_iceberg_source_id: Option<SourceId> = None;
2071
2072 match streaming_job {
2074 StreamingJob::Table(_, table, _table_job_type) => {
2075 let original_column_catalogs =
2076 get_table_columns(txn, original_job_id.as_mv_table_id()).await?;
2077 let schema_changed = original_column_catalogs.to_protobuf() != table.columns;
2078 let is_iceberg = table
2079 .engine
2080 .and_then(|engine| PbEngine::try_from(engine).ok())
2081 == Some(PbEngine::Iceberg);
2082 if is_iceberg && schema_changed {
2083 let iceberg_source_name = format!("{}{}", ICEBERG_SOURCE_PREFIX, table.name);
2084 let source = Source::find()
2085 .inner_join(Object)
2086 .filter(
2087 object::Column::DatabaseId
2088 .eq(table.database_id)
2089 .and(object::Column::SchemaId.eq(table.schema_id))
2090 .and(source::Column::Name.eq(&iceberg_source_name)),
2091 )
2092 .one(txn)
2093 .await?
2094 .ok_or_else(|| {
2095 MetaError::catalog_id_not_found("source", iceberg_source_name)
2096 })?;
2097
2098 let source_id = source.source_id;
2099 let source_version = source.version;
2100 let source_columns: Vec<ColumnCatalog> = source
2101 .columns
2102 .to_protobuf()
2103 .into_iter()
2104 .map(ColumnCatalog::from)
2105 .collect();
2106 let source_row_id_index = source
2107 .row_id_index
2108 .and_then(|idx| usize::try_from(idx).ok());
2109 let table_columns: Vec<ColumnCatalog> = table
2110 .columns
2111 .iter()
2112 .cloned()
2113 .map(ColumnCatalog::from)
2114 .collect();
2115 let (updated_columns, updated_row_id_index) =
2116 Self::update_iceberg_source_columns(
2117 &source_columns,
2118 source_row_id_index,
2119 &table_columns,
2120 );
2121 let updated_columns: Vec<PbColumnCatalog> = updated_columns
2122 .into_iter()
2123 .map(|col| col.to_protobuf())
2124 .collect();
2125
2126 let mut source_active = source.into_active_model();
2127 source_active.columns = Set(ColumnCatalogArray::from(updated_columns));
2128 source_active.row_id_index = Set(updated_row_id_index.map(|idx| idx as i32));
2129 source_active.version = Set(source_version + 1);
2130 source_active.update(txn).await?;
2131 updated_iceberg_source_id = Some(source_id);
2132 }
2133
2134 index_item_rewriter = Some({
2135 let original_columns = original_column_catalogs
2136 .to_protobuf()
2137 .into_iter()
2138 .map(|c| c.column_desc.unwrap())
2139 .collect_vec();
2140 let new_columns = table
2141 .columns
2142 .iter()
2143 .map(|c| c.column_desc.clone().unwrap())
2144 .collect_vec();
2145
2146 IndexItemRewriter {
2147 original_columns,
2148 new_columns,
2149 }
2150 });
2151
2152 for sink_id in updated_sink_catalogs {
2154 Sink::update(sink::ActiveModel {
2155 sink_id: Set(sink_id as _),
2156 original_target_columns: Set(Some(original_column_catalogs.clone())),
2157 ..Default::default()
2158 })
2159 .exec(txn)
2160 .await?;
2161 }
2162 let mut table = table::ActiveModel::from(table);
2164 if let Some(drop_table_connector_ctx) = drop_table_connector_ctx
2165 && drop_table_connector_ctx.to_change_streaming_job_id == original_job_id
2166 {
2167 table.optional_associated_source_id = Set(None);
2169 }
2170
2171 Table::update(table).exec(txn).await?;
2172 }
2173 StreamingJob::Source(source) => {
2174 let source = source::ActiveModel::from(source);
2176 Source::update(source).exec(txn).await?;
2177 }
2178 StreamingJob::MaterializedView(table) => {
2179 let table = table::ActiveModel::from(table);
2181 Table::update(table).exec(txn).await?;
2182 }
2183 _ => unreachable!(
2184 "invalid streaming job type: {:?}",
2185 streaming_job.job_type_str()
2186 ),
2187 }
2188
2189 async fn finish_fragments(
2190 txn: &DatabaseTransaction,
2191 tmp_id: JobId,
2192 original_job_id: JobId,
2193 replace_upstream: FragmentReplaceUpstream,
2194 ) -> MetaResult<()> {
2195 let fragment_info: Vec<(FragmentId, I32Array)> = Fragment::find()
2199 .select_only()
2200 .columns([
2201 fragment::Column::FragmentId,
2202 fragment::Column::StateTableIds,
2203 ])
2204 .filter(fragment::Column::JobId.eq(tmp_id))
2205 .into_tuple()
2206 .all(txn)
2207 .await?;
2208 for (fragment_id, state_table_ids) in fragment_info {
2209 for state_table_id in state_table_ids.into_inner() {
2210 let state_table_id = TableId::new(state_table_id as _);
2211 Table::update(table::ActiveModel {
2212 table_id: Set(state_table_id),
2213 fragment_id: Set(Some(fragment_id)),
2214 ..Default::default()
2216 })
2217 .exec(txn)
2218 .await?;
2219 }
2220 }
2221
2222 Fragment::delete_many()
2224 .filter(fragment::Column::JobId.eq(original_job_id))
2225 .exec(txn)
2226 .await?;
2227 Fragment::update_many()
2228 .col_expr(fragment::Column::JobId, SimpleExpr::from(original_job_id))
2229 .filter(fragment::Column::JobId.eq(tmp_id))
2230 .exec(txn)
2231 .await?;
2232
2233 for (fragment_id, fragment_replace_map) in replace_upstream {
2236 let (fragment_id, mut stream_node) =
2237 Fragment::find_by_id(fragment_id as FragmentId)
2238 .select_only()
2239 .columns([fragment::Column::FragmentId, fragment::Column::StreamNode])
2240 .into_tuple::<(FragmentId, StreamNode)>()
2241 .one(txn)
2242 .await?
2243 .map(|(id, node)| (id, node.to_protobuf()))
2244 .ok_or_else(|| MetaError::catalog_id_not_found("fragment", fragment_id))?;
2245
2246 visit_stream_node_mut(&mut stream_node, |body| {
2247 if let PbNodeBody::Merge(m) = body
2248 && let Some(new_fragment_id) =
2249 fragment_replace_map.get(&m.upstream_fragment_id)
2250 {
2251 m.upstream_fragment_id = *new_fragment_id;
2252 }
2253 });
2254 Fragment::update(fragment::ActiveModel {
2255 fragment_id: Set(fragment_id),
2256 stream_node: Set(StreamNode::from(&stream_node)),
2257 ..Default::default()
2258 })
2259 .exec(txn)
2260 .await?;
2261 }
2262
2263 Object::delete_by_id(tmp_id).exec(txn).await?;
2265
2266 Ok(())
2267 }
2268
2269 finish_fragments(txn, tmp_id, original_job_id, replace_upstream).await?;
2270
2271 let mut objects = vec![];
2273 match job_type {
2274 StreamingJobType::Table(_) | StreamingJobType::MaterializedView => {
2275 let (table, table_obj) = Table::find_by_id(original_job_id.as_mv_table_id())
2276 .find_also_related(Object)
2277 .one(txn)
2278 .await?
2279 .ok_or_else(|| MetaError::catalog_id_not_found("object", original_job_id))?;
2280 let streaming_job = streaming_job::Entity::find_by_id(table.job_id())
2281 .one(txn)
2282 .await?;
2283 objects.push(PbObject {
2284 object_info: Some(PbObjectInfo::Table(
2285 ObjectModel(table, table_obj.unwrap(), streaming_job).into(),
2286 )),
2287 })
2288 }
2289 StreamingJobType::Source => {
2290 let (source, source_obj) =
2291 Source::find_by_id(original_job_id.as_shared_source_id())
2292 .find_also_related(Object)
2293 .one(txn)
2294 .await?
2295 .ok_or_else(|| {
2296 MetaError::catalog_id_not_found("object", original_job_id)
2297 })?;
2298 objects.push(PbObject {
2299 object_info: Some(PbObjectInfo::Source(
2300 ObjectModel(source, source_obj.unwrap(), None).into(),
2301 )),
2302 })
2303 }
2304 _ => unreachable!("invalid streaming job type for replace: {:?}", job_type),
2305 }
2306
2307 if let Some(source_id) = updated_iceberg_source_id {
2308 let (source, source_obj) = Source::find_by_id(source_id)
2309 .find_also_related(Object)
2310 .one(txn)
2311 .await?
2312 .ok_or_else(|| MetaError::catalog_id_not_found("source", source_id))?;
2313 objects.push(PbObject {
2314 object_info: Some(PbObjectInfo::Source(
2315 ObjectModel(source, source_obj.unwrap(), None).into(),
2316 )),
2317 });
2318 }
2319
2320 if let Some(expr_rewriter) = index_item_rewriter {
2321 let index_items: Vec<(IndexId, ExprNodeArray)> = Index::find()
2322 .select_only()
2323 .columns([index::Column::IndexId, index::Column::IndexItems])
2324 .filter(index::Column::PrimaryTableId.eq(original_job_id))
2325 .into_tuple()
2326 .all(txn)
2327 .await?;
2328 for (index_id, nodes) in index_items {
2329 let mut pb_nodes = nodes.to_protobuf();
2330 pb_nodes
2331 .iter_mut()
2332 .for_each(|x| expr_rewriter.rewrite_expr(x));
2333 let index = index::ActiveModel {
2334 index_id: Set(index_id),
2335 index_items: Set(pb_nodes.into()),
2336 ..Default::default()
2337 }
2338 .update(txn)
2339 .await?;
2340 let (index_obj, streaming_job) = Object::find_by_id(index.index_id)
2341 .find_also_related(streaming_job::Entity)
2342 .one(txn)
2343 .await?
2344 .ok_or_else(|| MetaError::catalog_id_not_found("object", index.index_id))?;
2345 objects.push(PbObject {
2346 object_info: Some(PbObjectInfo::Index(
2347 ObjectModel(index, index_obj, streaming_job).into(),
2348 )),
2349 });
2350 }
2351 }
2352
2353 if let Some(sinks) = auto_refresh_schema_sinks {
2354 for finish_sink_context in sinks {
2355 finish_fragments(
2356 txn,
2357 finish_sink_context.tmp_sink_id.as_job_id(),
2358 finish_sink_context.original_sink_id.as_job_id(),
2359 Default::default(),
2360 )
2361 .await?;
2362 let (mut sink, sink_obj) = Sink::find_by_id(finish_sink_context.original_sink_id)
2363 .find_also_related(Object)
2364 .one(txn)
2365 .await?
2366 .ok_or_else(|| MetaError::catalog_id_not_found("sink", original_job_id))?;
2367 let sink_streaming_job =
2368 streaming_job::Entity::find_by_id(sink.sink_id.as_job_id())
2369 .one(txn)
2370 .await?;
2371 let columns = ColumnCatalogArray::from(finish_sink_context.columns);
2372 Sink::update(sink::ActiveModel {
2373 sink_id: Set(finish_sink_context.original_sink_id),
2374 columns: Set(columns.clone()),
2375 ..Default::default()
2376 })
2377 .exec(txn)
2378 .await?;
2379 sink.columns = columns;
2380 objects.push(PbObject {
2381 object_info: Some(PbObjectInfo::Sink(
2382 ObjectModel(sink, sink_obj.unwrap(), sink_streaming_job.clone()).into(),
2383 )),
2384 });
2385 if let Some(new_log_store_table) = finish_sink_context.new_log_store_table {
2386 let log_store_table_id = new_log_store_table.id;
2387 let new_log_store_table_columns: ColumnCatalogArray =
2388 new_log_store_table.columns.clone().into();
2389 let new_log_store_table_value_indices =
2390 new_log_store_table.value_indices.clone();
2391 let (mut table, table_obj) = Table::find_by_id(log_store_table_id)
2392 .find_also_related(Object)
2393 .one(txn)
2394 .await?
2395 .ok_or_else(|| MetaError::catalog_id_not_found("table", original_job_id))?;
2396 Table::update(table::ActiveModel {
2397 table_id: Set(log_store_table_id),
2398 columns: Set(new_log_store_table_columns.clone()),
2399 value_indices: Set(new_log_store_table_value_indices.clone().into()),
2400 ..Default::default()
2401 })
2402 .exec(txn)
2403 .await?;
2404 table.columns = new_log_store_table_columns;
2405 table.value_indices = new_log_store_table_value_indices.into();
2406 objects.push(PbObject {
2407 object_info: Some(PbObjectInfo::Table(
2408 ObjectModel(table, table_obj.unwrap(), sink_streaming_job.clone())
2409 .into(),
2410 )),
2411 });
2412 }
2413 }
2414 }
2415
2416 let mut notification_objs: Option<(Vec<PbUserInfo>, Vec<PartialObject>)> = None;
2417 if let Some(drop_table_connector_ctx) = drop_table_connector_ctx {
2418 notification_objs =
2419 Some(Self::drop_table_associated_source(txn, drop_table_connector_ctx).await?);
2420 }
2421
2422 Ok((
2423 objects,
2424 notification_objs,
2425 old_fragment_ids,
2426 new_fragment_ids,
2427 ))
2428 }
2429
2430 pub async fn try_abort_replacing_streaming_job(
2432 &self,
2433 tmp_job_id: JobId,
2434 tmp_sink_ids: Option<Vec<ObjectId>>,
2435 ) -> MetaResult<()> {
2436 let inner = self.inner.write().await;
2437 let txn = inner.db.begin().await?;
2438
2439 let mut root_ids: Vec<ObjectId> = vec![tmp_job_id.into()];
2440 if let Some(sink_ids) = &tmp_sink_ids {
2441 root_ids.extend(sink_ids.iter().copied());
2442 }
2443 let mut temporary_objects = Object::find()
2444 .filter(object::Column::Oid.is_in(root_ids.iter().copied()))
2445 .all(&txn)
2446 .await?;
2447 temporary_objects.extend(get_belong_objects_by_ids(&txn, root_ids.iter().copied()).await?);
2448
2449 let abort_fragment_ids: Vec<FragmentId> = Fragment::find()
2452 .select_only()
2453 .column(fragment::Column::FragmentId)
2454 .filter(
2455 fragment::Column::JobId.is_in(
2456 temporary_objects
2457 .iter()
2458 .map(|object| object.oid.as_job_id()),
2459 ),
2460 )
2461 .into_tuple()
2462 .all(&txn)
2463 .await?;
2464
2465 Object::delete_many()
2466 .filter(object::Column::Oid.is_in(root_ids))
2467 .exec(&txn)
2468 .await?;
2469 txn.commit().await?;
2470
2471 self.env
2473 .notification_manager()
2474 .notify_serving_fragment_mapping_delete(
2475 abort_fragment_ids.iter().map(|id| *id as _).collect(),
2476 );
2477
2478 Ok(())
2479 }
2480
2481 pub async fn update_source_rate_limit_by_source_id(
2484 &self,
2485 source_id: SourceId,
2486 rate_limit: Option<u32>,
2487 ) -> MetaResult<HashMap<FragmentId, PbStreamNode>> {
2488 let inner = self.inner.read().await;
2489 let txn = inner.db.begin().await?;
2490
2491 {
2492 let active_source = source::ActiveModel {
2493 source_id: Set(source_id),
2494 rate_limit: Set(rate_limit.map(|v| v as i32)),
2495 ..Default::default()
2496 };
2497 Source::update(active_source).exec(&txn).await?;
2498 }
2499
2500 let (source, obj) = Source::find_by_id(source_id)
2501 .find_also_related(Object)
2502 .one(&txn)
2503 .await?
2504 .ok_or_else(|| {
2505 MetaError::catalog_id_not_found(ObjectType::Source.as_str(), source_id)
2506 })?;
2507
2508 let streaming_job_ids: Vec<JobId> =
2509 if let Some(table_id) = source.optional_associated_table_id {
2510 vec![table_id.as_job_id()]
2511 } else if let Some(source_info) = &source.source_info
2512 && source_info.to_protobuf().is_shared()
2513 {
2514 vec![source_id.as_share_source_job_id()]
2515 } else {
2516 ObjectDependency::find()
2517 .select_only()
2518 .column(object_dependency::Column::UsedBy)
2519 .filter(object_dependency::Column::Oid.eq(source_id))
2520 .into_tuple()
2521 .all(&txn)
2522 .await?
2523 };
2524
2525 if streaming_job_ids.is_empty() {
2526 return Err(MetaError::invalid_parameter(format!(
2527 "source id {source_id} not used by any streaming job"
2528 )));
2529 }
2530
2531 let fragments: Vec<(FragmentId, JobId, i32, StreamNode)> = Fragment::find()
2532 .select_only()
2533 .columns([
2534 fragment::Column::FragmentId,
2535 fragment::Column::JobId,
2536 fragment::Column::FragmentTypeMask,
2537 fragment::Column::StreamNode,
2538 ])
2539 .filter(fragment::Column::JobId.is_in(streaming_job_ids))
2540 .into_tuple()
2541 .all(&txn)
2542 .await?;
2543 let mut fragments = fragments
2544 .into_iter()
2545 .map(|(id, job_id, mask, stream_node)| {
2546 (
2547 id,
2548 job_id,
2549 FragmentTypeMask::from(mask as u32),
2550 stream_node.to_protobuf(),
2551 )
2552 })
2553 .collect_vec();
2554
2555 fragments.retain_mut(|(_, _, fragment_type_mask, stream_node)| {
2556 let mut found = false;
2557 if fragment_type_mask.contains(FragmentTypeFlag::Source) {
2558 visit_stream_node_mut(stream_node, |node| {
2559 if let PbNodeBody::Source(node) = node
2560 && let Some(node_inner) = &mut node.source_inner
2561 && node_inner.source_id == source_id
2562 {
2563 node_inner.rate_limit = rate_limit;
2564 found = true;
2565 }
2566 });
2567 }
2568 visit_stream_node_mut(stream_node, |node| {
2571 if let PbNodeBody::StreamFsFetch(node) = node {
2572 fragment_type_mask.add(FragmentTypeFlag::FsFetch);
2573 if let Some(node_inner) = &mut node.node_inner
2574 && node_inner.source_id == source_id
2575 {
2576 node_inner.rate_limit = rate_limit;
2577 found = true;
2578 }
2579 }
2580 });
2581 found
2582 });
2583
2584 assert!(
2585 !fragments.is_empty(),
2586 "source id should be used by at least one fragment"
2587 );
2588
2589 let fragment_nodes = fragments
2590 .iter()
2591 .map(|(fragment_id, _, _, stream_node)| (*fragment_id, stream_node.clone()))
2592 .collect();
2593
2594 for (fragment_id, _, fragment_type_mask, stream_node) in fragments {
2595 Fragment::update(fragment::ActiveModel {
2596 fragment_id: Set(fragment_id),
2597 fragment_type_mask: Set(fragment_type_mask.into()),
2598 stream_node: Set(StreamNode::from(&stream_node)),
2599 ..Default::default()
2600 })
2601 .exec(&txn)
2602 .await?;
2603 }
2604
2605 txn.commit().await?;
2606
2607 let relation_info = PbObjectInfo::Source(ObjectModel(source, obj.unwrap(), None).into());
2608 let _version = self
2609 .notify_frontend(
2610 NotificationOperation::Update,
2611 NotificationInfo::ObjectGroup(PbObjectGroup {
2612 objects: vec![PbObject {
2613 object_info: Some(relation_info),
2614 }],
2615 dependencies: vec![],
2616 }),
2617 )
2618 .await;
2619
2620 Ok(fragment_nodes)
2621 }
2622
2623 pub async fn mutate_fragments_by_job_id(
2626 &self,
2627 job_id: JobId,
2628 mut fragments_mutation_fn: impl FnMut(FragmentTypeMask, &mut PbStreamNode) -> MetaResult<bool>,
2630 err_msg: &'static str,
2632 ) -> MetaResult<HashMap<FragmentId, PbStreamNode>> {
2633 let inner = self.inner.read().await;
2634 let txn = inner.db.begin().await?;
2635
2636 let fragments: Vec<(FragmentId, i32, StreamNode)> = Fragment::find()
2637 .select_only()
2638 .columns([
2639 fragment::Column::FragmentId,
2640 fragment::Column::FragmentTypeMask,
2641 fragment::Column::StreamNode,
2642 ])
2643 .filter(fragment::Column::JobId.eq(job_id))
2644 .into_tuple()
2645 .all(&txn)
2646 .await?;
2647 let mut fragments = fragments
2648 .into_iter()
2649 .map(|(id, mask, stream_node)| {
2650 (id, FragmentTypeMask::from(mask), stream_node.to_protobuf())
2651 })
2652 .collect_vec();
2653
2654 let fragments = fragments
2655 .iter_mut()
2656 .map(|(_, fragment_type_mask, stream_node)| {
2657 fragments_mutation_fn(*fragment_type_mask, stream_node)
2658 })
2659 .collect::<MetaResult<Vec<bool>>>()?
2660 .into_iter()
2661 .zip_eq_debug(std::mem::take(&mut fragments))
2662 .filter_map(|(keep, fragment)| if keep { Some(fragment) } else { None })
2663 .collect::<Vec<_>>();
2664
2665 if fragments.is_empty() {
2666 return Err(MetaError::invalid_parameter(format!(
2667 "job id {job_id}: {}",
2668 err_msg
2669 )));
2670 }
2671
2672 let fragment_nodes = fragments
2673 .iter()
2674 .map(|(id, _, stream_node)| (*id, stream_node.clone()))
2675 .collect();
2676 for (id, _, stream_node) in fragments {
2677 Fragment::update(fragment::ActiveModel {
2678 fragment_id: Set(id),
2679 stream_node: Set(StreamNode::from(&stream_node)),
2680 ..Default::default()
2681 })
2682 .exec(&txn)
2683 .await?;
2684 }
2685
2686 txn.commit().await?;
2687
2688 Ok(fragment_nodes)
2689 }
2690
2691 async fn mutate_fragment_by_fragment_id(
2692 &self,
2693 fragment_id: FragmentId,
2694 mut fragment_mutation_fn: impl FnMut(FragmentTypeMask, &mut PbStreamNode) -> MetaResult<bool>,
2695 err_msg: &'static str,
2696 ) -> MetaResult<PbStreamNode> {
2697 let inner = self.inner.read().await;
2698 let txn = inner.db.begin().await?;
2699
2700 let (fragment_type_mask, stream_node): (i32, StreamNode) =
2701 Fragment::find_by_id(fragment_id)
2702 .select_only()
2703 .columns([
2704 fragment::Column::FragmentTypeMask,
2705 fragment::Column::StreamNode,
2706 ])
2707 .into_tuple()
2708 .one(&txn)
2709 .await?
2710 .ok_or_else(|| MetaError::catalog_id_not_found("fragment", fragment_id))?;
2711 let mut pb_stream_node = stream_node.to_protobuf();
2712 let fragment_type_mask = FragmentTypeMask::from(fragment_type_mask);
2713
2714 if !fragment_mutation_fn(fragment_type_mask, &mut pb_stream_node)? {
2715 return Err(MetaError::invalid_parameter(format!(
2716 "fragment id {fragment_id}: {}",
2717 err_msg
2718 )));
2719 }
2720
2721 Fragment::update(fragment::ActiveModel {
2722 fragment_id: Set(fragment_id),
2723 stream_node: Set(StreamNode::from(&pb_stream_node)),
2724 ..Default::default()
2725 })
2726 .exec(&txn)
2727 .await?;
2728
2729 txn.commit().await?;
2730
2731 Ok(pb_stream_node)
2732 }
2733
2734 pub async fn update_backfill_orders_by_job_id(
2735 &self,
2736 job_id: JobId,
2737 backfill_orders: Option<BackfillOrders>,
2738 ) -> MetaResult<()> {
2739 let inner = self.inner.write().await;
2740 let txn = inner.db.begin().await?;
2741
2742 ensure_job_not_canceled(job_id, &txn).await?;
2743
2744 streaming_job::ActiveModel {
2745 job_id: Set(job_id),
2746 backfill_orders: Set(backfill_orders),
2747 ..Default::default()
2748 }
2749 .update(&txn)
2750 .await?;
2751
2752 txn.commit().await?;
2753
2754 Ok(())
2755 }
2756
2757 pub async fn update_backfill_rate_limit_by_job_id(
2760 &self,
2761 job_id: JobId,
2762 rate_limit: Option<u32>,
2763 ) -> MetaResult<HashMap<FragmentId, PbStreamNode>> {
2764 let update_backfill_rate_limit =
2765 |fragment_type_mask: FragmentTypeMask, stream_node: &mut PbStreamNode| {
2766 let mut found = false;
2767 if fragment_type_mask
2768 .contains_any(FragmentTypeFlag::backfill_rate_limit_fragments())
2769 {
2770 visit_stream_node_mut(stream_node, |node| match node {
2771 PbNodeBody::StreamCdcScan(node) => {
2772 node.rate_limit = rate_limit;
2773 found = true;
2774 }
2775 PbNodeBody::StreamScan(node) => {
2776 node.rate_limit = rate_limit;
2777 found = true;
2778 }
2779 PbNodeBody::SourceBackfill(node) => {
2780 node.rate_limit = rate_limit;
2781 found = true;
2782 }
2783 PbNodeBody::LocalityProvider(node) => {
2784 node.rate_limit = rate_limit;
2785 found = true;
2786 }
2787 _ => {}
2788 });
2789 }
2790 Ok(found)
2791 };
2792
2793 self.mutate_fragments_by_job_id(
2794 job_id,
2795 update_backfill_rate_limit,
2796 "stream scan node or source node not found",
2797 )
2798 .await
2799 }
2800
2801 pub async fn update_sink_rate_limit_by_job_id(
2804 &self,
2805 sink_id: SinkId,
2806 rate_limit: Option<u32>,
2807 ) -> MetaResult<HashMap<FragmentId, PbStreamNode>> {
2808 let update_sink_rate_limit =
2809 |fragment_type_mask: FragmentTypeMask, stream_node: &mut PbStreamNode| {
2810 let mut found = Ok(false);
2811 if fragment_type_mask.contains_any(FragmentTypeFlag::sink_rate_limit_fragments()) {
2812 visit_stream_node_mut(stream_node, |node| {
2813 if found.is_err() {
2814 return;
2815 }
2816 match update_sink_node_rate_limit(node, rate_limit) {
2817 Ok(true) => found = Ok(true),
2818 Ok(false) => {}
2819 Err(err) => found = Err(err),
2820 }
2821 });
2822 }
2823 found
2824 };
2825
2826 self.mutate_fragments_by_job_id(
2827 sink_id.as_job_id(),
2828 update_sink_rate_limit,
2829 "sink node not found",
2830 )
2831 .await
2832 }
2833
2834 pub async fn update_dml_rate_limit_by_job_id(
2835 &self,
2836 job_id: JobId,
2837 rate_limit: Option<u32>,
2838 ) -> MetaResult<HashMap<FragmentId, PbStreamNode>> {
2839 let update_dml_rate_limit =
2840 |fragment_type_mask: FragmentTypeMask, stream_node: &mut PbStreamNode| {
2841 let mut found = false;
2842 if fragment_type_mask.contains_any(FragmentTypeFlag::dml_rate_limit_fragments()) {
2843 visit_stream_node_mut(stream_node, |node| {
2844 if let PbNodeBody::Dml(node) = node {
2845 node.rate_limit = rate_limit;
2846 found = true;
2847 }
2848 });
2849 }
2850 Ok(found)
2851 };
2852
2853 self.mutate_fragments_by_job_id(job_id, update_dml_rate_limit, "dml node not found")
2854 .await
2855 }
2856
2857 pub async fn update_source_props_by_source_id(
2858 &self,
2859 source_id: SourceId,
2860 alter_props: BTreeMap<String, String>,
2861 alter_secret_refs: BTreeMap<String, PbSecretRef>,
2862 skip_alter_on_fly_check: bool,
2863 ) -> MetaResult<WithOptionsSecResolved> {
2864 let inner = self.inner.read().await;
2865 let txn = inner.db.begin().await?;
2866
2867 let (source, _obj) = Source::find_by_id(source_id)
2868 .find_also_related(Object)
2869 .one(&txn)
2870 .await?
2871 .ok_or_else(|| {
2872 MetaError::catalog_id_not_found(ObjectType::Source.as_str(), source_id)
2873 })?;
2874 ensure_source_props_not_set_by_connection(&txn, &source, &alter_props, &alter_secret_refs)
2875 .await?;
2876 let connector = source.with_properties.0.get_connector().unwrap();
2877 let is_shared_source = source.is_shared();
2878
2879 let mut dep_source_job_ids: Vec<JobId> = Vec::new();
2880 if !is_shared_source {
2881 dep_source_job_ids = ObjectDependency::find()
2883 .select_only()
2884 .column(object_dependency::Column::UsedBy)
2885 .filter(object_dependency::Column::Oid.eq(source_id))
2886 .into_tuple()
2887 .all(&txn)
2888 .await?;
2889 }
2890
2891 if let Some(new_connector) = alter_props.get(UPSTREAM_SOURCE_KEY)
2893 && new_connector != &connector
2894 {
2895 return Err(MetaError::invalid_parameter(format!(
2896 "Cannot change connector type from '{}' to '{}'. Drop and recreate the source instead.",
2897 connector, new_connector
2898 )));
2899 }
2900
2901 if !skip_alter_on_fly_check {
2903 let prop_keys: Vec<String> = alter_props
2904 .keys()
2905 .chain(alter_secret_refs.keys())
2906 .cloned()
2907 .collect();
2908 risingwave_connector::allow_alter_on_fly_fields::check_source_allow_alter_on_fly_fields(
2909 &connector, &prop_keys,
2910 )?;
2911 }
2912
2913 let mut options_with_secret = WithOptionsSecResolved::new(
2914 source.with_properties.0.clone(),
2915 source
2916 .secret_ref
2917 .map(|secret_ref| secret_ref.to_protobuf())
2918 .unwrap_or_default(),
2919 );
2920 let altered_options_with_secret =
2921 WithOptionsSecResolved::new(alter_props.clone(), alter_secret_refs.clone());
2922 let (to_add_secret_dep, to_remove_secret_dep) =
2923 options_with_secret.handle_update(alter_props, alter_secret_refs)?;
2924
2925 tracing::info!(
2926 "applying new properties to source: source_id={}, options_with_secret={:?}",
2927 source_id,
2928 options_with_secret
2929 );
2930 let _ = ConnectorProperties::extract(options_with_secret.clone(), true)?;
2932 let mut associate_table_id = None;
2935
2936 let mut preferred_id = source_id.as_object_id();
2940 let rewrite_sql = {
2941 let definition = source.definition.clone();
2942
2943 let [mut stmt]: [_; 1] = Parser::parse_sql(&definition)
2944 .map_err(|e| {
2945 MetaError::from(MetaErrorInner::Connector(ConnectorError::from(
2946 anyhow!(e).context("Failed to parse source definition SQL"),
2947 )))
2948 })?
2949 .try_into()
2950 .unwrap();
2951
2952 match &mut stmt {
2953 Statement::CreateSource { stmt } => {
2954 let altered_sql_options =
2955 format_with_option_secret_resolved(&txn, &altered_options_with_secret)
2956 .await?;
2957 merge_with_options(&mut stmt.with_properties.0, altered_sql_options);
2958 }
2959 Statement::CreateTable { with_options, .. } => {
2960 let altered_sql_options =
2961 format_with_option_secret_resolved(&txn, &altered_options_with_secret)
2962 .await?;
2963 merge_with_options(with_options, altered_sql_options);
2964 associate_table_id = source.optional_associated_table_id;
2965 preferred_id = associate_table_id.unwrap().as_object_id();
2966 }
2967 _ => unreachable!(),
2968 }
2969
2970 stmt.to_string()
2971 };
2972
2973 {
2974 if !to_add_secret_dep.is_empty() {
2977 ObjectDependency::insert_many(to_add_secret_dep.into_iter().map(|secret_id| {
2978 object_dependency::ActiveModel {
2979 oid: Set(secret_id.into()),
2980 used_by: Set(preferred_id),
2981 ..Default::default()
2982 }
2983 }))
2984 .exec(&txn)
2985 .await?;
2986 }
2987 if !to_remove_secret_dep.is_empty() {
2990 let _ = ObjectDependency::delete_many()
2991 .filter(
2992 object_dependency::Column::Oid
2993 .is_in(to_remove_secret_dep)
2994 .and(object_dependency::Column::UsedBy.eq(preferred_id)),
2995 )
2996 .exec(&txn)
2997 .await?;
2998 }
2999 }
3000
3001 let active_source_model = source::ActiveModel {
3002 source_id: Set(source_id),
3003 definition: Set(rewrite_sql.clone()),
3004 with_properties: Set(options_with_secret.as_plaintext().clone().into()),
3005 secret_ref: Set((!options_with_secret.as_secret().is_empty())
3006 .then(|| SecretRef::from(options_with_secret.as_secret().clone()))),
3007 ..Default::default()
3008 };
3009 Source::update(active_source_model).exec(&txn).await?;
3010
3011 if let Some(associate_table_id) = associate_table_id {
3012 let active_table_model = table::ActiveModel {
3014 table_id: Set(associate_table_id),
3015 definition: Set(rewrite_sql),
3016 ..Default::default()
3017 };
3018 Table::update(active_table_model).exec(&txn).await?;
3019 }
3020
3021 let to_check_job_ids = vec![if let Some(associate_table_id) = associate_table_id {
3022 associate_table_id.as_job_id()
3024 } else {
3025 source_id.as_share_source_job_id()
3026 }]
3027 .into_iter()
3028 .chain(dep_source_job_ids)
3029 .collect_vec();
3030
3031 update_connector_props_fragments(
3033 &txn,
3034 to_check_job_ids,
3035 FragmentTypeFlag::Source,
3036 |node, found| {
3037 if let PbNodeBody::Source(node) = node
3038 && let Some(source_inner) = &mut node.source_inner
3039 {
3040 source_inner.with_properties = options_with_secret.as_plaintext().clone();
3041 source_inner.secret_refs = options_with_secret.as_secret().clone();
3042 *found = true;
3043 }
3044 },
3045 is_shared_source,
3046 )
3047 .await?;
3048
3049 let mut to_update_objs = Vec::with_capacity(2);
3050 let (source, obj) = Source::find_by_id(source_id)
3051 .find_also_related(Object)
3052 .one(&txn)
3053 .await?
3054 .ok_or_else(|| {
3055 MetaError::catalog_id_not_found(ObjectType::Source.as_str(), source_id)
3056 })?;
3057 to_update_objs.push(PbObject {
3058 object_info: Some(PbObjectInfo::Source(
3059 ObjectModel(source, obj.unwrap(), None).into(),
3060 )),
3061 });
3062
3063 if let Some(associate_table_id) = associate_table_id {
3064 let (table, obj) = Table::find_by_id(associate_table_id)
3065 .find_also_related(Object)
3066 .one(&txn)
3067 .await?
3068 .ok_or_else(|| MetaError::catalog_id_not_found("table", associate_table_id))?;
3069 let streaming_job = streaming_job::Entity::find_by_id(table.job_id())
3070 .one(&txn)
3071 .await?;
3072 to_update_objs.push(PbObject {
3073 object_info: Some(PbObjectInfo::Table(
3074 ObjectModel(table, obj.unwrap(), streaming_job).into(),
3075 )),
3076 });
3077 }
3078
3079 txn.commit().await?;
3080
3081 self.notify_frontend(
3082 NotificationOperation::Update,
3083 NotificationInfo::ObjectGroup(PbObjectGroup {
3084 objects: to_update_objs,
3085 dependencies: vec![],
3086 }),
3087 )
3088 .await;
3089
3090 Ok(options_with_secret)
3091 }
3092
3093 pub async fn update_sink_props_by_sink_id(
3094 &self,
3095 sink_id: SinkId,
3096 props: BTreeMap<String, String>,
3097 ) -> MetaResult<HashMap<String, String>> {
3098 let inner = self.inner.read().await;
3099 let txn = inner.db.begin().await?;
3100
3101 let (sink, _obj) = Sink::find_by_id(sink_id)
3102 .find_also_related(Object)
3103 .one(&txn)
3104 .await?
3105 .ok_or_else(|| MetaError::catalog_id_not_found(ObjectType::Sink.as_str(), sink_id))?;
3106 validate_sink_props(&sink, &props, &props.keys().cloned().collect_vec())?;
3107 let definition = sink.definition.clone();
3108 let [mut stmt]: [_; 1] = Parser::parse_sql(&definition)
3109 .map_err(|e| SinkError::Config(anyhow!(e)))?
3110 .try_into()
3111 .unwrap();
3112 if let Statement::CreateSink { stmt } = &mut stmt {
3113 update_stmt_with_props(&mut stmt.with_properties.0, &props)?;
3114 } else {
3115 panic!("definition is not a create sink statement")
3116 }
3117 let mut new_config = sink.properties.clone().into_inner();
3118 new_config.extend(props.clone());
3119
3120 let definition = stmt.to_string();
3121 let active_sink = sink::ActiveModel {
3122 sink_id: Set(sink_id),
3123 properties: Set(risingwave_meta_model::Property(new_config.clone())),
3124 definition: Set(definition),
3125 ..Default::default()
3126 };
3127 Sink::update(active_sink).exec(&txn).await?;
3128
3129 update_sink_fragment_props(&txn, sink_id, new_config).await?;
3130 let (sink, obj) = Sink::find_by_id(sink_id)
3131 .find_also_related(Object)
3132 .one(&txn)
3133 .await?
3134 .ok_or_else(|| MetaError::catalog_id_not_found(ObjectType::Sink.as_str(), sink_id))?;
3135 let streaming_job = streaming_job::Entity::find_by_id(sink.sink_id.as_job_id())
3136 .one(&txn)
3137 .await?;
3138 txn.commit().await?;
3139 let relation_infos = vec![PbObject {
3140 object_info: Some(PbObjectInfo::Sink(
3141 ObjectModel(sink, obj.unwrap(), streaming_job).into(),
3142 )),
3143 }];
3144
3145 let _version = self
3146 .notify_frontend(
3147 NotificationOperation::Update,
3148 NotificationInfo::ObjectGroup(PbObjectGroup {
3149 objects: relation_infos,
3150 dependencies: vec![],
3151 }),
3152 )
3153 .await;
3154
3155 Ok(props.into_iter().collect())
3156 }
3157
3158 pub async fn update_iceberg_table_props_by_table_id(
3159 &self,
3160 table_id: TableId,
3161 props: BTreeMap<String, String>,
3162 alter_iceberg_table_props: Option<
3163 risingwave_pb::meta::alter_connector_props_request::PbExtraOptions,
3164 >,
3165 ) -> MetaResult<(HashMap<String, String>, SinkId)> {
3166 let risingwave_pb::meta::alter_connector_props_request::PbExtraOptions::AlterIcebergTableIds(AlterIcebergTableIds { sink_id, source_id }) = alter_iceberg_table_props.
3167 ok_or_else(|| MetaError::invalid_parameter("alter_iceberg_table_props is required"))?;
3168 let inner = self.inner.read().await;
3169 let txn = inner.db.begin().await?;
3170
3171 let (sink, _obj) = Sink::find_by_id(sink_id)
3172 .find_also_related(Object)
3173 .one(&txn)
3174 .await?
3175 .ok_or_else(|| MetaError::catalog_id_not_found(ObjectType::Sink.as_str(), sink_id))?;
3176 validate_sink_props(&sink, &props, &props.keys().cloned().collect_vec())?;
3177
3178 let definition = sink.definition.clone();
3179 let [mut stmt]: [_; 1] = Parser::parse_sql(&definition)
3180 .map_err(|e| SinkError::Config(anyhow!(e)))?
3181 .try_into()
3182 .unwrap();
3183 if let Statement::CreateTable {
3184 with_options,
3185 engine,
3186 ..
3187 } = &mut stmt
3188 {
3189 if !matches!(engine, Engine::Iceberg) {
3190 return Err(SinkError::Config(anyhow!(
3191 "only iceberg table can be altered as sink"
3192 ))
3193 .into());
3194 }
3195 update_stmt_with_props(with_options, &props)?;
3196 } else {
3197 panic!("definition is not a create iceberg table statement")
3198 }
3199 let mut new_config = sink.properties.clone().into_inner();
3200 new_config.extend(props.clone());
3201
3202 let definition = stmt.to_string();
3203 let active_sink = sink::ActiveModel {
3204 sink_id: Set(sink_id),
3205 properties: Set(risingwave_meta_model::Property(new_config.clone())),
3206 definition: Set(definition.clone()),
3207 ..Default::default()
3208 };
3209 let active_source = source::ActiveModel {
3210 source_id: Set(source_id),
3211 definition: Set(definition.clone()),
3212 ..Default::default()
3213 };
3214 let active_table = table::ActiveModel {
3215 table_id: Set(table_id),
3216 definition: Set(definition),
3217 ..Default::default()
3218 };
3219 Sink::update(active_sink).exec(&txn).await?;
3220 Source::update(active_source).exec(&txn).await?;
3221 Table::update(active_table).exec(&txn).await?;
3222
3223 update_sink_fragment_props(&txn, sink_id, new_config).await?;
3224
3225 let (sink, sink_obj) = Sink::find_by_id(sink_id)
3226 .find_also_related(Object)
3227 .one(&txn)
3228 .await?
3229 .ok_or_else(|| MetaError::catalog_id_not_found(ObjectType::Sink.as_str(), sink_id))?;
3230 let sink_streaming_job = streaming_job::Entity::find_by_id(sink.sink_id.as_job_id())
3231 .one(&txn)
3232 .await?;
3233 let (source, source_obj) = Source::find_by_id(source_id)
3234 .find_also_related(Object)
3235 .one(&txn)
3236 .await?
3237 .ok_or_else(|| {
3238 MetaError::catalog_id_not_found(ObjectType::Source.as_str(), source_id)
3239 })?;
3240 let (table, table_obj) = Table::find_by_id(table_id)
3241 .find_also_related(Object)
3242 .one(&txn)
3243 .await?
3244 .ok_or_else(|| MetaError::catalog_id_not_found(ObjectType::Table.as_str(), table_id))?;
3245 let table_streaming_job = streaming_job::Entity::find_by_id(table.job_id())
3246 .one(&txn)
3247 .await?;
3248 txn.commit().await?;
3249 let relation_infos = vec![
3250 PbObject {
3251 object_info: Some(PbObjectInfo::Sink(
3252 ObjectModel(sink, sink_obj.unwrap(), sink_streaming_job).into(),
3253 )),
3254 },
3255 PbObject {
3256 object_info: Some(PbObjectInfo::Source(
3257 ObjectModel(source, source_obj.unwrap(), None).into(),
3258 )),
3259 },
3260 PbObject {
3261 object_info: Some(PbObjectInfo::Table(
3262 ObjectModel(table, table_obj.unwrap(), table_streaming_job).into(),
3263 )),
3264 },
3265 ];
3266 let _version = self
3267 .notify_frontend(
3268 NotificationOperation::Update,
3269 NotificationInfo::ObjectGroup(PbObjectGroup {
3270 objects: relation_infos,
3271 dependencies: vec![],
3272 }),
3273 )
3274 .await;
3275
3276 Ok((props.into_iter().collect(), sink_id))
3277 }
3278
3279 pub async fn update_connection_and_dependent_objects_props(
3281 &self,
3282 connection_id: ConnectionId,
3283 alter_props: BTreeMap<String, String>,
3284 alter_secret_refs: BTreeMap<String, PbSecretRef>,
3285 ) -> MetaResult<(
3286 WithOptionsSecResolved, Vec<(SourceId, HashMap<String, String>)>, Vec<(SinkId, HashMap<String, String>)>, )> {
3290 let inner = self.inner.read().await;
3291 let txn = inner.db.begin().await?;
3292
3293 let dependent_sources: Vec<SourceId> = Source::find()
3295 .select_only()
3296 .column(source::Column::SourceId)
3297 .filter(source::Column::ConnectionId.eq(connection_id))
3298 .into_tuple()
3299 .all(&txn)
3300 .await?;
3301
3302 let dependent_sinks: Vec<SinkId> = Sink::find()
3303 .select_only()
3304 .column(sink::Column::SinkId)
3305 .filter(sink::Column::ConnectionId.eq(connection_id))
3306 .into_tuple()
3307 .all(&txn)
3308 .await?;
3309
3310 let (connection_catalog, _obj) = Connection::find_by_id(connection_id)
3311 .find_also_related(Object)
3312 .one(&txn)
3313 .await?
3314 .ok_or_else(|| {
3315 MetaError::catalog_id_not_found(ObjectType::Connection.as_str(), connection_id)
3316 })?;
3317
3318 let prop_keys: Vec<String> = alter_props
3320 .keys()
3321 .chain(alter_secret_refs.keys())
3322 .cloned()
3323 .collect();
3324
3325 let connection_type_str = pb_connection_type_to_connection_type(
3327 &connection_catalog.params.to_protobuf().connection_type(),
3328 )
3329 .ok_or_else(|| MetaError::invalid_parameter("Unspecified connection type"))?;
3330
3331 risingwave_connector::allow_alter_on_fly_fields::check_connection_allow_alter_on_fly_fields(
3332 connection_type_str, &prop_keys,
3333 )?;
3334
3335 let connection_pb = connection_catalog.params.to_protobuf();
3336 let mut connection_options_with_secret = WithOptionsSecResolved::new(
3337 connection_pb.properties.into_iter().collect(),
3338 connection_pb.secret_refs.into_iter().collect(),
3339 );
3340
3341 let (to_add_secret_dep, to_remove_secret_dep) = connection_options_with_secret
3342 .handle_update(alter_props.clone(), alter_secret_refs.clone())?;
3343
3344 tracing::debug!(
3345 "applying new properties to connection and dependents: connection_id={}, sources={:?}, sinks={:?}",
3346 connection_id,
3347 dependent_sources,
3348 dependent_sinks
3349 );
3350
3351 {
3353 let conn_params_pb = risingwave_pb::catalog::ConnectionParams {
3354 connection_type: connection_pb.connection_type,
3355 properties: connection_options_with_secret
3356 .as_plaintext()
3357 .clone()
3358 .into_iter()
3359 .collect(),
3360 secret_refs: connection_options_with_secret
3361 .as_secret()
3362 .clone()
3363 .into_iter()
3364 .collect(),
3365 };
3366 let connection = PbConnection {
3367 id: connection_id as _,
3368 info: Some(risingwave_pb::catalog::connection::Info::ConnectionParams(
3369 conn_params_pb,
3370 )),
3371 ..Default::default()
3372 };
3373 validate_connection(&connection).await?;
3374 }
3375
3376 if !to_add_secret_dep.is_empty() {
3378 ObjectDependency::insert_many(to_add_secret_dep.into_iter().map(|secret_id| {
3379 object_dependency::ActiveModel {
3380 oid: Set(secret_id.into()),
3381 used_by: Set(connection_id.as_object_id()),
3382 ..Default::default()
3383 }
3384 }))
3385 .exec(&txn)
3386 .await?;
3387 }
3388 if !to_remove_secret_dep.is_empty() {
3389 let _ = ObjectDependency::delete_many()
3390 .filter(
3391 object_dependency::Column::Oid
3392 .is_in(to_remove_secret_dep)
3393 .and(object_dependency::Column::UsedBy.eq(connection_id.as_object_id())),
3394 )
3395 .exec(&txn)
3396 .await?;
3397 }
3398
3399 let updated_connection_params = risingwave_pb::catalog::ConnectionParams {
3401 connection_type: connection_pb.connection_type,
3402 properties: connection_options_with_secret
3403 .as_plaintext()
3404 .clone()
3405 .into_iter()
3406 .collect(),
3407 secret_refs: connection_options_with_secret
3408 .as_secret()
3409 .clone()
3410 .into_iter()
3411 .collect(),
3412 };
3413 let active_connection_model = connection::ActiveModel {
3414 connection_id: Set(connection_id),
3415 params: Set(ConnectionParams::from(&updated_connection_params)),
3416 ..Default::default()
3417 };
3418 Connection::update(active_connection_model)
3419 .exec(&txn)
3420 .await?;
3421
3422 let mut updated_sources_with_props: Vec<(SourceId, HashMap<String, String>)> = Vec::new();
3424
3425 if !dependent_sources.is_empty() {
3426 let sources_with_objs = Source::find()
3428 .find_also_related(Object)
3429 .filter(source::Column::SourceId.is_in(dependent_sources.iter().cloned()))
3430 .all(&txn)
3431 .await?;
3432
3433 let mut source_updates = Vec::new();
3435 let mut fragment_updates: Vec<DependentSourceFragmentUpdate> = Vec::new();
3436
3437 for (source, _obj) in sources_with_objs {
3438 let source_id = source.source_id;
3439
3440 let mut source_options_with_secret = WithOptionsSecResolved::new(
3441 source.with_properties.0.clone(),
3442 source
3443 .secret_ref
3444 .clone()
3445 .map(|secret_ref| secret_ref.to_protobuf())
3446 .unwrap_or_default(),
3447 );
3448 let (source_to_add_secret_dep, source_to_remove_secret_dep) =
3449 source_options_with_secret
3450 .handle_update(alter_props.clone(), alter_secret_refs.clone())?;
3451
3452 let _ = ConnectorProperties::extract(source_options_with_secret.clone(), true)?;
3454
3455 let source_used_by_id = source
3458 .optional_associated_table_id
3459 .map(|table_id| table_id.as_object_id())
3460 .unwrap_or_else(|| source_id.as_object_id());
3461 if !source_to_add_secret_dep.is_empty() {
3462 ObjectDependency::insert_many(source_to_add_secret_dep.into_iter().map(
3463 |secret_id| object_dependency::ActiveModel {
3464 oid: Set(secret_id.into()),
3465 used_by: Set(source_used_by_id),
3466 ..Default::default()
3467 },
3468 ))
3469 .exec(&txn)
3470 .await?;
3471 }
3472 if !source_to_remove_secret_dep.is_empty() {
3473 let _ = ObjectDependency::delete_many()
3474 .filter(
3475 object_dependency::Column::Oid
3476 .is_in(source_to_remove_secret_dep)
3477 .and(object_dependency::Column::UsedBy.eq(source_used_by_id)),
3478 )
3479 .exec(&txn)
3480 .await?;
3481 }
3482
3483 let active_source = source::ActiveModel {
3485 source_id: Set(source_id),
3486 with_properties: Set(Property(
3487 source_options_with_secret.as_plaintext().clone(),
3488 )),
3489 secret_ref: Set((!source_options_with_secret.as_secret().is_empty()).then(
3490 || {
3491 risingwave_meta_model::SecretRef::from(
3492 source_options_with_secret.as_secret().clone(),
3493 )
3494 },
3495 )),
3496 ..Default::default()
3497 };
3498 source_updates.push(active_source);
3499
3500 let is_shared_source = source.is_shared();
3505 let mut dep_source_job_ids: Vec<JobId> = Vec::new();
3506 if !is_shared_source {
3507 dep_source_job_ids = ObjectDependency::find()
3508 .select_only()
3509 .column(object_dependency::Column::UsedBy)
3510 .filter(object_dependency::Column::Oid.eq(source_id))
3511 .into_tuple()
3512 .all(&txn)
3513 .await?;
3514 }
3515
3516 let base_job_id =
3517 if let Some(associate_table_id) = source.optional_associated_table_id {
3518 associate_table_id.as_job_id()
3519 } else {
3520 source_id.as_share_source_job_id()
3521 };
3522 let job_ids = vec![base_job_id]
3523 .into_iter()
3524 .chain(dep_source_job_ids)
3525 .collect_vec();
3526
3527 fragment_updates.push(DependentSourceFragmentUpdate {
3528 job_ids,
3529 with_properties: source_options_with_secret.as_plaintext().clone(),
3530 secret_refs: source_options_with_secret.as_secret().clone(),
3531 is_shared_source,
3532 });
3533
3534 let complete_source_props = LocalSecretManager::global()
3536 .fill_secrets(
3537 source_options_with_secret.as_plaintext().clone(),
3538 source_options_with_secret.as_secret().clone(),
3539 )
3540 .map_err(MetaError::from)?
3541 .into_iter()
3542 .collect::<HashMap<String, String>>();
3543 updated_sources_with_props.push((source_id, complete_source_props));
3544 }
3545
3546 for source_update in source_updates {
3547 Source::update(source_update).exec(&txn).await?;
3548 }
3549
3550 for DependentSourceFragmentUpdate {
3552 job_ids,
3553 with_properties,
3554 secret_refs,
3555 is_shared_source,
3556 } in fragment_updates
3557 {
3558 update_connector_props_fragments(
3559 &txn,
3560 job_ids,
3561 FragmentTypeFlag::Source,
3562 |node, found| {
3563 if let PbNodeBody::Source(node) = node
3564 && let Some(source_inner) = &mut node.source_inner
3565 {
3566 source_inner.with_properties = with_properties.clone();
3567 source_inner.secret_refs = secret_refs.clone();
3568 *found = true;
3569 }
3570 },
3571 is_shared_source,
3572 )
3573 .await?;
3574 }
3575 }
3576
3577 let mut updated_sinks_with_props: Vec<(SinkId, HashMap<String, String>)> = Vec::new();
3579
3580 if !dependent_sinks.is_empty() {
3581 let sinks_with_objs = Sink::find()
3583 .find_also_related(Object)
3584 .filter(sink::Column::SinkId.is_in(dependent_sinks.iter().cloned()))
3585 .all(&txn)
3586 .await?;
3587
3588 let mut sink_updates = Vec::new();
3590 let mut sink_fragment_updates = Vec::new();
3591
3592 for (sink, _obj) in sinks_with_objs {
3593 let sink_id = sink.sink_id;
3594
3595 validate_sink_props(&sink, &alter_props, &prop_keys)?;
3596
3597 let mut new_sink_props = sink.properties.0.clone();
3598 new_sink_props.extend(alter_props.clone());
3599
3600 let active_sink = sink::ActiveModel {
3602 sink_id: Set(sink_id),
3603 properties: Set(risingwave_meta_model::Property(new_sink_props.clone())),
3604 ..Default::default()
3605 };
3606 sink_updates.push(active_sink);
3607
3608 sink_fragment_updates.push((sink_id, new_sink_props.clone()));
3610
3611 let complete_sink_props: HashMap<String, String> =
3613 new_sink_props.into_iter().collect();
3614 updated_sinks_with_props.push((sink_id, complete_sink_props));
3615 }
3616
3617 for sink_update in sink_updates {
3619 Sink::update(sink_update).exec(&txn).await?;
3620 }
3621
3622 for (sink_id, new_sink_props) in sink_fragment_updates {
3624 update_connector_props_fragments(
3625 &txn,
3626 vec![sink_id.as_job_id()],
3627 FragmentTypeFlag::Sink,
3628 |node, found| {
3629 if let PbNodeBody::Sink(node) = node
3630 && let Some(sink_desc) = &mut node.sink_desc
3631 && sink_desc.id == sink_id.as_raw_id()
3632 {
3633 sink_desc.properties = new_sink_props.clone();
3634 *found = true;
3635 }
3636 },
3637 true,
3638 )
3639 .await?;
3640 }
3641 }
3642
3643 let mut updated_objects = Vec::new();
3645
3646 let (connection, obj) = Connection::find_by_id(connection_id)
3648 .find_also_related(Object)
3649 .one(&txn)
3650 .await?
3651 .ok_or_else(|| {
3652 MetaError::catalog_id_not_found(ObjectType::Connection.as_str(), connection_id)
3653 })?;
3654 updated_objects.push(PbObject {
3655 object_info: Some(PbObjectInfo::Connection(
3656 ObjectModel(connection, obj.unwrap(), None).into(),
3657 )),
3658 });
3659
3660 for source_id in &dependent_sources {
3662 let (source, obj) = Source::find_by_id(*source_id)
3663 .find_also_related(Object)
3664 .one(&txn)
3665 .await?
3666 .ok_or_else(|| {
3667 MetaError::catalog_id_not_found(ObjectType::Source.as_str(), *source_id)
3668 })?;
3669 updated_objects.push(PbObject {
3670 object_info: Some(PbObjectInfo::Source(
3671 ObjectModel(source, obj.unwrap(), None).into(),
3672 )),
3673 });
3674 }
3675
3676 for sink_id in &dependent_sinks {
3678 let (sink, obj) = Sink::find_by_id(*sink_id)
3679 .find_also_related(Object)
3680 .one(&txn)
3681 .await?
3682 .ok_or_else(|| {
3683 MetaError::catalog_id_not_found(ObjectType::Sink.as_str(), *sink_id)
3684 })?;
3685 let streaming_job = streaming_job::Entity::find_by_id(sink.sink_id.as_job_id())
3686 .one(&txn)
3687 .await?;
3688 updated_objects.push(PbObject {
3689 object_info: Some(PbObjectInfo::Sink(
3690 ObjectModel(sink, obj.unwrap(), streaming_job).into(),
3691 )),
3692 });
3693 }
3694
3695 txn.commit().await?;
3697
3698 if !updated_objects.is_empty() {
3700 self.notify_frontend(
3701 NotificationOperation::Update,
3702 NotificationInfo::ObjectGroup(PbObjectGroup {
3703 objects: updated_objects,
3704 dependencies: vec![],
3705 }),
3706 )
3707 .await;
3708 }
3709
3710 Ok((
3711 connection_options_with_secret,
3712 updated_sources_with_props,
3713 updated_sinks_with_props,
3714 ))
3715 }
3716
3717 pub async fn update_fragment_rate_limit_by_fragment_id(
3718 &self,
3719 fragment_id: FragmentId,
3720 throttle_type: ThrottleType,
3721 rate_limit: Option<u32>,
3722 ) -> MetaResult<PbStreamNode> {
3723 let update_rate_limit = |fragment_type_mask: FragmentTypeMask,
3724 stream_node: &mut PbStreamNode| {
3725 let mut found = Ok(false);
3726 match throttle_type {
3727 ThrottleType::Source => {
3728 visit_stream_node_mut(stream_node, |node| match node {
3729 PbNodeBody::Source(node) => {
3730 if let Some(node_inner) = &mut node.source_inner {
3731 node_inner.rate_limit = rate_limit;
3732 found = Ok(true);
3733 }
3734 }
3735 PbNodeBody::StreamFsFetch(node) => {
3736 if let Some(node_inner) = &mut node.node_inner {
3737 node_inner.rate_limit = rate_limit;
3738 found = Ok(true);
3739 }
3740 }
3741 _ => {}
3742 });
3743 }
3744 ThrottleType::Backfill => {
3745 if fragment_type_mask
3746 .contains_any(FragmentTypeFlag::backfill_rate_limit_fragments())
3747 {
3748 visit_stream_node_mut(stream_node, |node| match node {
3749 PbNodeBody::StreamCdcScan(node) => {
3750 node.rate_limit = rate_limit;
3751 found = Ok(true);
3752 }
3753 PbNodeBody::StreamScan(node) => {
3754 node.rate_limit = rate_limit;
3755 found = Ok(true);
3756 }
3757 PbNodeBody::SourceBackfill(node) => {
3758 node.rate_limit = rate_limit;
3759 found = Ok(true);
3760 }
3761 PbNodeBody::LocalityProvider(node) => {
3762 node.rate_limit = rate_limit;
3763 found = Ok(true);
3764 }
3765 _ => {}
3766 });
3767 }
3768 }
3769 ThrottleType::Sink => {
3770 if fragment_type_mask
3771 .contains_any(FragmentTypeFlag::sink_rate_limit_fragments())
3772 {
3773 visit_stream_node_mut(stream_node, |node| {
3774 if found.is_err() {
3775 return;
3776 }
3777 match update_sink_node_rate_limit(node, rate_limit) {
3778 Ok(true) => found = Ok(true),
3779 Ok(false) => {}
3780 Err(err) => found = Err(err),
3781 }
3782 });
3783 }
3784 }
3785 ThrottleType::Dml => {
3786 if fragment_type_mask.contains_any(FragmentTypeFlag::dml_rate_limit_fragments())
3787 {
3788 visit_stream_node_mut(stream_node, |node| {
3789 if let PbNodeBody::Dml(node) = node {
3790 node.rate_limit = rate_limit;
3791 found = Ok(true);
3792 }
3793 });
3794 }
3795 }
3796 ThrottleType::Unspecified => {}
3797 }
3798 found
3799 };
3800 self.mutate_fragment_by_fragment_id(
3801 fragment_id,
3802 update_rate_limit,
3803 "rate limit node not found",
3804 )
3805 .await
3806 }
3807
3808 pub async fn list_rate_limits(&self) -> MetaResult<Vec<RateLimitInfo>> {
3811 let inner = self.inner.read().await;
3812 let txn = inner.db.begin().await?;
3813
3814 let fragments: Vec<(FragmentId, JobId, i32, StreamNode)> = Fragment::find()
3815 .select_only()
3816 .columns([
3817 fragment::Column::FragmentId,
3818 fragment::Column::JobId,
3819 fragment::Column::FragmentTypeMask,
3820 fragment::Column::StreamNode,
3821 ])
3822 .filter(FragmentTypeMask::intersects_any(
3823 FragmentTypeFlag::rate_limit_fragments(),
3824 ))
3825 .into_tuple()
3826 .all(&txn)
3827 .await?;
3828
3829 let mut rate_limits = Vec::new();
3830 for (fragment_id, job_id, fragment_type_mask, stream_node) in fragments {
3831 let stream_node = stream_node.to_protobuf();
3832 visit_stream_node_body(&stream_node, |node| {
3833 let mut rate_limit = None;
3834 let mut node_name = None;
3835
3836 match node {
3837 PbNodeBody::Source(node) => {
3839 if let Some(node_inner) = &node.source_inner {
3840 rate_limit = node_inner.rate_limit;
3841 node_name = Some("SOURCE");
3842 }
3843 }
3844 PbNodeBody::StreamFsFetch(node) => {
3845 if let Some(node_inner) = &node.node_inner {
3846 rate_limit = node_inner.rate_limit;
3847 node_name = Some("FS_FETCH");
3848 }
3849 }
3850 PbNodeBody::SourceBackfill(node) => {
3852 rate_limit = node.rate_limit;
3853 node_name = Some("SOURCE_BACKFILL");
3854 }
3855 PbNodeBody::StreamScan(node) => {
3856 rate_limit = node.rate_limit;
3857 node_name = Some("STREAM_SCAN");
3858 }
3859 PbNodeBody::StreamCdcScan(node) => {
3860 rate_limit = node.rate_limit;
3861 node_name = Some("STREAM_CDC_SCAN");
3862 }
3863 PbNodeBody::LocalityProvider(node) => {
3864 rate_limit = node.rate_limit;
3865 node_name = Some("LOCALITY_PROVIDER");
3866 }
3867 PbNodeBody::Sink(node) => {
3868 rate_limit = node.rate_limit;
3869 node_name = Some("SINK");
3870 }
3871 PbNodeBody::Dml(node) => {
3872 rate_limit = node.rate_limit;
3873 node_name = Some("DML");
3874 }
3875 _ => {}
3876 }
3877
3878 if let Some(rate_limit) = rate_limit {
3879 rate_limits.push(RateLimitInfo {
3880 fragment_id,
3881 job_id,
3882 fragment_type_mask: fragment_type_mask as u32,
3883 rate_limit,
3884 node_name: node_name.unwrap().to_owned(),
3885 });
3886 }
3887 });
3888 }
3889
3890 Ok(rate_limits)
3891 }
3892}
3893
3894fn validate_sink_props(
3897 sink: &sink::Model,
3898 alter_props: &BTreeMap<String, String>,
3899 altered_field_names: &[String],
3900) -> MetaResult<()> {
3901 let Some(connector) = sink.properties.inner_ref().get(CONNECTOR_TYPE_KEY) else {
3902 return Err(SinkError::Config(anyhow!("connector not specified when alter sink")).into());
3903 };
3904 let connector_type = connector.to_lowercase();
3905 check_sink_allow_alter_on_fly_fields(&connector_type, altered_field_names)
3906 .map_err(|e| SinkError::Config(anyhow!(e)))?;
3907
3908 let mut new_props = LocalSecretManager::global()
3911 .fill_secrets(
3912 sink.properties.0.clone(),
3913 sink.secret_ref
3914 .as_ref()
3915 .map(|secret_ref| secret_ref.to_protobuf())
3916 .unwrap_or_default(),
3917 )
3918 .map_err(MetaError::from)?;
3919 new_props.extend(alter_props.clone());
3920
3921 match_sink_name_str!(
3922 connector_type.as_str(),
3923 SinkType,
3924 SinkType::validate_alter_config_change(&new_props, alter_props),
3925 |sink: &str| Err(SinkError::Config(anyhow!("unsupported sink type {}", sink)))
3926 )?;
3927 Ok(())
3928}
3929
3930fn update_stmt_with_props(
3931 with_properties: &mut Vec<SqlOption>,
3932 props: &BTreeMap<String, String>,
3933) -> MetaResult<()> {
3934 let mut new_sql_options = with_properties
3935 .iter()
3936 .map(|sql_option| (&sql_option.name, sql_option))
3937 .collect::<IndexMap<_, _>>();
3938 let add_sql_options = props
3939 .iter()
3940 .map(|(k, v)| SqlOption::try_from((k, v)))
3941 .collect::<Result<Vec<SqlOption>, ParserError>>()
3942 .map_err(|e| SinkError::Config(anyhow!(e)))?;
3943 new_sql_options.extend(
3944 add_sql_options
3945 .iter()
3946 .map(|sql_option| (&sql_option.name, sql_option)),
3947 );
3948 *with_properties = new_sql_options.into_values().cloned().collect();
3949 Ok(())
3950}
3951
3952fn merge_with_options(with_properties: &mut Vec<SqlOption>, altered_options: Vec<SqlOption>) {
3953 for altered_option in altered_options {
3954 if let Some(existing_option) = with_properties
3955 .iter_mut()
3956 .find(|option| option.name.real_value() == altered_option.name.real_value())
3957 {
3958 *existing_option = altered_option;
3959 } else {
3960 with_properties.push(altered_option);
3961 }
3962 }
3963}
3964
3965async fn ensure_source_props_not_set_by_connection(
3966 txn: &DatabaseTransaction,
3967 source: &source::Model,
3968 alter_props: &BTreeMap<String, String>,
3969 alter_secret_refs: &BTreeMap<String, PbSecretRef>,
3970) -> MetaResult<()> {
3971 let Some(connection_id) = source.connection_id else {
3972 return Ok(());
3973 };
3974
3975 let connection = Connection::find_by_id(connection_id)
3976 .one(txn)
3977 .await?
3978 .ok_or_else(|| {
3979 MetaError::catalog_id_not_found(ObjectType::Connection.as_str(), connection_id)
3980 })?;
3981 let connection_params = connection.params.to_protobuf();
3982
3983 if let Some(key) = alter_props
3984 .keys()
3985 .chain(alter_secret_refs.keys())
3986 .find(|key| {
3987 connection_params.properties.contains_key(*key)
3988 || connection_params.secret_refs.contains_key(*key)
3989 })
3990 {
3991 return Err(MetaError::invalid_parameter(format!(
3992 "Cannot alter source connector property `{key}` because it is set by CONNECTION `{}`. Use ALTER CONNECTION instead.",
3993 connection.name
3994 )));
3995 }
3996
3997 Ok(())
3998}
3999
4000async fn update_sink_fragment_props(
4001 txn: &DatabaseTransaction,
4002 sink_id: SinkId,
4003 props: BTreeMap<String, String>,
4004) -> MetaResult<()> {
4005 let fragments: Vec<(FragmentId, i32, StreamNode)> = Fragment::find()
4006 .select_only()
4007 .columns([
4008 fragment::Column::FragmentId,
4009 fragment::Column::FragmentTypeMask,
4010 fragment::Column::StreamNode,
4011 ])
4012 .filter(fragment::Column::JobId.eq(sink_id))
4013 .into_tuple()
4014 .all(txn)
4015 .await?;
4016 let fragments = fragments
4017 .into_iter()
4018 .filter(|(_, fragment_type_mask, _)| {
4019 *fragment_type_mask & FragmentTypeFlag::Sink as i32 != 0
4020 })
4021 .filter_map(|(id, _, stream_node)| {
4022 let mut stream_node = stream_node.to_protobuf();
4023 let mut found = false;
4024 visit_stream_node_mut(&mut stream_node, |node| {
4025 if let PbNodeBody::Sink(node) = node
4026 && let Some(sink_desc) = &mut node.sink_desc
4027 && sink_desc.id == sink_id
4028 {
4029 sink_desc.properties.extend(props.clone());
4030 found = true;
4031 }
4032 });
4033 if found { Some((id, stream_node)) } else { None }
4034 })
4035 .collect_vec();
4036 assert!(
4037 !fragments.is_empty(),
4038 "sink id should be used by at least one fragment"
4039 );
4040 for (id, stream_node) in fragments {
4041 Fragment::update(fragment::ActiveModel {
4042 fragment_id: Set(id),
4043 stream_node: Set(StreamNode::from(&stream_node)),
4044 ..Default::default()
4045 })
4046 .exec(txn)
4047 .await?;
4048 }
4049 Ok(())
4050}
4051
4052pub struct SinkIntoTableContext {
4053 pub updated_sink_catalogs: Vec<SinkId>,
4056}
4057
4058pub struct FinishAutoRefreshSchemaSinkContext {
4059 pub tmp_sink_id: SinkId,
4060 pub original_sink_id: SinkId,
4061 pub columns: Vec<PbColumnCatalog>,
4062 pub new_log_store_table: Option<Box<PbTable>>,
4063}
4064
4065async fn update_connector_props_fragments<F>(
4066 txn: &DatabaseTransaction,
4067 job_ids: Vec<JobId>,
4068 expect_flag: FragmentTypeFlag,
4069 mut alter_stream_node_fn: F,
4070 is_shared_source: bool,
4071) -> MetaResult<()>
4072where
4073 F: FnMut(&mut PbNodeBody, &mut bool),
4074{
4075 let fragments: Vec<(FragmentId, StreamNode)> = Fragment::find()
4076 .select_only()
4077 .columns([fragment::Column::FragmentId, fragment::Column::StreamNode])
4078 .filter(
4079 fragment::Column::JobId
4080 .is_in(job_ids.clone())
4081 .and(FragmentTypeMask::intersects(expect_flag)),
4082 )
4083 .into_tuple()
4084 .all(txn)
4085 .await?;
4086 let fragments = fragments
4087 .into_iter()
4088 .filter_map(|(id, stream_node)| {
4089 let mut stream_node = stream_node.to_protobuf();
4090 let mut found = false;
4091 visit_stream_node_mut(&mut stream_node, |node| {
4092 alter_stream_node_fn(node, &mut found);
4093 });
4094 if found { Some((id, stream_node)) } else { None }
4095 })
4096 .collect_vec();
4097 if is_shared_source || job_ids.len() > 1 {
4098 assert!(
4102 !fragments.is_empty(),
4103 "job ids {:?} (type: {:?}) should be used by at least one fragment",
4104 job_ids,
4105 expect_flag
4106 );
4107 }
4108
4109 for (id, stream_node) in fragments {
4110 Fragment::update(fragment::ActiveModel {
4111 fragment_id: Set(id),
4112 stream_node: Set(StreamNode::from(&stream_node)),
4113 ..Default::default()
4114 })
4115 .exec(txn)
4116 .await?;
4117 }
4118
4119 Ok(())
4120}
4121
4122#[cfg(test)]
4123mod tests {
4124 use risingwave_sqlparser::ast::{SqlOption, Statement};
4125
4126 use super::{Parser, merge_with_options};
4127
4128 #[test]
4129 fn test_merge_with_options_normalizes_altered_option_name() {
4130 let mut statements = Parser::parse_sql(
4131 "CREATE SOURCE s WITH (properties.receive.message.max.bytes = 'old', \
4132 connection = kafka_conn) FORMAT PLAIN ENCODE JSON",
4133 )
4134 .unwrap();
4135 let Statement::CreateSource { stmt } = statements.remove(0) else {
4136 unreachable!()
4137 };
4138 let mut with_properties = stmt.with_properties.0;
4139 let altered_name = "properties.receive.message.max.bytes".to_owned();
4140 let altered_value = "new".to_owned();
4141
4142 merge_with_options(
4143 &mut with_properties,
4144 vec![SqlOption::try_from((&altered_name, &altered_value)).unwrap()],
4145 );
4146
4147 assert_eq!(with_properties.len(), 2);
4148 assert_eq!(
4149 with_properties[0].to_string(),
4150 "properties.receive.\"message\".\"max\".bytes = 'new'"
4151 );
4152 assert_eq!(with_properties[1].to_string(), "connection = kafka_conn");
4153 }
4154}