1use std::collections::{HashMap, HashSet};
16use std::sync::Arc;
17
18use anyhow::Context;
19use futures::StreamExt;
20use futures::stream::FuturesUnordered;
21use risingwave_common::catalog::{DatabaseId, FragmentTypeFlag, TableId};
22use risingwave_common::id::{JobId, PartialGraphId, SinkId};
23use risingwave_hummock_sdk::change_log::TableChangeLogs;
24use risingwave_meta_model::ActorId;
25use risingwave_meta_model::streaming_job::BackfillOrders;
26use risingwave_pb::common::WorkerNode;
27use risingwave_pb::hummock::HummockVersionStats;
28use risingwave_pb::id::SourceId;
29use risingwave_pb::meta::PbTableRefillRuntimeConfig;
30use risingwave_pb::meta::subscribe_response::{Info, Operation};
31use risingwave_pb::stream_service::barrier_complete_response::{
32 PbIcebergPkIndexSinkMetadata, PbListFinishedSource, PbLoadFinishedSource,
33};
34use risingwave_pb::stream_service::streaming_control_stream_request::PbInitRequest;
35use risingwave_rpc_client::StreamingControlHandle;
36use thiserror_ext::AsReport;
37
38use crate::barrier::cdc_progress::CdcTableBackfillTracker;
39use crate::barrier::checkpoint::independent_job::BatchRefreshJobTriggerContext;
40use crate::barrier::command::{
41 PostCollectCommand, ResumeBackfillTarget, SinceTimestampResolvedEpoch,
42};
43use crate::barrier::context::{GlobalBarrierWorkerContext, GlobalBarrierWorkerContextImpl};
44use crate::barrier::progress::TrackingJob;
45use crate::barrier::schedule::MarkReadyOptions;
46use crate::barrier::{
47 BarrierManagerStatus, BarrierWorkerRuntimeInfoSnapshot, BatchRefreshInfo, Command,
48 CreateStreamingJobCommandInfo, CreateStreamingJobType, DatabaseRuntimeInfoSnapshot,
49 RecoveryReason, ReplaceStreamJobPlan, Scheduled,
50};
51use crate::hummock::CommitEpochInfo;
52use crate::manager::LocalNotification;
53use crate::model::FragmentDownstreamRelation;
54use crate::serving::{fetch_serving_infos, sync_serving_table_vnode_mappings_to_hummock};
55use crate::stream::{SourceChange, cleanup_dropped_streaming_jobs};
56use crate::{MetaError, MetaResult};
57
58fn resolve_since_timestamp_log_store_epoch(
59 table_id: TableId,
60 since_epoch: u64,
61 upstream_committed_epoch: u64,
62 table_change_log: &TableChangeLogs,
63) -> MetaResult<SinceTimestampResolvedEpoch> {
64 let change_log = table_change_log.get(&table_id).ok_or_else(|| {
65 anyhow::anyhow!(
66 "no table changelog found for upstream table {} when resolving since_timestamp",
67 table_id
68 )
69 })?;
70 let Some(first_log) = change_log.first() else {
71 return Err(anyhow::anyhow!(
72 "empty table changelog found for upstream table {} when resolving since_timestamp",
73 table_id
74 )
75 .into());
76 };
77 let first_checkpoint_epoch = first_log.checkpoint_epoch;
78 if since_epoch < first_checkpoint_epoch {
79 return Err(anyhow::anyhow!(
80 "since_timestamp is earlier than the retained changelog of upstream table {}: requested epoch {}, first retained checkpoint epoch {}",
81 table_id,
82 since_epoch,
83 first_checkpoint_epoch,
84 )
85 .into());
86 }
87 if since_epoch >= upstream_committed_epoch {
88 return Err(anyhow::anyhow!(
89 "since_timestamp is not before the committed epoch of upstream table {}: requested epoch {}, committed epoch {}",
90 table_id,
91 since_epoch,
92 upstream_committed_epoch,
93 )
94 .into());
95 }
96 let latest_log = change_log.last().expect("checked non-empty");
97 if upstream_committed_epoch != latest_log.checkpoint_epoch {
98 return Err(anyhow::anyhow!(
99 "upstream committed epoch {} does not match latest changelog epoch {} for upstream table {}",
100 upstream_committed_epoch,
101 latest_log.checkpoint_epoch,
102 table_id,
103 )
104 .into());
105 }
106
107 let snapshot_epoch_index = match change_log.binary_search_by_checkpoint_epoch(since_epoch) {
115 Ok(index) => index,
116 Err(index) => index,
117 };
118 let snapshot_epoch = change_log
119 .get(snapshot_epoch_index)
120 .ok_or_else(|| {
121 anyhow::anyhow!(
122 "since_timestamp is later than the latest changelog of upstream table {}: requested epoch {}, latest changelog epoch {}",
123 table_id,
124 since_epoch,
125 upstream_committed_epoch,
126 )
127 })?
128 .checkpoint_epoch;
129 if snapshot_epoch >= upstream_committed_epoch {
133 return Err(anyhow::anyhow!(
134 "since_timestamp is too new for upstream table {}: requested epoch {}, resolved snapshot epoch {}, latest changelog epoch {}",
135 table_id,
136 since_epoch,
137 snapshot_epoch,
138 upstream_committed_epoch,
139 )
140 .into());
141 }
142 let epochs = change_log
143 .range((snapshot_epoch_index + 1)..)
144 .map(|epoch_log| {
145 (
146 epoch_log.non_checkpoint_epochs.clone(),
147 epoch_log.checkpoint_epoch,
148 )
149 })
150 .collect::<Vec<_>>();
151
152 Ok((snapshot_epoch, epochs))
153}
154
155impl GlobalBarrierWorkerContext for GlobalBarrierWorkerContextImpl {
156 #[await_tree::instrument]
157 async fn commit_epoch(&self, commit_info: CommitEpochInfo) -> MetaResult<HummockVersionStats> {
158 self.hummock_manager.commit_epoch(commit_info).await?;
159 Ok(self.hummock_manager.get_version_stats().await)
160 }
161
162 #[await_tree::instrument("next_scheduled_barrier")]
163 async fn next_scheduled(&self) -> Scheduled {
164 self.scheduled_barriers.next_scheduled().await
165 }
166
167 fn abort_and_mark_blocked(
168 &self,
169 database_id: Option<DatabaseId>,
170 recovery_reason: RecoveryReason,
171 ) {
172 if database_id.is_none() {
173 self.set_status(BarrierManagerStatus::Recovering(recovery_reason));
174 }
175
176 self.scheduled_barriers
178 .abort_and_mark_blocked(database_id, "cluster is under recovering");
179 }
180
181 fn mark_ready(&self, options: MarkReadyOptions) {
182 let is_global = matches!(&options, MarkReadyOptions::Global { .. });
183 self.scheduled_barriers.mark_ready(options);
184 if is_global {
185 self.set_status(BarrierManagerStatus::Running);
186 }
187 }
188
189 async fn resolve_log_store_epoch<'a>(
190 &'a self,
191 upstream_table_ids: impl Iterator<Item = TableId> + Send + 'a,
192 since_epoch: u64,
193 ) -> MetaResult<SinceTimestampResolvedEpoch> {
194 let upstream_table_ids = upstream_table_ids.collect::<Vec<_>>();
195 if upstream_table_ids.is_empty() {
196 return Err(
197 anyhow::anyhow!("since_timestamp requires at least one upstream table").into(),
198 );
199 }
200
201 self.hummock_manager
202 .on_current_version_and_table_change_log(|version, table_change_log| {
203 let mut unified_log_epochs = None;
204 for &upstream_table_id in &upstream_table_ids {
205 let upstream_committed_epoch = version
206 .state_table_info
207 .info()
208 .get(&upstream_table_id)
209 .map(|info| info.committed_epoch)
210 .ok_or_else(|| {
211 anyhow::anyhow!(
212 "cannot get committed epoch for upstream table {}",
213 upstream_table_id
214 )
215 })?;
216 let table_log_epochs = resolve_since_timestamp_log_store_epoch(
217 upstream_table_id,
218 since_epoch,
219 upstream_committed_epoch,
220 table_change_log,
221 )?;
222 if let Some(unified_log_epochs) = &unified_log_epochs {
223 if unified_log_epochs != &table_log_epochs {
224 return Err(anyhow::anyhow!(
225 "resolved since_timestamp log epochs for upstream table {} do not match previous upstream table log epochs: {:?} vs {:?}",
226 upstream_table_id,
227 table_log_epochs,
228 unified_log_epochs,
229 )
230 .into());
231 }
232 } else {
233 unified_log_epochs = Some(table_log_epochs);
234 }
235 }
236 unified_log_epochs.ok_or_else(|| {
237 anyhow::anyhow!("since_timestamp requires at least one upstream table").into()
238 })
239 })
240 .await
241 }
242
243 async fn refresh_table_refill_runtime_state_after_recovery(&self) -> MetaResult<()> {
244 let policies = self
245 .metadata_manager
246 .catalog_controller
247 .table_cache_refill_policies_snapshot()
248 .await?;
249 let (serving_workers, fragment_serving_infos) =
250 fetch_serving_infos(&self.metadata_manager).await?;
251
252 self.env
253 .notification_manager()
254 .notify_hummock(
255 Operation::Update,
256 Info::TableRefillRuntimeConfig(PbTableRefillRuntimeConfig {
257 table_cache_refill_policies: Some(policies),
258 ..Default::default()
259 }),
260 )
261 .await;
262 sync_serving_table_vnode_mappings_to_hummock(
263 &self.env.notification_manager_ref(),
264 &self.serving_vnode_mapping,
265 &serving_workers,
266 &fragment_serving_infos,
267 )
268 .await;
269 Ok(())
270 }
271
272 #[await_tree::instrument("post_collect_command({command})")]
273 async fn post_collect_command(&self, command: PostCollectCommand) -> MetaResult<()> {
274 Box::pin(command.post_collect(self)).await
275 }
276
277 async fn notify_creating_job_failed(&self, database_id: Option<DatabaseId>, err: String) {
278 self.metadata_manager
279 .notify_finish_failed(database_id, err)
280 .await
281 }
282
283 #[await_tree::instrument("finish_creating_job({job})")]
284 async fn finish_creating_job(&self, job: TrackingJob) -> MetaResult<()> {
285 let job_id = job.job_id();
286 job.finish(&self.metadata_manager, &self.source_manager)
287 .await?;
288 self.env
289 .notification_manager()
290 .notify_local_subscribers(LocalNotification::StreamingJobBackfillFinished(job_id));
291 Ok(())
292 }
293
294 #[await_tree::instrument("finish_cdc_table_backfill({job})")]
295 async fn finish_cdc_table_backfill(&self, job: JobId) -> MetaResult<()> {
296 CdcTableBackfillTracker::mark_complete_job(&self.env.meta_store().conn, job).await
297 }
298
299 #[await_tree::instrument("new_control_stream({})", node.id)]
300 async fn new_control_stream(
301 &self,
302 node: &WorkerNode,
303 init_request: &PbInitRequest,
304 ) -> MetaResult<StreamingControlHandle> {
305 self.new_control_stream_impl(node, init_request).await
306 }
307
308 async fn reload_runtime_info(&self) -> MetaResult<BarrierWorkerRuntimeInfoSnapshot> {
309 self.reload_runtime_info_impl().await
310 }
311
312 async fn reload_database_runtime_info(
313 &self,
314 database_id: DatabaseId,
315 ) -> MetaResult<DatabaseRuntimeInfoSnapshot> {
316 self.reload_database_runtime_info_impl(database_id).await
317 }
318
319 async fn handle_list_finished_source_ids(
320 &self,
321 list_finished: Vec<PbListFinishedSource>,
322 ) -> MetaResult<()> {
323 let mut list_finished_info: HashMap<(TableId, SourceId), HashSet<ActorId>> = HashMap::new();
324
325 for list_finished in list_finished {
326 let table_id = list_finished.table_id;
327 let associated_source_id = list_finished.associated_source_id;
328 list_finished_info
329 .entry((table_id, associated_source_id))
330 .or_default()
331 .insert(list_finished.reporter_actor_id);
332 }
333
334 for ((table_id, associated_source_id), actors) in list_finished_info {
335 let allow_yield = self
336 .refresh_manager
337 .mark_list_stage_finished(table_id, &actors)?;
338
339 if !allow_yield {
340 continue;
341 }
342
343 let Some(database_id) = self
344 .get_source_database_id_for_refresh_stage(table_id, associated_source_id, "list")
345 .await?
346 else {
347 continue;
348 };
349
350 let list_finish_command = Command::ListFinish {
352 table_id,
353 associated_source_id,
354 };
355
356 self.barrier_scheduler
358 .run_command_no_wait(database_id, list_finish_command)
359 .context("Failed to schedule ListFinish command")?;
360
361 tracing::info!(
362 %table_id,
363 %associated_source_id,
364 "ListFinish command scheduled successfully"
365 );
366 }
367 Ok(())
368 }
369
370 async fn handle_load_finished_source_ids(
371 &self,
372 load_finished: Vec<PbLoadFinishedSource>,
373 ) -> MetaResult<()> {
374 let mut load_finished_info: HashMap<(TableId, SourceId), HashSet<ActorId>> = HashMap::new();
375
376 for load_finished in load_finished {
377 let table_id = load_finished.table_id;
378 let associated_source_id = load_finished.associated_source_id;
379 load_finished_info
380 .entry((table_id, associated_source_id))
381 .or_default()
382 .insert(load_finished.reporter_actor_id);
383 }
384
385 for ((table_id, associated_source_id), actors) in load_finished_info {
386 let allow_yield = self
387 .refresh_manager
388 .mark_load_stage_finished(table_id, &actors)?;
389
390 if !allow_yield {
391 continue;
392 }
393
394 let Some(database_id) = self
395 .get_source_database_id_for_refresh_stage(table_id, associated_source_id, "load")
396 .await?
397 else {
398 continue;
399 };
400
401 let load_finish_command = Command::LoadFinish {
403 table_id,
404 associated_source_id,
405 };
406
407 self.barrier_scheduler
409 .run_command_no_wait(database_id, load_finish_command)
410 .context("Failed to schedule LoadFinish command")?;
411
412 tracing::info!(
413 %table_id,
414 %associated_source_id,
415 "LoadFinish command scheduled successfully"
416 );
417 }
418
419 Ok(())
420 }
421
422 async fn handle_refresh_finished_table_ids(
423 &self,
424 refresh_finished_table_job_ids: Vec<JobId>,
425 ) -> MetaResult<()> {
426 for job_id in refresh_finished_table_job_ids {
427 let table_id = job_id.as_mv_table_id();
428
429 self.refresh_manager.mark_refresh_complete(table_id).await?;
430 }
431
432 Ok(())
433 }
434
435 async fn load_batch_refresh_trigger_context(
436 &self,
437 job_id: JobId,
438 database_id: DatabaseId,
439 last_committed_epoch: u64,
440 ) -> MetaResult<BatchRefreshJobTriggerContext> {
441 self.load_batch_refresh_trigger_context_impl(job_id, database_id, last_committed_epoch)
442 .await
443 }
444
445 #[await_tree::instrument]
446 async fn pre_commit_iceberg_pk_index_sink_metadata(
447 &self,
448 reports: Vec<PbIcebergPkIndexSinkMetadata>,
449 ) -> MetaResult<Vec<SinkId>> {
450 let grouped = group_reports_by_sink(reports)?;
451 let success_ids: Vec<SinkId> = grouped.keys().cloned().collect();
452 let futs = FuturesUnordered::new();
453 for (sink_id, (prev_epoch, reports)) in grouped {
454 if reports.is_empty() {
455 continue;
456 }
457 let manager = &self.iceberg_pk_index_sink_manager;
458 futs.push(async move {
459 (
460 sink_id,
461 manager.pre_commit_epoch(sink_id, prev_epoch, reports).await,
462 )
463 });
464 }
465
466 let results: Vec<(SinkId, anyhow::Result<()>)> = futs.collect().await;
469 let errs: Vec<(SinkId, anyhow::Error)> = results
470 .into_iter()
471 .filter_map(|(id, r)| r.err().map(|e| (id, e)))
472 .collect();
473 if errs.is_empty() {
474 Ok(success_ids)
475 } else {
476 Err(aggregate_sink_errors("pre-commit", errs).into())
477 }
478 }
479
480 #[await_tree::instrument]
481 async fn commit_iceberg_pk_index_sink_metadata(&self, sink_ids: Vec<SinkId>) -> MetaResult<()> {
482 let futs = FuturesUnordered::new();
483 for sink_id in sink_ids {
484 let manager = &self.iceberg_pk_index_sink_manager;
485 futs.push(async move { (sink_id, manager.commit_epoch(sink_id).await) });
486 }
487
488 let results: Vec<(SinkId, anyhow::Result<()>)> = futs.collect().await;
489 let errs: Vec<(SinkId, anyhow::Error)> = results
490 .into_iter()
491 .filter_map(|(id, r)| r.err().map(|e| (id, e)))
492 .collect();
493 if errs.is_empty() {
494 Ok(())
495 } else {
496 Err(aggregate_sink_errors("commit", errs).into())
497 }
498 }
499
500 fn advance_iceberg_pk_index_sink_committed_epochs(
501 &self,
502 epochs: impl IntoIterator<Item = (PartialGraphId, u64)>,
503 ) {
504 self.iceberg_pk_index_sink_manager
505 .advance_committed_epochs(epochs);
506 }
507}
508
509fn aggregate_sink_errors(
513 phase: &'static str,
514 mut errs: Vec<(SinkId, anyhow::Error)>,
515) -> anyhow::Error {
516 debug_assert!(!errs.is_empty());
517 let sink_ids: Vec<String> = errs.iter().map(|(id, _)| id.to_string()).collect();
518 let details: Vec<String> = errs
519 .iter()
520 .map(|(id, e)| format!("sink {}: {}", id, e.as_report()))
521 .collect();
522 let (_, first_err) = errs.remove(0);
524 first_err.context(format!(
525 "iceberg v3 sink {} failed for sink_id(s) [{}]: {}",
526 phase,
527 sink_ids.join(", "),
528 details.join("; ")
529 ))
530}
531
532fn group_reports_by_sink(
533 reports: Vec<PbIcebergPkIndexSinkMetadata>,
534) -> MetaResult<HashMap<SinkId, (u64, Vec<PbIcebergPkIndexSinkMetadata>)>> {
535 let mut grouped: HashMap<SinkId, (u64, Vec<PbIcebergPkIndexSinkMetadata>)> = HashMap::new();
536 for r in reports {
537 let sink_id = r.sink_id;
538 let prev_epoch = r.prev_epoch;
539 let entry = grouped.entry(sink_id).or_insert((prev_epoch, Vec::new()));
540 if entry.0 != prev_epoch {
541 return Err(anyhow::anyhow!(
542 "iceberg v3 sink {} reports disagree on prev_epoch: {} vs {}",
543 sink_id,
544 entry.0,
545 prev_epoch
546 )
547 .into());
548 }
549 entry.1.push(r);
550 }
551 Ok(grouped)
552}
553
554impl GlobalBarrierWorkerContextImpl {
555 async fn get_source_database_id_for_refresh_stage(
556 &self,
557 table_id: TableId,
558 associated_source_id: SourceId,
559 stage: &'static str,
560 ) -> MetaResult<Option<DatabaseId>> {
561 match self
562 .metadata_manager
563 .catalog_controller
564 .get_object_database_id(associated_source_id)
565 .await
566 {
567 Ok(database_id) => Ok(Some(database_id)),
568 Err(err) if err.is_catalog_id_not_found("object") => {
569 tracing::warn!(
570 %table_id,
571 %associated_source_id,
572 stage,
573 "skip refresh finish command because associated source is already dropped"
574 );
575 Ok(None)
576 }
577 Err(err) => Err(err)
578 .with_context(|| {
579 format!(
580 "failed to get database id for refresh stage: table_id={}, associated_source_id={}, stage={stage}",
581 table_id, associated_source_id
582 )
583 })
584 .map_err(Into::into),
585 }
586 }
587
588 fn set_status(&self, new_status: BarrierManagerStatus) {
589 self.status.store(Arc::new(new_status));
590 }
591
592 async fn load_batch_refresh_trigger_context_impl(
594 &self,
595 job_id: JobId,
596 database_id: DatabaseId,
597 last_committed_epoch: u64,
598 ) -> MetaResult<BatchRefreshJobTriggerContext> {
599 use itertools::Itertools;
600 use sea_orm::TransactionTrait;
601
602 use crate::controller::scale::load_fragment_context_for_jobs;
603
604 let inner = self
606 .metadata_manager
607 .catalog_controller
608 .get_inner_read_guard()
609 .await;
610 let txn = inner.db.begin().await?;
611
612 let fragment_context =
614 load_fragment_context_for_jobs(&txn, HashSet::from([job_id])).await?;
615
616 let streaming_job_model = fragment_context
617 .job_map
618 .get(&job_id)
619 .ok_or_else(|| anyhow::anyhow!("streaming job model not found for job {}", job_id))?
620 .clone();
621
622 let database_model = fragment_context
623 .database_map
624 .get(&database_id)
625 .ok_or_else(|| {
626 anyhow::anyhow!("database model not found for database {}", database_id)
627 })?;
628 let database_resource_group = database_model.resource_group.clone();
629
630 let mut job_extra_info = self
632 .metadata_manager
633 .catalog_controller
634 .get_streaming_job_extra_info_in_txn(&txn, vec![job_id])
635 .await?;
636 let definition = job_extra_info
637 .remove(&job_id)
638 .ok_or_else(|| anyhow::anyhow!("extra info not found for job {}", job_id))?
639 .job_definition;
640
641 let fragments = fragment_context
643 .job_fragments
644 .get(&job_id)
645 .ok_or_else(|| anyhow::anyhow!("fragments not found for job {}", job_id))?
646 .clone();
647
648 let upstream_table_ids: HashSet<TableId> = {
650 use crate::stream::StreamFragmentGraph;
651 let snapshot_backfill_info = StreamFragmentGraph::collect_snapshot_backfill_info_impl(
652 fragments.values().map(|f| (&f.nodes, f.fragment_type_mask)),
653 )?
654 .0
655 .ok_or_else(|| {
656 anyhow::anyhow!("batch refresh job {} has no snapshot backfill info", job_id)
657 })?;
658 snapshot_backfill_info
659 .upstream_mv_table_id_to_backfill_epoch
660 .into_keys()
661 .collect()
662 };
663
664 let fragment_ids: Vec<_> = fragments.keys().copied().collect();
665 let downstreams = self
666 .metadata_manager
667 .catalog_controller
668 .get_fragment_downstream_relations_in_txn(&txn, fragment_ids)
669 .await?;
670
671 txn.commit().await?;
672 drop(inner);
673
674 let (upstream_table_log_epochs, target_upstream_epoch) = self
676 .hummock_manager
677 .on_current_version_and_table_change_log(|version, table_change_log| {
678 let mut target_upstream_epoch = last_committed_epoch;
679 let mut log_epochs: HashMap<TableId, Vec<(Vec<u64>, u64)>> = HashMap::new();
680
681 for &upstream_table_id in &upstream_table_ids {
682 let upstream_committed_epoch = version
683 .state_table_info
684 .info()
685 .get(&upstream_table_id)
686 .map(|info| info.committed_epoch)
687 .ok_or_else(|| {
688 anyhow::anyhow!(
689 "cannot get committed epoch for upstream table {}",
690 upstream_table_id
691 )
692 })?;
693
694 target_upstream_epoch =
695 std::cmp::max(target_upstream_epoch, upstream_committed_epoch);
696
697 if upstream_committed_epoch <= last_committed_epoch {
698 continue;
699 }
700
701 if let Some(change_log) = table_change_log.get(&upstream_table_id) {
702 let epochs = change_log
703 .filter_epoch((last_committed_epoch, upstream_committed_epoch))
704 .map(|epoch_log| {
705 (
706 epoch_log.non_checkpoint_epochs.clone(),
707 epoch_log.checkpoint_epoch,
708 )
709 })
710 .collect_vec();
711 if !epochs.is_empty() {
712 log_epochs.insert(upstream_table_id, epochs);
713 }
714 } else {
715 anyhow::bail!(
716 "upstream table {} has lagged downstream on epoch {} but no table change log (upstream committed: {})",
717 upstream_table_id,
718 last_committed_epoch,
719 upstream_committed_epoch,
720 );
721 }
722 }
723
724 Ok((log_epochs, target_upstream_epoch))
725 })
726 .await?;
727
728 Ok(BatchRefreshJobTriggerContext {
729 fragments,
730 downstreams,
731 streaming_job_model,
732 definition,
733 database_resource_group,
734 upstream_table_log_epochs,
735 target_upstream_epoch,
736 })
737 }
738}
739
740impl PostCollectCommand {
741 pub async fn post_collect(
744 self,
745 barrier_manager_context: &GlobalBarrierWorkerContextImpl,
746 ) -> MetaResult<()> {
747 match self {
748 PostCollectCommand::Command(_) => {}
749 PostCollectCommand::SourceChangeSplit {
750 split_assignment: assignment,
751 } => {
752 barrier_manager_context
753 .metadata_manager
754 .update_fragment_splits(&assignment)
755 .await?;
756 }
757
758 PostCollectCommand::DropStreamingJobs => {}
759 PostCollectCommand::ConnectorPropsChange(obj_id_map_props) => {
760 barrier_manager_context
762 .source_manager
763 .apply_source_change(SourceChange::UpdateSourceProps {
764 source_id_map_new_props: obj_id_map_props
767 .iter()
768 .map(|(object_id, props)| (object_id.as_source_id(), props.clone()))
769 .collect(),
770 })
771 .await;
772 }
773 PostCollectCommand::ResumeBackfill { target } => match target {
774 ResumeBackfillTarget::Job(job_id) => {
775 barrier_manager_context
776 .metadata_manager
777 .catalog_controller
778 .update_backfill_orders_by_job_id(job_id, None)
779 .await?;
780 }
781 ResumeBackfillTarget::Fragment(fragment_id) => {
782 let mut job_ids = barrier_manager_context
783 .metadata_manager
784 .catalog_controller
785 .get_fragment_job_id(vec![fragment_id])
786 .await?;
787 let job_id = job_ids.pop().ok_or_else(|| {
788 MetaError::invalid_parameter("fragment not found".to_owned())
789 })?;
790 let job_id = JobId::new(job_id.as_raw_id());
791
792 let extra_info = barrier_manager_context
793 .metadata_manager
794 .catalog_controller
795 .get_streaming_job_extra_info(vec![job_id])
796 .await?;
797 let mut backfill_orders: BackfillOrders = extra_info
798 .get(&job_id)
799 .cloned()
800 .ok_or_else(|| MetaError::invalid_parameter("job not found".to_owned()))?
801 .backfill_orders
802 .unwrap_or_default();
803
804 let resumed_fragment_id = fragment_id.as_raw_id();
805 for children in backfill_orders.0.values_mut() {
806 children.retain(|child| *child != resumed_fragment_id);
807 }
808 backfill_orders.0.retain(|_, children| !children.is_empty());
809
810 barrier_manager_context
811 .metadata_manager
812 .catalog_controller
813 .update_backfill_orders_by_job_id(job_id, Some(backfill_orders))
814 .await?;
815 }
816 },
817 PostCollectCommand::CreateStreamingJob {
818 info,
819 job_type,
820 cross_db_snapshot_backfill_info,
821 resolved_split_assignment,
822 } => {
823 match &job_type {
824 CreateStreamingJobType::SinkIntoTable(_) | CreateStreamingJobType::Normal => {
825 barrier_manager_context
826 .metadata_manager
827 .catalog_controller
828 .fill_snapshot_backfill_epoch(
829 info.stream_job_fragments.fragments.iter().filter_map(
830 |(fragment_id, fragment)| {
831 if fragment.fragment_type_mask.contains(
832 FragmentTypeFlag::CrossDbSnapshotBackfillStreamScan,
833 ) {
834 Some(*fragment_id as _)
835 } else {
836 None
837 }
838 },
839 ),
840 None,
841 &cross_db_snapshot_backfill_info,
842 )
843 .await?
844 }
845 CreateStreamingJobType::SnapshotBackfill {
846 snapshot_backfill_info,
847 ..
848 }
849 | CreateStreamingJobType::BatchRefresh(BatchRefreshInfo {
850 snapshot_backfill_info,
851 ..
852 }) => {
853 barrier_manager_context
854 .metadata_manager
855 .catalog_controller
856 .fill_snapshot_backfill_epoch(
857 info.stream_job_fragments.fragments.iter().filter_map(
858 |(fragment_id, fragment)| {
859 if fragment.fragment_type_mask.contains_any([
860 FragmentTypeFlag::SnapshotBackfillStreamScan,
861 FragmentTypeFlag::CrossDbSnapshotBackfillStreamScan,
862 ]) {
863 Some(*fragment_id as _)
864 } else {
865 None
866 }
867 },
868 ),
869 Some(snapshot_backfill_info),
870 &cross_db_snapshot_backfill_info,
871 )
872 .await?
873 }
874 }
875
876 let CreateStreamingJobCommandInfo {
879 stream_job_fragments,
880 upstream_fragment_downstreams,
881 replace_sink,
882 ..
883 } = info;
884 let new_job_id = stream_job_fragments.stream_job_id();
885 let new_sink_downstream =
886 if let CreateStreamingJobType::SinkIntoTable(ctx) = job_type {
887 let new_downstreams = ctx.new_sink_downstream.clone();
888 let new_downstreams = FragmentDownstreamRelation::from([(
889 ctx.sink_fragment_id,
890 vec![new_downstreams],
891 )]);
892 Some(new_downstreams)
893 } else {
894 None
895 };
896
897 let old_state_table_ids = barrier_manager_context
898 .metadata_manager
899 .catalog_controller
900 .post_collect_job_fragments(
901 new_job_id,
902 &upstream_fragment_downstreams,
903 new_sink_downstream,
904 Some(&resolved_split_assignment),
905 replace_sink.as_ref(),
906 replace_sink.is_none(),
907 )
908 .await?;
909
910 let added_source_fragments = stream_job_fragments.stream_source_fragments();
911 let added_backfill_fragments = stream_job_fragments.source_backfill_fragments();
912 if !added_source_fragments.is_empty() || !added_backfill_fragments.is_empty() {
914 barrier_manager_context
915 .source_manager
916 .apply_source_change(SourceChange::CreateJob {
917 added_source_fragments,
918 added_backfill_fragments,
919 })
920 .await;
921 }
922
923 if let Some(old_sink_id) = replace_sink {
924 barrier_manager_context
925 .sink_manager
926 .stop_sink_coordinator(vec![old_sink_id])
927 .await;
928 cleanup_dropped_streaming_jobs(
929 &barrier_manager_context.refresh_manager,
930 &barrier_manager_context.hummock_manager,
931 &barrier_manager_context.metadata_manager,
932 [old_sink_id.as_job_id()],
933 old_state_table_ids.expect("replace sink should return old state tables"),
934 "replace_sink",
935 )
936 .await?;
937 }
938 }
939 PostCollectCommand::Reschedule { reschedules, .. } => {
940 let fragment_splits = reschedules
941 .iter()
942 .map(|(fragment_id, reschedule)| {
943 (*fragment_id, reschedule.actor_splits.clone())
944 })
945 .collect();
946
947 barrier_manager_context
948 .metadata_manager
949 .update_fragment_splits(&fragment_splits)
950 .await?;
951 }
952
953 PostCollectCommand::ReplaceStreamJob {
954 plan: replace_plan,
955 resolved_split_assignment,
956 } => {
957 let ReplaceStreamJobPlan {
958 old_fragments,
959 new_fragments,
960 upstream_fragment_downstreams,
961 to_drop_state_table_ids,
962 auto_refresh_schema_sinks,
963 ..
964 } = &replace_plan;
965 barrier_manager_context
967 .metadata_manager
968 .catalog_controller
969 .post_collect_job_fragments(
970 new_fragments.stream_job_id,
971 upstream_fragment_downstreams,
972 None,
973 Some(&resolved_split_assignment),
974 None,
975 false,
976 )
977 .await?;
978
979 if let Some(sinks) = auto_refresh_schema_sinks {
980 for sink in sinks {
981 barrier_manager_context
982 .metadata_manager
983 .catalog_controller
984 .post_collect_job_fragments(
985 sink.tmp_sink_id.as_job_id(),
986 &Default::default(), None, None, None,
990 false,
991 )
992 .await?;
993 }
994 }
995
996 barrier_manager_context
998 .source_manager
999 .handle_replace_job(
1000 old_fragments,
1001 new_fragments.stream_source_fragments(),
1002 &replace_plan,
1003 )
1004 .await;
1005 cleanup_dropped_streaming_jobs(
1006 &barrier_manager_context.refresh_manager,
1007 &barrier_manager_context.hummock_manager,
1008 &barrier_manager_context.metadata_manager,
1009 [],
1010 to_drop_state_table_ids.clone(),
1011 "replace_streaming_job",
1012 )
1013 .await?;
1014 }
1015
1016 PostCollectCommand::CreateSubscription { subscription_id } => {
1017 barrier_manager_context
1018 .metadata_manager
1019 .catalog_controller
1020 .finish_create_subscription_catalog(subscription_id)
1021 .await?
1022 }
1023 }
1024
1025 Ok(())
1026 }
1027}
1028
1029#[cfg(test)]
1030mod tests {
1031 use risingwave_hummock_sdk::change_log::{EpochNewChangeLog, TableChangeLog};
1032
1033 use super::*;
1034
1035 fn test_change_logs(table_id: TableId) -> TableChangeLogs {
1036 TableChangeLogs::from_iter([(
1037 table_id,
1038 TableChangeLog::new([
1039 EpochNewChangeLog {
1040 new_value: vec![],
1041 old_value: vec![],
1042 non_checkpoint_epochs: vec![10],
1043 checkpoint_epoch: 20,
1044 },
1045 EpochNewChangeLog {
1046 new_value: vec![],
1047 old_value: vec![],
1048 non_checkpoint_epochs: vec![30],
1049 checkpoint_epoch: 40,
1050 },
1051 EpochNewChangeLog {
1052 new_value: vec![],
1053 old_value: vec![],
1054 non_checkpoint_epochs: vec![50],
1055 checkpoint_epoch: 60,
1056 },
1057 EpochNewChangeLog {
1058 new_value: vec![],
1059 old_value: vec![],
1060 non_checkpoint_epochs: vec![],
1061 checkpoint_epoch: 80,
1062 },
1063 ]),
1064 )])
1065 }
1066
1067 #[test]
1068 fn test_resolve_since_timestamp_log_store_epoch() {
1069 let table_id = TableId::new(233);
1070 let change_logs = test_change_logs(table_id);
1071
1072 let resolved =
1073 resolve_since_timestamp_log_store_epoch(table_id, 45, 80, &change_logs).unwrap();
1074
1075 assert_eq!(resolved, (60, vec![(vec![], 80)]));
1076 }
1077
1078 #[test]
1079 fn test_resolve_since_timestamp_log_store_epoch_rejects_invalid_range() {
1080 let table_id = TableId::new(233);
1081 let change_logs = test_change_logs(table_id);
1082
1083 let err = resolve_since_timestamp_log_store_epoch(table_id, 9, 80, &change_logs)
1084 .unwrap_err()
1085 .to_string();
1086 assert!(err.contains("earlier than the retained changelog"));
1087
1088 let err = resolve_since_timestamp_log_store_epoch(table_id, 80, 80, &change_logs)
1089 .unwrap_err()
1090 .to_string();
1091 assert!(err.contains("is not before the committed epoch"));
1092
1093 let err = resolve_since_timestamp_log_store_epoch(table_id, 45, 60, &change_logs)
1094 .unwrap_err()
1095 .to_string();
1096 assert!(err.contains("does not match latest changelog epoch"));
1097 }
1098
1099 #[test]
1100 fn test_resolve_since_timestamp_log_store_epoch_rejects_empty_or_missing_log() {
1101 let table_id = TableId::new(233);
1102 let empty_change_logs = TableChangeLogs::from_iter([(table_id, TableChangeLog::new([]))]);
1103
1104 let err = resolve_since_timestamp_log_store_epoch(table_id, 30, 60, &empty_change_logs)
1105 .unwrap_err()
1106 .to_string();
1107 assert!(err.contains("empty table changelog"));
1108
1109 let err = resolve_since_timestamp_log_store_epoch(table_id, 30, 60, &Default::default())
1110 .unwrap_err()
1111 .to_string();
1112 assert!(err.contains("no table changelog"));
1113 }
1114
1115 #[test]
1116 fn test_resolve_since_timestamp_log_store_epoch_keeps_latest_non_checkpoint_epochs() {
1117 let table_id = TableId::new(233);
1118 let change_logs = TableChangeLogs::from_iter([(
1119 table_id,
1120 TableChangeLog::new([
1121 EpochNewChangeLog {
1122 new_value: vec![],
1123 old_value: vec![],
1124 non_checkpoint_epochs: vec![],
1125 checkpoint_epoch: 20,
1126 },
1127 EpochNewChangeLog {
1128 new_value: vec![],
1129 old_value: vec![],
1130 non_checkpoint_epochs: vec![30],
1131 checkpoint_epoch: 40,
1132 },
1133 ]),
1134 )]);
1135
1136 let resolved =
1137 resolve_since_timestamp_log_store_epoch(table_id, 20, 40, &change_logs).unwrap();
1138 assert_eq!(resolved, (20, vec![(vec![30], 40)]));
1139
1140 let err = resolve_since_timestamp_log_store_epoch(table_id, 30, 40, &change_logs)
1141 .unwrap_err()
1142 .to_string();
1143 assert!(err.contains("since_timestamp is too new"));
1144 }
1145
1146 #[test]
1147 fn test_skip_refresh_finish_when_associated_source_missing() {
1148 let err = MetaError::catalog_id_not_found("object", 42);
1149 assert!(err.is_catalog_id_not_found("object"));
1150 }
1151
1152 #[test]
1153 fn test_do_not_skip_refresh_finish_for_other_not_found_types() {
1154 let err = MetaError::catalog_id_not_found("table", 42);
1155 assert!(!err.is_catalog_id_not_found("object"));
1156 }
1157}