1use std::collections::{HashMap, HashSet};
16use std::pin::pin;
17use std::sync::Arc;
18
19use anyhow::{Context, anyhow};
20use futures::future::select;
21use rand::rng as thread_rng;
22use rand::seq::IndexedRandom;
23use replace_job_plan::{ReplaceSource, ReplaceTable};
24use risingwave_common::catalog::cdc_type_compatibility::cdc_auto_schema_change_existing_type_compatible;
25use risingwave_common::catalog::{AlterDatabaseParam, ColumnCatalog};
26use risingwave_common::id::{ObjectId, TableId};
27use risingwave_common::system_param::adaptive_parallelism_strategy::parse_strategy;
28use risingwave_common::types::DataType;
29use risingwave_common::util::stream_graph_visitor;
30use risingwave_connector::sink::catalog::SinkId;
31use risingwave_connector::sink::iceberg::ENABLE_PK_INDEX;
32use risingwave_meta::barrier::{BarrierScheduler, Command, ResumeBackfillTarget};
33use risingwave_meta::manager::{EventLogManagerRef, MetadataManager, iceberg_compaction};
34use risingwave_meta::model::TableParallelism as ModelTableParallelism;
35use risingwave_meta::rpc::metrics::MetaMetrics;
36use risingwave_meta::stream::{ParallelismPolicy, ReschedulePolicy, ResourceGroupPolicy};
37use risingwave_meta::{MetaResult, bail_invalid_parameter, bail_unavailable};
38use risingwave_meta_model::StreamingParallelism;
39use risingwave_pb::catalog::connection::Info as ConnectionInfo;
40use risingwave_pb::catalog::table::{CdcTableType as PbCdcTableType, OptionalAssociatedSourceId};
41use risingwave_pb::catalog::{Comment, Connection, PbCreateType, Secret, Table};
42use risingwave_pb::common::WorkerType;
43use risingwave_pb::common::worker_node::State;
44use risingwave_pb::ddl_service::create_iceberg_table_request::{PbSinkJobInfo, PbTableJobInfo};
45use risingwave_pb::ddl_service::ddl_service_server::DdlService;
46use risingwave_pb::ddl_service::drop_table_request::PbSourceId;
47use risingwave_pb::ddl_service::replace_job_plan::ReplaceMaterializedView;
48use risingwave_pb::ddl_service::{streaming_job_resource_type, *};
49use risingwave_pb::frontend_service::GetTableReplacePlanRequest;
50use risingwave_pb::meta::event_log;
51use risingwave_pb::meta::table_parallelism::{FixedParallelism, Parallelism};
52use risingwave_pb::stream_plan::stream_node::NodeBody;
53use risingwave_pb::stream_plan::throttle_mutation::ThrottleConfig;
54use thiserror_ext::AsReport;
55use tokio::sync::oneshot::Sender;
56use tokio::task::JoinHandle;
57use tonic::{Request, Response, Status};
58
59use crate::MetaError;
60use crate::barrier::BarrierManagerRef;
61use crate::manager::iceberg_pk_index_sink::IcebergPkIndexSinkManager;
62use crate::manager::sink_coordination::SinkCoordinatorManager;
63use crate::manager::{MetaSrvEnv, StreamingJob};
64use crate::rpc::ddl_controller::{
65 DdlCommand, DdlController, DropMode, ReplaceStreamJobInfo, StreamingJobId,
66};
67use crate::stream::{GlobalStreamManagerRef, SourceManagerRef};
68
69#[derive(Clone)]
70pub struct DdlServiceImpl {
71 env: MetaSrvEnv,
72
73 metadata_manager: MetadataManager,
74 sink_manager: SinkCoordinatorManager,
75 ddl_controller: DdlController,
76 meta_metrics: Arc<MetaMetrics>,
77 iceberg_compaction_manager: iceberg_compaction::IcebergCompactionManagerRef,
78 barrier_scheduler: BarrierScheduler,
79 iceberg_pk_index_sink_manager: IcebergPkIndexSinkManager,
80}
81
82impl DdlServiceImpl {
83 pub async fn new(
84 env: MetaSrvEnv,
85 metadata_manager: MetadataManager,
86 stream_manager: GlobalStreamManagerRef,
87 source_manager: SourceManagerRef,
88 barrier_manager: BarrierManagerRef,
89 sink_manager: SinkCoordinatorManager,
90 meta_metrics: Arc<MetaMetrics>,
91 iceberg_compaction_manager: iceberg_compaction::IcebergCompactionManagerRef,
92 barrier_scheduler: BarrierScheduler,
93 iceberg_pk_index_sink_manager: IcebergPkIndexSinkManager,
94 ) -> Self {
95 let ddl_controller = DdlController::new(
96 env.clone(),
97 metadata_manager.clone(),
98 stream_manager,
99 source_manager,
100 barrier_manager,
101 sink_manager.clone(),
102 iceberg_compaction_manager.clone(),
103 iceberg_pk_index_sink_manager.clone(),
104 )
105 .await;
106 Self {
107 env,
108 metadata_manager,
109 sink_manager,
110 ddl_controller,
111 meta_metrics,
112 iceberg_compaction_manager,
113 barrier_scheduler,
114 iceberg_pk_index_sink_manager,
115 }
116 }
117
118 fn extract_replace_table_info(
119 ReplaceJobPlan {
120 fragment_graph,
121 replace_job,
122 }: ReplaceJobPlan,
123 ) -> ReplaceStreamJobInfo {
124 let replace_streaming_job: StreamingJob = match replace_job.unwrap() {
125 replace_job_plan::ReplaceJob::ReplaceTable(ReplaceTable {
126 table,
127 source,
128 job_type,
129 }) => StreamingJob::Table(
130 source,
131 table.unwrap(),
132 TableJobType::try_from(job_type).unwrap(),
133 ),
134 replace_job_plan::ReplaceJob::ReplaceSource(ReplaceSource { source }) => {
135 StreamingJob::Source(source.unwrap())
136 }
137 replace_job_plan::ReplaceJob::ReplaceMaterializedView(ReplaceMaterializedView {
138 table,
139 }) => StreamingJob::MaterializedView(table.unwrap()),
140 replace_job_plan::ReplaceJob::ReplaceSink(_) => unreachable!("use replace sink path"),
141 };
142
143 ReplaceStreamJobInfo {
144 streaming_job: replace_streaming_job,
145 fragment_graph: fragment_graph.unwrap(),
146 }
147 }
148
149 fn default_streaming_job_resource_type() -> streaming_job_resource_type::ResourceType {
150 streaming_job_resource_type::ResourceType::Regular(true)
151 }
152
153 pub fn start_migrate_table_fragments(&self) -> (JoinHandle<()>, Sender<()>) {
154 tracing::info!("start migrate legacy table fragments task");
155 let env = self.env.clone();
156 let metadata_manager = self.metadata_manager.clone();
157 let ddl_controller = self.ddl_controller.clone();
158
159 let (shutdown_tx, shutdown_rx) = tokio::sync::oneshot::channel();
160 let join_handle = tokio::spawn(async move {
161 async fn migrate_inner(
162 env: &MetaSrvEnv,
163 metadata_manager: &MetadataManager,
164 ddl_controller: &DdlController,
165 ) -> MetaResult<()> {
166 let tables = metadata_manager
167 .catalog_controller
168 .list_unmigrated_tables()
169 .await?;
170
171 if tables.is_empty() {
172 tracing::info!("no legacy table fragments need migration");
173 return Ok(());
174 }
175
176 let client = {
177 let workers = metadata_manager
178 .list_worker_node(Some(WorkerType::Frontend), Some(State::Running))
179 .await?;
180 if workers.is_empty() {
181 return Err(anyhow::anyhow!("no active frontend nodes found").into());
182 }
183 let worker = workers.choose(&mut thread_rng()).unwrap();
184 env.frontend_client_pool().get(worker).await?
185 };
186
187 for table in tables {
188 let start = tokio::time::Instant::now();
189 let req = GetTableReplacePlanRequest {
190 database_id: table.database_id,
191 table_id: table.id,
192 cdc_table_change: None,
193 };
194 let resp = client
195 .get_table_replace_plan(req)
196 .await
197 .context("failed to get table replace plan from frontend")?;
198
199 let plan = resp.into_inner().replace_plan.unwrap();
200 let replace_info = DdlServiceImpl::extract_replace_table_info(plan);
201 ddl_controller
202 .run_command(DdlCommand::ReplaceStreamJob(replace_info))
203 .await?;
204 tracing::info!(elapsed=?start.elapsed(), table_id=%table.id, "migrated table fragments");
205 }
206 tracing::info!("successfully migrated all legacy table fragments");
207
208 Ok(())
209 }
210
211 let migrate_future = async move {
212 let mut attempt = 0;
213 loop {
214 match migrate_inner(&env, &metadata_manager, &ddl_controller).await {
215 Ok(_) => break,
216 Err(e) => {
217 attempt += 1;
218 tracing::error!(
219 "failed to migrate legacy table fragments: {}, attempt {}, retrying in 5 secs",
220 e.as_report(),
221 attempt
222 );
223 tokio::time::sleep(std::time::Duration::from_secs(5)).await;
224 }
225 }
226 }
227 };
228
229 select(pin!(migrate_future), shutdown_rx).await;
230 });
231
232 (join_handle, shutdown_tx)
233 }
234}
235
236#[async_trait::async_trait]
237impl DdlService for DdlServiceImpl {
238 async fn create_database(
239 &self,
240 request: Request<CreateDatabaseRequest>,
241 ) -> Result<Response<CreateDatabaseResponse>, Status> {
242 let req = request.into_inner();
243 let database = req.get_db()?.clone();
244 let version = self
245 .ddl_controller
246 .run_command(DdlCommand::CreateDatabase(database))
247 .await?;
248
249 Ok(Response::new(CreateDatabaseResponse {
250 status: None,
251 version,
252 }))
253 }
254
255 async fn drop_database(
256 &self,
257 request: Request<DropDatabaseRequest>,
258 ) -> Result<Response<DropDatabaseResponse>, Status> {
259 let req = request.into_inner();
260 let database_id = req.get_database_id();
261
262 let version = self
263 .ddl_controller
264 .run_command(DdlCommand::DropDatabase(database_id))
265 .await?;
266
267 Ok(Response::new(DropDatabaseResponse {
268 status: None,
269 version,
270 }))
271 }
272
273 async fn create_secret(
274 &self,
275 request: Request<CreateSecretRequest>,
276 ) -> Result<Response<CreateSecretResponse>, Status> {
277 let req = request.into_inner();
278 let pb_secret = Secret {
279 id: 0.into(),
280 name: req.get_name().clone(),
281 database_id: req.get_database_id(),
282 value: req.get_value().clone(),
283 owner: req.get_owner_id(),
284 schema_id: req.get_schema_id(),
285 };
286 let version = self
287 .ddl_controller
288 .run_command(DdlCommand::CreateSecret(pb_secret))
289 .await?;
290
291 Ok(Response::new(CreateSecretResponse { version }))
292 }
293
294 async fn drop_secret(
295 &self,
296 request: Request<DropSecretRequest>,
297 ) -> Result<Response<DropSecretResponse>, Status> {
298 let req = request.into_inner();
299 let secret_id = req.get_secret_id();
300 let drop_mode = DropMode::from_request_setting(req.cascade);
301 let version = self
302 .ddl_controller
303 .run_command(DdlCommand::DropSecret(secret_id, drop_mode))
304 .await?;
305
306 Ok(Response::new(DropSecretResponse {
307 status: None,
308 version,
309 }))
310 }
311
312 async fn alter_secret(
313 &self,
314 request: Request<AlterSecretRequest>,
315 ) -> Result<Response<AlterSecretResponse>, Status> {
316 let req = request.into_inner();
317 let pb_secret = Secret {
318 id: req.get_secret_id(),
319 name: req.get_name().clone(),
320 database_id: req.get_database_id(),
321 value: req.get_value().clone(),
322 owner: req.get_owner_id(),
323 schema_id: req.get_schema_id(),
324 };
325 let version = self
326 .ddl_controller
327 .run_command(DdlCommand::AlterSecret(pb_secret))
328 .await?;
329
330 Ok(Response::new(AlterSecretResponse { version }))
331 }
332
333 async fn create_schema(
334 &self,
335 request: Request<CreateSchemaRequest>,
336 ) -> Result<Response<CreateSchemaResponse>, Status> {
337 let req = request.into_inner();
338 let schema = req.get_schema()?.clone();
339 let version = self
340 .ddl_controller
341 .run_command(DdlCommand::CreateSchema(schema))
342 .await?;
343
344 Ok(Response::new(CreateSchemaResponse {
345 status: None,
346 version,
347 }))
348 }
349
350 async fn drop_schema(
351 &self,
352 request: Request<DropSchemaRequest>,
353 ) -> Result<Response<DropSchemaResponse>, Status> {
354 let req = request.into_inner();
355 let schema_id = req.get_schema_id();
356 let drop_mode = DropMode::from_request_setting(req.cascade);
357 let version = self
358 .ddl_controller
359 .run_command(DdlCommand::DropSchema(schema_id, drop_mode))
360 .await?;
361 Ok(Response::new(DropSchemaResponse {
362 status: None,
363 version,
364 }))
365 }
366
367 async fn create_source(
368 &self,
369 request: Request<CreateSourceRequest>,
370 ) -> Result<Response<CreateSourceResponse>, Status> {
371 let req = request.into_inner();
372 let source = req.get_source()?.clone();
373
374 match req.fragment_graph {
375 None => {
376 let version = self
377 .ddl_controller
378 .run_command(DdlCommand::CreateNonSharedSource(source, None))
379 .await?;
380 Ok(Response::new(CreateSourceResponse {
381 status: None,
382 version,
383 }))
384 }
385 Some(fragment_graph) => {
386 let stream_job = StreamingJob::Source(source);
388 let version = self
389 .ddl_controller
390 .run_command(DdlCommand::CreateStreamingJob {
391 stream_job,
392 fragment_graph,
393 dependencies: HashSet::new(),
394 resource_type: Self::default_streaming_job_resource_type(),
395 if_not_exists: req.if_not_exists,
396 refresh_interval_sec: None,
397 replace_sink: None,
398 since_timestamp_epoch: None,
399 })
400 .await?;
401 Ok(Response::new(CreateSourceResponse {
402 status: None,
403 version,
404 }))
405 }
406 }
407 }
408
409 async fn drop_source(
410 &self,
411 request: Request<DropSourceRequest>,
412 ) -> Result<Response<DropSourceResponse>, Status> {
413 let request = request.into_inner();
414 let source_id = request.source_id;
415 let drop_mode = DropMode::from_request_setting(request.cascade);
416 let version = self
417 .ddl_controller
418 .run_command(DdlCommand::DropSource(source_id, drop_mode))
419 .await?;
420
421 Ok(Response::new(DropSourceResponse {
422 status: None,
423 version,
424 }))
425 }
426
427 async fn reset_source(
428 &self,
429 request: Request<ResetSourceRequest>,
430 ) -> Result<Response<ResetSourceResponse>, Status> {
431 let request = request.into_inner();
432 let source_id = request.source_id;
433
434 tracing::info!(
435 source_id = %source_id,
436 "Received RESET SOURCE request, routing to DDL controller"
437 );
438
439 let version = self
441 .ddl_controller
442 .run_command(DdlCommand::ResetSource(source_id))
443 .await?;
444
445 Ok(Response::new(ResetSourceResponse {
446 status: None,
447 version,
448 }))
449 }
450
451 async fn create_sink(
452 &self,
453 request: Request<CreateSinkRequest>,
454 ) -> Result<Response<CreateSinkResponse>, Status> {
455 self.env.idle_manager().record_activity();
456
457 let req = request.into_inner();
458
459 let sink = req.get_sink()?.clone();
460 let fragment_graph = req.get_fragment_graph()?.clone();
461 let dependencies = req.get_dependencies().iter().copied().collect();
462 let resource_type = req
463 .resource_type
464 .and_then(|resource_type| resource_type.resource_type)
465 .unwrap_or_else(Self::default_streaming_job_resource_type);
466 let since_timestamp_epoch = req.since_timestamp_epoch;
467
468 let stream_job = StreamingJob::Sink(sink, None);
469
470 let command = DdlCommand::CreateStreamingJob {
471 stream_job,
472 fragment_graph,
473 dependencies,
474 resource_type,
475 if_not_exists: req.if_not_exists,
476 refresh_interval_sec: None,
477 replace_sink: None,
478 since_timestamp_epoch,
479 };
480
481 let version = self.ddl_controller.run_command(command).await?;
482
483 Ok(Response::new(CreateSinkResponse {
484 status: None,
485 version,
486 }))
487 }
488
489 async fn drop_sink(
490 &self,
491 request: Request<DropSinkRequest>,
492 ) -> Result<Response<DropSinkResponse>, Status> {
493 let request = request.into_inner();
494 let sink_id = request.sink_id;
495 let drop_mode = DropMode::from_request_setting(request.cascade);
496
497 let command = DdlCommand::DropStreamingJob {
498 job_id: StreamingJobId::Sink(sink_id),
499 drop_mode,
500 };
501
502 let version = self.ddl_controller.run_command(command).await?;
503
504 self.sink_manager
505 .stop_sink_coordinator(vec![SinkId::from(sink_id)])
506 .await;
507 self.iceberg_compaction_manager
508 .clear_iceberg_maintenance_by_sink_id(SinkId::from(sink_id));
509
510 Ok(Response::new(DropSinkResponse {
511 status: None,
512 version,
513 }))
514 }
515
516 async fn create_subscription(
517 &self,
518 request: Request<CreateSubscriptionRequest>,
519 ) -> Result<Response<CreateSubscriptionResponse>, Status> {
520 self.env.idle_manager().record_activity();
521
522 let req = request.into_inner();
523
524 let subscription = req.get_subscription()?.clone();
525 let command = DdlCommand::CreateSubscription(subscription);
526
527 let version = self.ddl_controller.run_command(command).await?;
528
529 Ok(Response::new(CreateSubscriptionResponse {
530 status: None,
531 version,
532 }))
533 }
534
535 async fn drop_subscription(
536 &self,
537 request: Request<DropSubscriptionRequest>,
538 ) -> Result<Response<DropSubscriptionResponse>, Status> {
539 let request = request.into_inner();
540 let subscription_id = request.subscription_id;
541 let drop_mode = DropMode::from_request_setting(request.cascade);
542
543 let command = DdlCommand::DropSubscription(subscription_id, drop_mode);
544
545 let version = self.ddl_controller.run_command(command).await?;
546
547 Ok(Response::new(DropSubscriptionResponse {
548 status: None,
549 version,
550 }))
551 }
552
553 async fn create_materialized_view(
554 &self,
555 request: Request<CreateMaterializedViewRequest>,
556 ) -> Result<Response<CreateMaterializedViewResponse>, Status> {
557 self.env.idle_manager().record_activity();
558
559 let req = request.into_inner();
560 let mview = req.get_materialized_view()?.clone();
561 let fragment_graph = req.get_fragment_graph()?.clone();
562 let dependencies = req.get_dependencies().iter().copied().collect();
563 let resource_type = req.resource_type.unwrap().resource_type.unwrap();
564
565 let stream_job = StreamingJob::MaterializedView(mview);
566 let version = self
567 .ddl_controller
568 .run_command(DdlCommand::CreateStreamingJob {
569 stream_job,
570 fragment_graph,
571 dependencies,
572 resource_type,
573 if_not_exists: req.if_not_exists,
574 refresh_interval_sec: req.refresh_interval_sec,
575 replace_sink: None,
576 since_timestamp_epoch: None,
577 })
578 .await?;
579
580 Ok(Response::new(CreateMaterializedViewResponse {
581 status: None,
582 version,
583 }))
584 }
585
586 async fn drop_materialized_view(
587 &self,
588 request: Request<DropMaterializedViewRequest>,
589 ) -> Result<Response<DropMaterializedViewResponse>, Status> {
590 self.env.idle_manager().record_activity();
591
592 let request = request.into_inner();
593 let table_id = request.table_id;
594 let drop_mode = DropMode::from_request_setting(request.cascade);
595
596 let version = self
597 .ddl_controller
598 .run_command(DdlCommand::DropStreamingJob {
599 job_id: StreamingJobId::MaterializedView(table_id),
600 drop_mode,
601 })
602 .await?;
603
604 Ok(Response::new(DropMaterializedViewResponse {
605 status: None,
606 version,
607 }))
608 }
609
610 async fn create_index(
611 &self,
612 request: Request<CreateIndexRequest>,
613 ) -> Result<Response<CreateIndexResponse>, Status> {
614 self.env.idle_manager().record_activity();
615
616 let req = request.into_inner();
617 let index = req.get_index()?.clone();
618 let index_table = req.get_index_table()?.clone();
619 let fragment_graph = req.get_fragment_graph()?.clone();
620 let resource_type = req
621 .resource_type
622 .and_then(|resource_type| resource_type.resource_type)
623 .unwrap_or_else(Self::default_streaming_job_resource_type);
624
625 let stream_job = StreamingJob::Index(index, index_table);
626 let version = self
627 .ddl_controller
628 .run_command(DdlCommand::CreateStreamingJob {
629 stream_job,
630 fragment_graph,
631 dependencies: HashSet::new(),
632 resource_type,
633 if_not_exists: req.if_not_exists,
634 replace_sink: None,
635 refresh_interval_sec: None,
636 since_timestamp_epoch: None,
637 })
638 .await?;
639
640 Ok(Response::new(CreateIndexResponse {
641 status: None,
642 version,
643 }))
644 }
645
646 async fn drop_index(
647 &self,
648 request: Request<DropIndexRequest>,
649 ) -> Result<Response<DropIndexResponse>, Status> {
650 self.env.idle_manager().record_activity();
651
652 let request = request.into_inner();
653 let index_id = request.index_id;
654 let drop_mode = DropMode::from_request_setting(request.cascade);
655 let version = self
656 .ddl_controller
657 .run_command(DdlCommand::DropStreamingJob {
658 job_id: StreamingJobId::Index(index_id),
659 drop_mode,
660 })
661 .await?;
662
663 Ok(Response::new(DropIndexResponse {
664 status: None,
665 version,
666 }))
667 }
668
669 async fn create_function(
670 &self,
671 request: Request<CreateFunctionRequest>,
672 ) -> Result<Response<CreateFunctionResponse>, Status> {
673 let req = request.into_inner();
674 let function = req.get_function()?.clone();
675
676 let version = self
677 .ddl_controller
678 .run_command(DdlCommand::CreateFunction(function))
679 .await?;
680
681 Ok(Response::new(CreateFunctionResponse {
682 status: None,
683 version,
684 }))
685 }
686
687 async fn drop_function(
688 &self,
689 request: Request<DropFunctionRequest>,
690 ) -> Result<Response<DropFunctionResponse>, Status> {
691 let request = request.into_inner();
692
693 let version = self
694 .ddl_controller
695 .run_command(DdlCommand::DropFunction(
696 request.function_id,
697 DropMode::from_request_setting(request.cascade),
698 ))
699 .await?;
700
701 Ok(Response::new(DropFunctionResponse {
702 status: None,
703 version,
704 }))
705 }
706
707 async fn create_table(
708 &self,
709 request: Request<CreateTableRequest>,
710 ) -> Result<Response<CreateTableResponse>, Status> {
711 let request = request.into_inner();
712 let job_type = request.get_job_type().unwrap_or_default();
713 let dependencies = request.get_dependencies().iter().copied().collect();
714 let source = request.source;
715 let mview = request.materialized_view.unwrap();
716 let fragment_graph = request.fragment_graph.unwrap();
717
718 let stream_job = StreamingJob::Table(source, mview, job_type);
719 let version = self
720 .ddl_controller
721 .run_command(DdlCommand::CreateStreamingJob {
722 stream_job,
723 fragment_graph,
724 dependencies,
725 resource_type: Self::default_streaming_job_resource_type(),
726 if_not_exists: request.if_not_exists,
727 refresh_interval_sec: None,
728 replace_sink: None,
729 since_timestamp_epoch: None,
730 })
731 .await?;
732
733 Ok(Response::new(CreateTableResponse {
734 status: None,
735 version,
736 }))
737 }
738
739 async fn drop_table(
740 &self,
741 request: Request<DropTableRequest>,
742 ) -> Result<Response<DropTableResponse>, Status> {
743 let request = request.into_inner();
744 let source_id = request.source_id;
745 let table_id = request.table_id;
746
747 let drop_mode = DropMode::from_request_setting(request.cascade);
748 let version = self
749 .ddl_controller
750 .run_command(DdlCommand::DropStreamingJob {
751 job_id: StreamingJobId::Table(source_id.map(|PbSourceId::Id(id)| id), table_id),
752 drop_mode,
753 })
754 .await?;
755
756 Ok(Response::new(DropTableResponse {
757 status: None,
758 version,
759 }))
760 }
761
762 async fn create_view(
763 &self,
764 request: Request<CreateViewRequest>,
765 ) -> Result<Response<CreateViewResponse>, Status> {
766 let req = request.into_inner();
767 let view = req.get_view()?.clone();
768 let dependencies = req
769 .get_dependencies()
770 .iter()
771 .copied()
772 .collect::<HashSet<_>>();
773
774 let version = self
775 .ddl_controller
776 .run_command(DdlCommand::CreateView(view, dependencies))
777 .await?;
778
779 Ok(Response::new(CreateViewResponse {
780 status: None,
781 version,
782 }))
783 }
784
785 async fn drop_view(
786 &self,
787 request: Request<DropViewRequest>,
788 ) -> Result<Response<DropViewResponse>, Status> {
789 let request = request.into_inner();
790 let view_id = request.get_view_id();
791 let drop_mode = DropMode::from_request_setting(request.cascade);
792 let version = self
793 .ddl_controller
794 .run_command(DdlCommand::DropView(view_id, drop_mode))
795 .await?;
796 Ok(Response::new(DropViewResponse {
797 status: None,
798 version,
799 }))
800 }
801
802 async fn risectl_list_state_tables(
803 &self,
804 _request: Request<RisectlListStateTablesRequest>,
805 ) -> Result<Response<RisectlListStateTablesResponse>, Status> {
806 let tables = self
807 .metadata_manager
808 .catalog_controller
809 .list_all_state_tables()
810 .await?;
811 Ok(Response::new(RisectlListStateTablesResponse { tables }))
812 }
813
814 async fn risectl_resume_backfill(
815 &self,
816 request: Request<RisectlResumeBackfillRequest>,
817 ) -> Result<Response<RisectlResumeBackfillResponse>, Status> {
818 let request = request.into_inner();
819 let target = request
820 .target
821 .ok_or_else(|| Status::invalid_argument("missing resume backfill target"))?;
822
823 match target {
824 risectl_resume_backfill_request::Target::JobId(job_id) => {
825 let database_id = self
826 .metadata_manager
827 .catalog_controller
828 .get_object_database_id(ObjectId::new(job_id.as_raw_id()))
829 .await?;
830 self.barrier_scheduler
831 .run_command(
832 database_id,
833 Command::ResumeBackfill {
834 target: ResumeBackfillTarget::Job(job_id),
835 },
836 )
837 .await?;
838 }
839 risectl_resume_backfill_request::Target::FragmentId(fragment_id) => {
840 let mut job_ids = self
841 .metadata_manager
842 .catalog_controller
843 .get_fragment_job_id(vec![fragment_id])
844 .await?;
845 let job_id = job_ids
846 .pop()
847 .ok_or_else(|| Status::invalid_argument("fragment not found"))?;
848 let database_id = self
849 .metadata_manager
850 .catalog_controller
851 .get_object_database_id(ObjectId::new(job_id.as_raw_id()))
852 .await?;
853 self.barrier_scheduler
854 .run_command(
855 database_id,
856 Command::ResumeBackfill {
857 target: ResumeBackfillTarget::Fragment(fragment_id),
858 },
859 )
860 .await?;
861 }
862 }
863
864 Ok(Response::new(RisectlResumeBackfillResponse {}))
865 }
866
867 async fn replace_job_plan(
868 &self,
869 request: Request<ReplaceJobPlanRequest>,
870 ) -> Result<Response<ReplaceJobPlanResponse>, Status> {
871 let ReplaceJobPlan {
872 fragment_graph,
873 replace_job,
874 } = request.into_inner().get_plan().cloned()?;
875
876 let command = match replace_job {
877 Some(replace_job_plan::ReplaceJob::ReplaceSink(replace_sink)) => {
878 let replace_job_plan::ReplaceSink {
879 sink,
880 old_sink_id,
881 dependencies,
882 resource_type,
883 } = replace_sink;
884 DdlCommand::CreateStreamingJob {
885 stream_job: StreamingJob::Sink(sink.unwrap(), None),
886 fragment_graph: fragment_graph.unwrap(),
887 dependencies: dependencies.into_iter().collect::<HashSet<_>>(),
888 resource_type: resource_type
889 .and_then(|resource_type| resource_type.resource_type)
890 .unwrap_or(streaming_job_resource_type::ResourceType::Regular(true)),
891 if_not_exists: false,
892 refresh_interval_sec: None,
893 replace_sink: Some(old_sink_id),
894 since_timestamp_epoch: None,
895 }
896 }
897 replace_job => {
898 DdlCommand::ReplaceStreamJob(Self::extract_replace_table_info(ReplaceJobPlan {
899 fragment_graph,
900 replace_job,
901 }))
902 }
903 };
904
905 let version = self.ddl_controller.run_command(command).await?;
906
907 Ok(Response::new(ReplaceJobPlanResponse {
908 status: None,
909 version,
910 }))
911 }
912
913 async fn get_table(
914 &self,
915 request: Request<GetTableRequest>,
916 ) -> Result<Response<GetTableResponse>, Status> {
917 let req = request.into_inner();
918 let table = self
919 .metadata_manager
920 .catalog_controller
921 .get_table_by_name(&req.database_name, &req.table_name)
922 .await?;
923
924 Ok(Response::new(GetTableResponse { table }))
925 }
926
927 async fn alter_name(
928 &self,
929 request: Request<AlterNameRequest>,
930 ) -> Result<Response<AlterNameResponse>, Status> {
931 let AlterNameRequest { object, new_name } = request.into_inner();
932 let version = self
933 .ddl_controller
934 .run_command(DdlCommand::AlterName(object.unwrap(), new_name))
935 .await?;
936 Ok(Response::new(AlterNameResponse {
937 status: None,
938 version,
939 }))
940 }
941
942 async fn alter_source(
944 &self,
945 request: Request<AlterSourceRequest>,
946 ) -> Result<Response<AlterSourceResponse>, Status> {
947 let AlterSourceRequest { source } = request.into_inner();
948 let version = self
949 .ddl_controller
950 .run_command(DdlCommand::AlterNonSharedSource(source.unwrap()))
951 .await?;
952 Ok(Response::new(AlterSourceResponse {
953 status: None,
954 version,
955 }))
956 }
957
958 async fn alter_owner(
959 &self,
960 request: Request<AlterOwnerRequest>,
961 ) -> Result<Response<AlterOwnerResponse>, Status> {
962 let AlterOwnerRequest { object, owner_id } = request.into_inner();
963 let version = self
964 .ddl_controller
965 .run_command(DdlCommand::AlterObjectOwner(object.unwrap(), owner_id as _))
966 .await?;
967 Ok(Response::new(AlterOwnerResponse {
968 status: None,
969 version,
970 }))
971 }
972
973 async fn alter_subscription_retention(
974 &self,
975 request: Request<AlterSubscriptionRetentionRequest>,
976 ) -> Result<Response<AlterSubscriptionRetentionResponse>, Status> {
977 let AlterSubscriptionRetentionRequest {
978 subscription_id,
979 retention_seconds,
980 definition,
981 } = request.into_inner();
982 let version = self
983 .ddl_controller
984 .run_command(DdlCommand::AlterSubscriptionRetention {
985 subscription_id,
986 retention_seconds,
987 definition,
988 })
989 .await?;
990 Ok(Response::new(AlterSubscriptionRetentionResponse {
991 status: None,
992 version,
993 }))
994 }
995
996 async fn alter_set_schema(
997 &self,
998 request: Request<AlterSetSchemaRequest>,
999 ) -> Result<Response<AlterSetSchemaResponse>, Status> {
1000 let AlterSetSchemaRequest {
1001 object,
1002 new_schema_id,
1003 } = request.into_inner();
1004 let version = self
1005 .ddl_controller
1006 .run_command(DdlCommand::AlterSetSchema(object.unwrap(), new_schema_id))
1007 .await?;
1008 Ok(Response::new(AlterSetSchemaResponse {
1009 status: None,
1010 version,
1011 }))
1012 }
1013
1014 async fn get_ddl_progress(
1015 &self,
1016 _request: Request<GetDdlProgressRequest>,
1017 ) -> Result<Response<GetDdlProgressResponse>, Status> {
1018 Ok(Response::new(GetDdlProgressResponse {
1019 ddl_progress: self.ddl_controller.get_ddl_progress().await?,
1020 }))
1021 }
1022
1023 async fn create_connection(
1024 &self,
1025 request: Request<CreateConnectionRequest>,
1026 ) -> Result<Response<CreateConnectionResponse>, Status> {
1027 let req = request.into_inner();
1028 if req.payload.is_none() {
1029 return Err(Status::invalid_argument("request is empty"));
1030 }
1031
1032 match req.payload.unwrap() {
1033 #[expect(deprecated)]
1034 create_connection_request::Payload::PrivateLink(_) => {
1035 panic!("Private Link Connection has been deprecated")
1036 }
1037 create_connection_request::Payload::ConnectionParams(params) => {
1038 let pb_connection = Connection {
1039 id: 0.into(),
1040 schema_id: req.schema_id,
1041 database_id: req.database_id,
1042 name: req.name,
1043 info: Some(ConnectionInfo::ConnectionParams(params)),
1044 owner: req.owner_id,
1045 };
1046 let version = self
1047 .ddl_controller
1048 .run_command(DdlCommand::CreateConnection(pb_connection))
1049 .await?;
1050 Ok(Response::new(CreateConnectionResponse { version }))
1051 }
1052 }
1053 }
1054
1055 async fn list_connections(
1056 &self,
1057 _request: Request<ListConnectionsRequest>,
1058 ) -> Result<Response<ListConnectionsResponse>, Status> {
1059 let conns = self
1060 .metadata_manager
1061 .catalog_controller
1062 .list_connections()
1063 .await?;
1064
1065 Ok(Response::new(ListConnectionsResponse {
1066 connections: conns,
1067 }))
1068 }
1069
1070 async fn drop_connection(
1071 &self,
1072 request: Request<DropConnectionRequest>,
1073 ) -> Result<Response<DropConnectionResponse>, Status> {
1074 let req = request.into_inner();
1075 let drop_mode = DropMode::from_request_setting(req.cascade);
1076
1077 let version = self
1078 .ddl_controller
1079 .run_command(DdlCommand::DropConnection(req.connection_id, drop_mode))
1080 .await?;
1081
1082 Ok(Response::new(DropConnectionResponse {
1083 status: None,
1084 version,
1085 }))
1086 }
1087
1088 async fn comment_on(
1089 &self,
1090 request: Request<CommentOnRequest>,
1091 ) -> Result<Response<CommentOnResponse>, Status> {
1092 let req = request.into_inner();
1093 let comment = req.get_comment()?.clone();
1094
1095 let version = self
1096 .ddl_controller
1097 .run_command(DdlCommand::CommentOn(Comment {
1098 table_id: comment.table_id,
1099 schema_id: comment.schema_id,
1100 database_id: comment.database_id,
1101 column_index: comment.column_index,
1102 description: comment.description,
1103 }))
1104 .await?;
1105
1106 Ok(Response::new(CommentOnResponse {
1107 status: None,
1108 version,
1109 }))
1110 }
1111
1112 async fn get_tables(
1113 &self,
1114 request: Request<GetTablesRequest>,
1115 ) -> Result<Response<GetTablesResponse>, Status> {
1116 let GetTablesRequest {
1117 table_ids,
1118 include_dropped_tables,
1119 } = request.into_inner();
1120 let ret = self
1121 .metadata_manager
1122 .catalog_controller
1123 .get_table_by_ids(table_ids, include_dropped_tables)
1124 .await?;
1125
1126 let mut tables = HashMap::default();
1127 for table in ret {
1128 tables.insert(table.id, table);
1129 }
1130 Ok(Response::new(GetTablesResponse { tables }))
1131 }
1132
1133 async fn wait(&self, request: Request<WaitRequest>) -> Result<Response<WaitResponse>, Status> {
1134 let req = request.into_inner();
1135 let version = self.ddl_controller.wait(req.job_id).await?;
1136 Ok(Response::new(WaitResponse {
1137 version: Some(version),
1138 }))
1139 }
1140
1141 async fn alter_cdc_table_backfill_parallelism(
1142 &self,
1143 request: Request<AlterCdcTableBackfillParallelismRequest>,
1144 ) -> Result<Response<AlterCdcTableBackfillParallelismResponse>, Status> {
1145 let req = request.into_inner();
1146 let job_id = req.get_table_id();
1147 let parallelism = *req.get_parallelism()?;
1148
1149 let table_parallelism = ModelTableParallelism::from(parallelism);
1150 let streaming_parallelism = match table_parallelism {
1151 ModelTableParallelism::Fixed(n) => StreamingParallelism::Fixed(n),
1152 _ => bail_invalid_parameter!(
1153 "CDC table backfill parallelism must be set to a fixed value"
1154 ),
1155 };
1156
1157 self.ddl_controller
1158 .reschedule_cdc_table_backfill(
1159 job_id,
1160 ReschedulePolicy::Parallelism(ParallelismPolicy {
1161 parallelism: streaming_parallelism,
1162 adaptive_parallelism_strategy: None,
1163 }),
1164 )
1165 .await?;
1166 Ok(Response::new(AlterCdcTableBackfillParallelismResponse {}))
1167 }
1168
1169 async fn alter_parallelism(
1170 &self,
1171 request: Request<AlterParallelismRequest>,
1172 ) -> Result<Response<AlterParallelismResponse>, Status> {
1173 let req = request.into_inner();
1174
1175 let job_id = req.get_table_id();
1176 let parallelism = *req.get_parallelism()?;
1177 let deferred = req.get_deferred();
1178
1179 let adaptive_parallelism_strategy = req.adaptive_parallelism_strategy;
1180 let (parallelism, adaptive_parallelism_strategy) = match parallelism.get_parallelism()? {
1181 Parallelism::Fixed(FixedParallelism { parallelism }) => {
1182 (StreamingParallelism::Fixed(*parallelism as _), None)
1183 }
1184 Parallelism::Auto(_) | Parallelism::Adaptive(_) => (
1185 StreamingParallelism::Adaptive,
1186 Some(adaptive_parallelism_strategy.unwrap_or_else(|| "AUTO".to_owned())),
1187 ),
1188 Parallelism::Custom(_) => (
1189 StreamingParallelism::Adaptive,
1190 Some(adaptive_parallelism_strategy.ok_or_else(|| {
1191 Status::invalid_argument(
1192 "adaptive_parallelism_strategy is required for custom parallelism",
1193 )
1194 })?),
1195 ),
1196 };
1197
1198 if let Some(strategy) = adaptive_parallelism_strategy.as_deref() {
1199 parse_strategy(strategy).map_err(|e| {
1200 Status::invalid_argument(format!(
1201 "invalid adaptive parallelism strategy: {}",
1202 e.as_report()
1203 ))
1204 })?;
1205 };
1206
1207 self.ddl_controller
1208 .reschedule_streaming_job(
1209 job_id,
1210 ReschedulePolicy::Parallelism(ParallelismPolicy {
1211 parallelism,
1212 adaptive_parallelism_strategy,
1213 }),
1214 deferred,
1215 )
1216 .await?;
1217
1218 Ok(Response::new(AlterParallelismResponse {}))
1219 }
1220
1221 async fn alter_backfill_parallelism(
1222 &self,
1223 request: Request<AlterBackfillParallelismRequest>,
1224 ) -> Result<Response<AlterBackfillParallelismResponse>, Status> {
1225 let req = request.into_inner();
1226
1227 let job_id = req.get_table_id();
1228 let deferred = req.get_deferred();
1229 let adaptive_parallelism_strategy = req.adaptive_parallelism_strategy;
1230
1231 let parallelism = match req.parallelism {
1232 None => None,
1233 Some(parallelism) => {
1234 let (parallelism, adaptive_parallelism_strategy) = match parallelism
1235 .get_parallelism()?
1236 {
1237 Parallelism::Fixed(FixedParallelism { parallelism }) => {
1238 (StreamingParallelism::Fixed(*parallelism as _), None)
1239 }
1240 Parallelism::Auto(_) | Parallelism::Adaptive(_) => (
1241 StreamingParallelism::Adaptive,
1242 Some(adaptive_parallelism_strategy.unwrap_or_else(|| "AUTO".to_owned())),
1243 ),
1244 Parallelism::Custom(_) => (
1245 StreamingParallelism::Adaptive,
1246 Some(adaptive_parallelism_strategy.ok_or_else(|| {
1247 Status::invalid_argument(
1248 "adaptive_parallelism_strategy is required for custom parallelism",
1249 )
1250 })?),
1251 ),
1252 };
1253
1254 if let Some(strategy) = adaptive_parallelism_strategy.as_deref() {
1255 parse_strategy(strategy).map_err(|e| {
1256 Status::invalid_argument(format!(
1257 "invalid adaptive parallelism strategy: {}",
1258 e.as_report()
1259 ))
1260 })?;
1261 };
1262
1263 Some(ParallelismPolicy {
1264 parallelism,
1265 adaptive_parallelism_strategy,
1266 })
1267 }
1268 };
1269
1270 self.ddl_controller
1271 .reschedule_streaming_job_backfill_parallelism(job_id, parallelism, deferred)
1272 .await?;
1273
1274 Ok(Response::new(AlterBackfillParallelismResponse {}))
1275 }
1276
1277 async fn alter_fragment_parallelism(
1278 &self,
1279 request: Request<AlterFragmentParallelismRequest>,
1280 ) -> Result<Response<AlterFragmentParallelismResponse>, Status> {
1281 let req = request.into_inner();
1282
1283 let fragment_ids = req.fragment_ids;
1284 if fragment_ids.is_empty() {
1285 return Err(Status::invalid_argument(
1286 "at least one fragment id must be provided",
1287 ));
1288 }
1289
1290 let parallelism = match req.parallelism {
1291 Some(parallelism) => {
1292 let streaming_parallelism = match parallelism.get_parallelism()? {
1293 Parallelism::Fixed(FixedParallelism { parallelism }) => {
1294 StreamingParallelism::Fixed(*parallelism as _)
1295 }
1296 Parallelism::Auto(_) | Parallelism::Adaptive(_) => {
1297 StreamingParallelism::Adaptive
1298 }
1299 _ => bail_unavailable!(),
1300 };
1301 Some(streaming_parallelism)
1302 }
1303 None => None,
1304 };
1305
1306 let fragment_targets = fragment_ids
1307 .into_iter()
1308 .map(|fragment_id| (fragment_id, parallelism.clone()))
1309 .collect();
1310
1311 self.ddl_controller
1312 .reschedule_fragments(fragment_targets)
1313 .await?;
1314
1315 Ok(Response::new(AlterFragmentParallelismResponse {}))
1316 }
1317
1318 async fn alter_streaming_job_config(
1319 &self,
1320 request: Request<AlterStreamingJobConfigRequest>,
1321 ) -> Result<Response<AlterStreamingJobConfigResponse>, Status> {
1322 let AlterStreamingJobConfigRequest {
1323 job_id,
1324 entries_to_add,
1325 keys_to_remove,
1326 } = request.into_inner();
1327
1328 self.ddl_controller
1329 .run_command(DdlCommand::AlterStreamingJobConfig(
1330 job_id,
1331 entries_to_add,
1332 keys_to_remove,
1333 ))
1334 .await?;
1335
1336 Ok(Response::new(AlterStreamingJobConfigResponse {}))
1337 }
1338
1339 async fn auto_schema_change(
1342 &self,
1343 request: Request<AutoSchemaChangeRequest>,
1344 ) -> Result<Response<AutoSchemaChangeResponse>, Status> {
1345 let req = request.into_inner();
1346
1347 let workers = self
1349 .metadata_manager
1350 .list_worker_node(Some(WorkerType::Frontend), Some(State::Running))
1351 .await?;
1352 let worker = workers
1353 .choose(&mut thread_rng())
1354 .ok_or_else(|| MetaError::from(anyhow!("no frontend worker available")))?;
1355
1356 let client = self
1357 .env
1358 .frontend_client_pool()
1359 .get(worker)
1360 .await
1361 .map_err(MetaError::from)?;
1362
1363 let Some(schema_change) = req.schema_change else {
1364 return Err(Status::invalid_argument(
1365 "schema change message is required",
1366 ));
1367 };
1368
1369 for table_change in schema_change.table_changes {
1370 for c in &table_change.columns {
1371 let c = ColumnCatalog::from(c.clone());
1372
1373 let invalid_col_type = |column_type: &str, c: &ColumnCatalog| {
1374 tracing::warn!(target: "auto_schema_change",
1375 cdc_table_id = table_change.cdc_table_id,
1376 upstraem_ddl = table_change.upstream_ddl,
1377 "invalid column type from cdc table change");
1378 Err(Status::invalid_argument(format!(
1379 "invalid column type: {} from cdc table change, column: {:?}",
1380 column_type, c
1381 )))
1382 };
1383 if c.is_generated() {
1384 return invalid_col_type("generated column", &c);
1385 }
1386 if c.is_rw_sys_column() {
1387 return invalid_col_type("rw system column", &c);
1388 }
1389 if c.is_hidden {
1390 return invalid_col_type("hidden column", &c);
1391 }
1392 }
1393
1394 let tables: Vec<Table> = self
1396 .metadata_manager
1397 .get_table_catalog_by_cdc_table_id(&table_change.cdc_table_id)
1398 .await?;
1399
1400 for table in tables {
1401 let original_columns_by_name: HashMap<String, ColumnCatalog> = table
1404 .columns
1405 .iter()
1406 .filter_map(|col| {
1407 let col = ColumnCatalog::from(col.clone());
1408 cdc_auto_schema_change_comparable_column(&col).map(|(name, _)| (name, col))
1409 })
1410 .collect();
1411
1412 let original_column_types: HashMap<String, DataType> = original_columns_by_name
1413 .iter()
1414 .map(|(name, col)| (name.clone(), col.data_type().clone()))
1415 .collect();
1416
1417 let original_column_names: HashSet<String> =
1418 HashSet::from_iter(original_column_types.keys().cloned());
1419
1420 let cdc_table_type =
1421 PbCdcTableType::try_from(table.cdc_table_type.unwrap_or_default())
1422 .unwrap_or(PbCdcTableType::Unspecified);
1423 let table_change = normalize_cdc_auto_schema_change_column_names(
1424 table_change.clone(),
1425 cdc_table_type,
1426 &original_columns_by_name,
1427 );
1428
1429 let collect_new_columns = |table_change: &TableSchemaChange| {
1430 let mut new_columns: HashSet<(String, DataType)> =
1431 HashSet::from_iter(table_change.columns.iter().filter_map(|col| {
1432 let col = ColumnCatalog::from(col.clone());
1433 cdc_auto_schema_change_comparable_column(&col)
1434 }));
1435
1436 for col in &table.columns {
1439 let col = ColumnCatalog::from(col.clone());
1440 if col.is_connector_additional_column()
1441 && !col.is_hidden()
1442 && !col.is_generated()
1443 {
1444 new_columns
1445 .insert((col.column_desc.name.clone(), col.data_type().clone()));
1446 }
1447 }
1448
1449 new_columns
1450 };
1451
1452 let new_columns = collect_new_columns(&table_change);
1453
1454 let new_column_names: HashSet<String> =
1455 HashSet::from_iter(new_columns.iter().map(|(name, _)| name.clone()));
1456 let is_add_or_drop_by_name = original_column_names.is_subset(&new_column_names)
1457 || original_column_names.is_superset(&new_column_names);
1458
1459 let table_change =
1462 if original_column_names != new_column_names && is_add_or_drop_by_name {
1463 normalize_cdc_auto_schema_change_existing_column_types(
1464 table_change,
1465 cdc_table_type,
1466 &original_columns_by_name,
1467 )
1468 } else {
1469 table_change
1473 };
1474
1475 let original_columns: HashSet<(String, DataType)> = HashSet::from_iter(
1476 original_column_types
1477 .iter()
1478 .map(|(name, data_type)| (name.clone(), data_type.clone())),
1479 );
1480
1481 let new_columns = collect_new_columns(&table_change);
1482
1483 if !(original_columns.is_subset(&new_columns)
1484 || original_columns.is_superset(&new_columns))
1485 {
1486 tracing::warn!(target: "auto_schema_change",
1487 table_id = %table.id,
1488 cdc_table_id = table.cdc_table_id,
1489 upstraem_ddl = table_change.upstream_ddl,
1490 original_columns = ?original_columns,
1491 new_columns = ?new_columns,
1492 "New columns should be a subset or superset of the original columns (including hidden columns), since only `ADD COLUMN` and `DROP COLUMN` is supported");
1493
1494 let fail_info = "New columns should be a subset or superset of the original columns (including hidden columns), since only `ADD COLUMN` and `DROP COLUMN` is supported".to_owned();
1495 add_auto_schema_change_fail_event_log(
1496 &self.meta_metrics,
1497 table.id,
1498 table.name.clone(),
1499 table_change.cdc_table_id.clone(),
1500 table_change.upstream_ddl.clone(),
1501 &self.env.event_log_manager_ref(),
1502 fail_info,
1503 );
1504
1505 return Err(Status::invalid_argument(
1506 "New columns should be a subset or superset of the original columns (including hidden columns)",
1507 ));
1508 }
1509 if original_columns == new_columns {
1511 tracing::warn!(target: "auto_schema_change",
1512 table_id = %table.id,
1513 cdc_table_id = table.cdc_table_id,
1514 upstraem_ddl = table_change.upstream_ddl,
1515 original_columns = ?original_columns,
1516 new_columns = ?new_columns,
1517 "No change to columns, skipping the schema change");
1518 continue;
1519 }
1520
1521 let latency_timer = self
1522 .meta_metrics
1523 .auto_schema_change_latency
1524 .with_label_values(&[&table.id.to_string(), &table.name])
1525 .start_timer();
1526 let resp = client
1529 .get_table_replace_plan(GetTableReplacePlanRequest {
1530 database_id: table.database_id,
1531 table_id: table.id,
1532 cdc_table_change: Some(table_change.clone()),
1533 })
1534 .await;
1535
1536 match resp {
1537 Ok(resp) => {
1538 let resp = resp.into_inner();
1539 if let Some(plan) = resp.replace_plan {
1540 let plan = Self::extract_replace_table_info(plan);
1541 plan.streaming_job.table().inspect(|t| {
1542 tracing::info!(
1543 target: "auto_schema_change",
1544 table_id = %t.id,
1545 cdc_table_id = t.cdc_table_id,
1546 upstraem_ddl = table_change.upstream_ddl,
1547 "Start the replace config change")
1548 });
1549 let replace_res = self
1551 .ddl_controller
1552 .run_command(DdlCommand::ReplaceStreamJob(plan))
1553 .await;
1554
1555 match replace_res {
1556 Ok(_) => {
1557 tracing::info!(
1558 target: "auto_schema_change",
1559 table_id = %table.id,
1560 cdc_table_id = table.cdc_table_id,
1561 "Table replaced success");
1562
1563 self.meta_metrics
1564 .auto_schema_change_success_cnt
1565 .with_label_values(&[&table.id.to_string(), &table.name])
1566 .inc();
1567 latency_timer.observe_duration();
1568 }
1569 Err(e) => {
1570 tracing::error!(
1571 target: "auto_schema_change",
1572 error = %e.as_report(),
1573 table_id = %table.id,
1574 cdc_table_id = table.cdc_table_id,
1575 upstraem_ddl = table_change.upstream_ddl,
1576 "failed to replace the table",
1577 );
1578 let fail_info =
1579 format!("failed to replace the table: {}", e.as_report());
1580 add_auto_schema_change_fail_event_log(
1581 &self.meta_metrics,
1582 table.id,
1583 table.name.clone(),
1584 table_change.cdc_table_id.clone(),
1585 table_change.upstream_ddl.clone(),
1586 &self.env.event_log_manager_ref(),
1587 fail_info,
1588 );
1589 }
1590 };
1591 }
1592 }
1593 Err(e) => {
1594 tracing::error!(
1595 target: "auto_schema_change",
1596 error = %e.as_report(),
1597 table_id = %table.id,
1598 cdc_table_id = table.cdc_table_id,
1599 "failed to get replace table plan",
1600 );
1601 let fail_info =
1602 format!("failed to get replace table plan: {}", e.as_report());
1603 add_auto_schema_change_fail_event_log(
1604 &self.meta_metrics,
1605 table.id,
1606 table.name.clone(),
1607 table_change.cdc_table_id.clone(),
1608 table_change.upstream_ddl.clone(),
1609 &self.env.event_log_manager_ref(),
1610 fail_info,
1611 );
1612 }
1613 };
1614 }
1615 }
1616
1617 Ok(Response::new(AutoSchemaChangeResponse {}))
1618 }
1619
1620 async fn wait_iceberg_pk_index_sink_epoch(
1621 &self,
1622 request: Request<WaitIcebergPkIndexSinkEpochRequest>,
1623 ) -> Result<Response<WaitIcebergPkIndexSinkEpochResponse>, Status> {
1624 let req = request.into_inner();
1625 let snapshot_id = self
1626 .iceberg_pk_index_sink_manager
1627 .wait_epoch(req.sink_id, req.epoch)
1628 .await
1629 .map_err(|e| {
1630 Status::internal(format!(
1631 "Failed to wait for pk-index sink epoch: {}",
1632 e.as_report()
1633 ))
1634 })?;
1635 Ok(Response::new(WaitIcebergPkIndexSinkEpochResponse {
1636 snapshot_id,
1637 }))
1638 }
1639
1640 async fn alter_swap_rename(
1641 &self,
1642 request: Request<AlterSwapRenameRequest>,
1643 ) -> Result<Response<AlterSwapRenameResponse>, Status> {
1644 let req = request.into_inner();
1645
1646 let version = self
1647 .ddl_controller
1648 .run_command(DdlCommand::AlterSwapRename(req.object.unwrap()))
1649 .await?;
1650
1651 Ok(Response::new(AlterSwapRenameResponse {
1652 status: None,
1653 version,
1654 }))
1655 }
1656
1657 async fn alter_resource_group(
1658 &self,
1659 request: Request<AlterResourceGroupRequest>,
1660 ) -> Result<Response<AlterResourceGroupResponse>, Status> {
1661 let req = request.into_inner();
1662
1663 let job_id = req.get_job_id();
1664 let deferred = req.get_deferred();
1665 let resource_group = req.resource_group;
1666
1667 self.ddl_controller
1668 .reschedule_streaming_job(
1669 job_id,
1670 ReschedulePolicy::ResourceGroup(ResourceGroupPolicy { resource_group }),
1671 deferred,
1672 )
1673 .await?;
1674
1675 Ok(Response::new(AlterResourceGroupResponse {}))
1676 }
1677
1678 async fn alter_database_resource_group(
1679 &self,
1680 request: Request<AlterDatabaseResourceGroupRequest>,
1681 ) -> Result<Response<AlterDatabaseResourceGroupResponse>, Status> {
1682 let req = request.into_inner();
1683
1684 let version = self
1685 .ddl_controller
1686 .run_command(DdlCommand::AlterDatabaseResourceGroup(
1687 req.database_id,
1688 req.resource_group,
1689 req.deferred,
1690 ))
1691 .await?;
1692
1693 Ok(Response::new(AlterDatabaseResourceGroupResponse {
1694 status: None,
1695 version,
1696 }))
1697 }
1698
1699 async fn alter_database_param(
1700 &self,
1701 request: Request<AlterDatabaseParamRequest>,
1702 ) -> Result<Response<AlterDatabaseParamResponse>, Status> {
1703 let req = request.into_inner();
1704 let database_id = req.database_id;
1705
1706 let param = match req.param.unwrap() {
1707 alter_database_param_request::Param::BarrierIntervalMs(value) => {
1708 AlterDatabaseParam::BarrierIntervalMs(value.value)
1709 }
1710 alter_database_param_request::Param::CheckpointFrequency(value) => {
1711 AlterDatabaseParam::CheckpointFrequency(value.value)
1712 }
1713 };
1714 let version = self
1715 .ddl_controller
1716 .run_command(DdlCommand::AlterDatabaseParam(database_id, param))
1717 .await?;
1718
1719 return Ok(Response::new(AlterDatabaseParamResponse {
1720 status: None,
1721 version,
1722 }));
1723 }
1724
1725 async fn compact_iceberg_table(
1726 &self,
1727 request: Request<CompactIcebergTableRequest>,
1728 ) -> Result<Response<CompactIcebergTableResponse>, Status> {
1729 let req = request.into_inner();
1730 let sink_id = req.sink_id;
1731
1732 let task_id = self
1734 .iceberg_compaction_manager
1735 .trigger_manual_compaction(sink_id)
1736 .await
1737 .map_err(|e| {
1738 Status::internal(format!("Failed to trigger compaction: {}", e.as_report()))
1739 })?;
1740
1741 Ok(Response::new(CompactIcebergTableResponse {
1742 status: None,
1743 task_id,
1744 }))
1745 }
1746
1747 async fn expire_iceberg_table_snapshots(
1748 &self,
1749 request: Request<ExpireIcebergTableSnapshotsRequest>,
1750 ) -> Result<Response<ExpireIcebergTableSnapshotsResponse>, Status> {
1751 let req = request.into_inner();
1752 let sink_id = req.sink_id;
1753
1754 self.iceberg_compaction_manager
1756 .check_and_expire_snapshots(sink_id)
1757 .await
1758 .map_err(|e| {
1759 Status::internal(format!("Failed to expire snapshots: {}", e.as_report()))
1760 })?;
1761
1762 Ok(Response::new(ExpireIcebergTableSnapshotsResponse {
1763 status: None,
1764 }))
1765 }
1766
1767 async fn rewrite_iceberg_table_manifests(
1768 &self,
1769 request: Request<RewriteIcebergTableManifestsRequest>,
1770 ) -> Result<Response<RewriteIcebergTableManifestsResponse>, Status> {
1771 let req = request.into_inner();
1772 let sink_id = req.sink_id;
1773
1774 self.iceberg_compaction_manager
1775 .check_and_rewrite_manifests(sink_id)
1776 .await
1777 .map_err(|e| {
1778 Status::internal(format!(
1779 "Failed to rewrite manifests for sink {}: {}",
1780 sink_id,
1781 e.as_report()
1782 ))
1783 })?;
1784
1785 Ok(Response::new(RewriteIcebergTableManifestsResponse {
1786 status: None,
1787 }))
1788 }
1789
1790 async fn create_iceberg_table(
1791 &self,
1792 request: Request<CreateIcebergTableRequest>,
1793 ) -> Result<Response<CreateIcebergTableResponse>, Status> {
1794 let req = request.into_inner();
1795 let CreateIcebergTableRequest {
1796 table_info,
1797 sink_info,
1798 iceberg_source,
1799 if_not_exists,
1800 } = req;
1801
1802 let PbTableJobInfo {
1804 source,
1805 table,
1806 fragment_graph,
1807 job_type,
1808 } = table_info.unwrap();
1809 let mut table = table.unwrap();
1810 let mut fragment_graph = fragment_graph.unwrap();
1811 let database_id = table.get_database_id();
1812 let schema_id = table.get_schema_id();
1813 let table_name = table.get_name().to_owned();
1814
1815 table.create_type = PbCreateType::Background as _;
1817
1818 let source_rate_limit = if let Some(source) = &source {
1820 for fragment in fragment_graph.fragments.values_mut() {
1821 stream_graph_visitor::visit_fragment_mut(fragment, |node| {
1822 if let NodeBody::Source(source_node) = node
1823 && let Some(inner) = &mut source_node.source_inner
1824 {
1825 inner.rate_limit = Some(0);
1826 }
1827 });
1828 }
1829 Some(source.rate_limit)
1830 } else {
1831 None
1832 };
1833
1834 let stream_job =
1835 StreamingJob::Table(source, table, PbTableJobType::try_from(job_type).unwrap());
1836 let _ = self
1837 .ddl_controller
1838 .run_command(DdlCommand::CreateStreamingJob {
1839 stream_job,
1840 fragment_graph,
1841 dependencies: HashSet::new(),
1842 resource_type: Self::default_streaming_job_resource_type(),
1843 if_not_exists,
1844 refresh_interval_sec: None,
1845 replace_sink: None,
1846 since_timestamp_epoch: None,
1847 })
1848 .await?;
1849
1850 let table_catalog = self
1851 .metadata_manager
1852 .catalog_controller
1853 .get_table_catalog_by_name(database_id, schema_id, &table_name)
1854 .await?
1855 .ok_or(Status::not_found("Internal error: table not found"))?;
1856
1857 let PbSinkJobInfo {
1859 sink,
1860 fragment_graph,
1861 } = sink_info.unwrap();
1862 let mut sink = sink.unwrap();
1863
1864 sink.create_type = PbCreateType::Background as _;
1866
1867 let enable_pk_index = sink
1868 .properties
1869 .get(ENABLE_PK_INDEX)
1870 .is_some_and(|v| v.eq_ignore_ascii_case("true"));
1871 sink.auto_refresh_schema_from_table = if enable_pk_index {
1874 None
1875 } else {
1876 Some(table_catalog.id)
1877 };
1878
1879 let mut fragment_graph = fragment_graph.unwrap();
1880
1881 assert_eq!(fragment_graph.dependent_table_ids.len(), 1);
1882 assert!(
1883 risingwave_common::catalog::TableId::from(fragment_graph.dependent_table_ids[0])
1884 .is_placeholder()
1885 );
1886 fragment_graph.dependent_table_ids[0] = table_catalog.id;
1887 for fragment in fragment_graph.fragments.values_mut() {
1888 stream_graph_visitor::visit_fragment_mut(fragment, |node| match node {
1889 NodeBody::StreamScan(scan) => {
1890 scan.table_id = table_catalog.id;
1891 if let Some(table_desc) = &mut scan.table_desc {
1892 assert!(
1893 risingwave_common::catalog::TableId::from(table_desc.table_id)
1894 .is_placeholder()
1895 );
1896 table_desc.table_id = table_catalog.id;
1897 table_desc.maybe_vnode_count = table_catalog.maybe_vnode_count;
1898 }
1899 if let Some(table) = &mut scan.arrangement_table {
1900 assert!(
1901 risingwave_common::catalog::TableId::from(table.id).is_placeholder()
1902 );
1903 *table = table_catalog.clone();
1904 }
1905 }
1906 NodeBody::BatchPlan(plan) => {
1907 if let Some(table_desc) = &mut plan.table_desc {
1908 assert!(
1909 risingwave_common::catalog::TableId::from(table_desc.table_id)
1910 .is_placeholder()
1911 );
1912 table_desc.table_id = table_catalog.id;
1913 table_desc.maybe_vnode_count = table_catalog.maybe_vnode_count;
1914 }
1915 }
1916 _ => {}
1917 });
1918 }
1919
1920 let table_id = table_catalog.id;
1921 let dependencies = HashSet::from_iter([table_id.into(), schema_id.into()]);
1922 let stream_job = StreamingJob::Sink(sink, Some(table_id));
1923 let res = self
1924 .ddl_controller
1925 .run_command(DdlCommand::CreateStreamingJob {
1926 stream_job,
1927 fragment_graph,
1928 dependencies,
1929 resource_type: Self::default_streaming_job_resource_type(),
1930 if_not_exists,
1931 refresh_interval_sec: None,
1932 replace_sink: None,
1933 since_timestamp_epoch: None,
1934 })
1935 .await;
1936
1937 if res.is_err() {
1938 let _ = self
1939 .ddl_controller
1940 .run_command(DdlCommand::DropStreamingJob {
1941 job_id: StreamingJobId::Table(None, table_id),
1942 drop_mode: DropMode::Cascade,
1943 })
1944 .await
1945 .inspect_err(|err| {
1946 tracing::error!(error = %err.as_report(),
1947 "Failed to clean up table after iceberg sink creation failure",
1948 );
1949 });
1950 res?;
1951 }
1952
1953 if let Some(source_rate_limit) = source_rate_limit
1955 && source_rate_limit != Some(0)
1956 {
1957 let OptionalAssociatedSourceId::AssociatedSourceId(source_id) =
1958 table_catalog.optional_associated_source_id.unwrap();
1959 let fragment_nodes = self
1960 .metadata_manager
1961 .update_source_rate_limit_by_source_id(source_id, source_rate_limit)
1962 .await?;
1963 let throttle_config = ThrottleConfig {
1964 throttle_type: risingwave_pb::common::ThrottleType::Source.into(),
1965 rate_limit: source_rate_limit,
1966 };
1967 let config = fragment_nodes
1968 .into_iter()
1969 .map(|(fragment_id, stream_node)| (fragment_id, (throttle_config, stream_node)))
1970 .collect();
1971 let _ = self
1972 .barrier_scheduler
1973 .run_command(database_id, Command::Throttle { config })
1974 .await?;
1975 }
1976
1977 let iceberg_source = iceberg_source.unwrap();
1979 let res = self
1980 .ddl_controller
1981 .run_command(DdlCommand::CreateNonSharedSource(
1982 iceberg_source,
1983 Some(table_id),
1984 ))
1985 .await;
1986 if res.is_err() {
1987 let _ = self
1988 .ddl_controller
1989 .run_command(DdlCommand::DropStreamingJob {
1990 job_id: StreamingJobId::Table(None, table_id),
1991 drop_mode: DropMode::Cascade,
1992 })
1993 .await
1994 .inspect_err(|err| {
1995 tracing::error!(
1996 error = %err.as_report(),
1997 "Failed to clean up table after iceberg source creation failure",
1998 );
1999 });
2000 }
2001
2002 Ok(Response::new(CreateIcebergTableResponse {
2003 status: None,
2004 version: res?,
2005 }))
2006 }
2007}
2008
2009fn add_auto_schema_change_fail_event_log(
2010 meta_metrics: &MetaMetrics,
2011 table_id: TableId,
2012 table_name: String,
2013 cdc_table_id: String,
2014 upstream_ddl: String,
2015 event_log_manager: &EventLogManagerRef,
2016 fail_info: String,
2017) {
2018 meta_metrics
2019 .auto_schema_change_failure_cnt
2020 .with_label_values(&[&table_id.to_string(), &table_name])
2021 .inc();
2022 let event = event_log::EventAutoSchemaChangeFail {
2023 table_id,
2024 table_name,
2025 cdc_table_id,
2026 upstream_ddl,
2027 fail_info,
2028 };
2029 event_log_manager.add_event_logs(vec![event_log::Event::AutoSchemaChangeFail(event)]);
2030}
2031
2032fn cdc_auto_schema_change_comparable_column(column: &ColumnCatalog) -> Option<(String, DataType)> {
2033 if column.is_generated() || column.is_hidden() {
2034 None
2035 } else {
2036 Some((column.column_desc.name.clone(), column.data_type().clone()))
2037 }
2038}
2039
2040fn normalize_cdc_auto_schema_change_column_names(
2041 mut table_change: TableSchemaChange,
2042 cdc_table_type: PbCdcTableType,
2043 original_columns_by_name: &HashMap<String, ColumnCatalog>,
2044) -> TableSchemaChange {
2045 if cdc_table_type != PbCdcTableType::Mysql {
2046 return table_change;
2047 }
2048
2049 let mut original_names_by_lowercase = HashMap::<String, Option<String>>::new();
2053 for original_name in original_columns_by_name.keys() {
2054 original_names_by_lowercase
2055 .entry(original_name.to_lowercase())
2056 .and_modify(|name| *name = None)
2057 .or_insert_with(|| Some(original_name.clone()));
2058 }
2059
2060 for column in &mut table_change.columns {
2061 let mut column_catalog = ColumnCatalog::from(column.clone());
2062 if column_catalog.is_generated() || column_catalog.is_hidden() {
2063 continue;
2064 }
2065
2066 let incoming_name = &column_catalog.column_desc.name;
2067 let lowercase_name = incoming_name.to_lowercase();
2068 let normalized_name = if original_columns_by_name.contains_key(incoming_name) {
2069 incoming_name.clone()
2070 } else {
2071 original_names_by_lowercase
2072 .get(&lowercase_name)
2073 .and_then(|name| name.clone())
2074 .unwrap_or(lowercase_name)
2075 };
2076
2077 column_catalog.column_desc.name = normalized_name;
2078 *column = column_catalog.to_protobuf();
2079 }
2080
2081 table_change
2082}
2083
2084fn normalize_cdc_auto_schema_change_existing_column_types(
2085 mut table_change: TableSchemaChange,
2086 cdc_table_type: PbCdcTableType,
2087 original_columns_by_name: &HashMap<String, ColumnCatalog>,
2088) -> TableSchemaChange {
2089 for column in &mut table_change.columns {
2090 let mut column_catalog = ColumnCatalog::from(column.clone());
2091 if column_catalog.is_generated() || column_catalog.is_hidden() {
2092 continue;
2093 }
2094
2095 if let Some(original_column) =
2096 original_columns_by_name.get(&column_catalog.column_desc.name)
2097 && cdc_auto_schema_change_existing_type_compatible(
2098 cdc_table_type,
2099 original_column.data_type(),
2100 column_catalog.data_type(),
2101 )
2102 {
2103 column_catalog.column_desc.data_type = original_column.data_type().clone();
2104 column_catalog.column_desc.generated_or_default_column = original_column
2105 .column_desc
2106 .generated_or_default_column
2107 .clone();
2108 *column = column_catalog.to_protobuf();
2109 }
2110 }
2111
2112 table_change
2113}
2114
2115#[cfg(test)]
2116mod tests {
2117 use risingwave_common::catalog::{ColumnDesc, ColumnId};
2118 use risingwave_common::types::ScalarImpl;
2119 use risingwave_pb::ddl_service::table_schema_change::TableChangeType;
2120 use risingwave_pb::plan_common::column_desc::GeneratedOrDefaultColumn;
2121
2122 use super::*;
2123
2124 fn pb_column(name: &str, data_type: DataType) -> risingwave_pb::plan_common::ColumnCatalog {
2125 ColumnCatalog::visible(ColumnDesc::named(name, ColumnId::placeholder(), data_type))
2126 .to_protobuf()
2127 }
2128
2129 fn pb_column_with_default_value(
2130 name: &str,
2131 data_type: DataType,
2132 snapshot_value: ScalarImpl,
2133 ) -> risingwave_pb::plan_common::ColumnCatalog {
2134 ColumnCatalog::visible(ColumnDesc::named_with_default_value(
2135 name,
2136 ColumnId::placeholder(),
2137 data_type,
2138 Some(snapshot_value),
2139 ))
2140 .to_protobuf()
2141 }
2142
2143 fn pb_table_change(
2144 columns: Vec<risingwave_pb::plan_common::ColumnCatalog>,
2145 ) -> TableSchemaChange {
2146 TableSchemaChange {
2147 change_type: TableChangeType::Alter as _,
2148 cdc_table_id: "1.db.t".to_owned(),
2149 columns,
2150 upstream_ddl: "ALTER TABLE t ADD COLUMN note VARCHAR(255)".to_owned(),
2151 }
2152 }
2153
2154 #[test]
2155 fn test_cdc_auto_schema_change_normalizes_compatible_existing_column_types() {
2156 let original_columns_by_name = HashMap::from([
2157 (
2158 "id".to_owned(),
2159 ColumnCatalog::visible(ColumnDesc::named(
2160 "id",
2161 ColumnId::placeholder(),
2162 DataType::Int64,
2163 )),
2164 ),
2165 (
2166 "v".to_owned(),
2167 ColumnCatalog::visible(ColumnDesc::named(
2168 "v",
2169 ColumnId::placeholder(),
2170 DataType::Varchar,
2171 )),
2172 ),
2173 ]);
2174 let table_change = pb_table_change(vec![
2175 pb_column("id", DataType::Int32),
2176 pb_column("v", DataType::Varchar),
2177 pb_column("note", DataType::Varchar),
2178 ]);
2179
2180 let normalized = normalize_cdc_auto_schema_change_existing_column_types(
2181 table_change,
2182 PbCdcTableType::Mysql,
2183 &original_columns_by_name,
2184 );
2185 let columns = normalized
2186 .columns
2187 .into_iter()
2188 .map(ColumnCatalog::from)
2189 .map(|column| (column.column_desc.name, column.column_desc.data_type))
2190 .collect::<HashMap<_, _>>();
2191
2192 assert_eq!(columns["id"], DataType::Int64);
2193 assert_eq!(columns["v"], DataType::Varchar);
2194 assert_eq!(columns["note"], DataType::Varchar);
2195 }
2196
2197 #[test]
2198 fn test_mysql_cdc_auto_schema_change_normalizes_column_names() {
2199 let original_columns_by_name = HashMap::from([
2200 (
2201 "Id".to_owned(),
2202 ColumnCatalog::visible(ColumnDesc::named(
2203 "Id",
2204 ColumnId::placeholder(),
2205 DataType::Int64,
2206 )),
2207 ),
2208 (
2209 "Details".to_owned(),
2210 ColumnCatalog::visible(ColumnDesc::named(
2211 "Details",
2212 ColumnId::placeholder(),
2213 DataType::Varchar,
2214 )),
2215 ),
2216 ]);
2217 let table_change = pb_table_change(vec![
2218 pb_column("ID", DataType::Int32),
2219 pb_column("details", DataType::Varchar),
2220 pb_column("NewCol", DataType::Varchar),
2221 ]);
2222
2223 let normalized = normalize_cdc_auto_schema_change_column_names(
2224 table_change,
2225 PbCdcTableType::Mysql,
2226 &original_columns_by_name,
2227 );
2228 let normalized = normalize_cdc_auto_schema_change_existing_column_types(
2229 normalized,
2230 PbCdcTableType::Mysql,
2231 &original_columns_by_name,
2232 );
2233 let columns = normalized
2234 .columns
2235 .into_iter()
2236 .map(ColumnCatalog::from)
2237 .map(|column| (column.column_desc.name, column.column_desc.data_type))
2238 .collect::<Vec<_>>();
2239
2240 assert_eq!(
2241 columns,
2242 vec![
2243 ("Id".to_owned(), DataType::Int64),
2244 ("Details".to_owned(), DataType::Varchar),
2245 ("newcol".to_owned(), DataType::Varchar),
2246 ]
2247 );
2248 }
2249
2250 #[test]
2251 fn test_non_mysql_cdc_auto_schema_change_preserves_column_names() {
2252 let original_columns_by_name = HashMap::from([(
2253 "Id".to_owned(),
2254 ColumnCatalog::visible(ColumnDesc::named(
2255 "Id",
2256 ColumnId::placeholder(),
2257 DataType::Int64,
2258 )),
2259 )]);
2260 let table_change = pb_table_change(vec![
2261 pb_column("ID", DataType::Int64),
2262 pb_column("NewCol", DataType::Varchar),
2263 ]);
2264
2265 let normalized = normalize_cdc_auto_schema_change_column_names(
2266 table_change,
2267 PbCdcTableType::Postgres,
2268 &original_columns_by_name,
2269 );
2270 let column_names = normalized
2271 .columns
2272 .into_iter()
2273 .map(ColumnCatalog::from)
2274 .map(|column| column.column_desc.name)
2275 .collect::<Vec<_>>();
2276
2277 assert_eq!(column_names, vec!["ID".to_owned(), "NewCol".to_owned()]);
2278 }
2279
2280 #[test]
2281 fn test_cdc_auto_schema_change_preserves_default_when_normalizing_type() {
2282 let original = ColumnCatalog::visible(ColumnDesc::named_with_default_value(
2283 "v",
2284 ColumnId::placeholder(),
2285 DataType::Int64,
2286 Some(ScalarImpl::Int64(1)),
2287 ));
2288 let original_columns_by_name = HashMap::from([("v".to_owned(), original.clone())]);
2289 let table_change = pb_table_change(vec![pb_column_with_default_value(
2290 "v",
2291 DataType::Int32,
2292 ScalarImpl::Int32(1),
2293 )]);
2294
2295 let normalized = normalize_cdc_auto_schema_change_existing_column_types(
2296 table_change,
2297 PbCdcTableType::Mysql,
2298 &original_columns_by_name,
2299 );
2300 let column = ColumnCatalog::from(normalized.columns.into_iter().next().unwrap());
2301
2302 assert_eq!(column.column_desc.data_type, DataType::Int64);
2303 assert_eq!(
2304 column.column_desc.generated_or_default_column,
2305 original.column_desc.generated_or_default_column
2306 );
2307 assert!(matches!(
2308 column.column_desc.generated_or_default_column,
2309 Some(GeneratedOrDefaultColumn::DefaultColumn(_))
2310 ));
2311 }
2312}