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