1use std::cmp::{Ordering, max, min};
16use std::collections::hash_map::Entry;
17use std::collections::{HashMap, HashSet};
18use std::sync::atomic::AtomicU32;
19
20use anyhow::{Context, anyhow};
21use futures::StreamExt;
22use futures::stream::FuturesUnordered;
23use itertools::Itertools;
24use risingwave_common::bail;
25use risingwave_common::catalog::{DatabaseId, TableId};
26use risingwave_common::id::JobId;
27use risingwave_common::util::stream_graph_visitor::visit_stream_node_cont;
28use risingwave_connector::source::SplitImpl;
29use risingwave_hummock_sdk::change_log::TableChangeLogs;
30use risingwave_hummock_sdk::version::HummockVersion;
31use risingwave_meta_model::SinkId;
32use risingwave_pb::stream_plan::stream_node::PbNodeBody;
33use sea_orm::TransactionTrait;
34use thiserror_ext::AsReport;
35use tracing::{info, warn};
36
37use super::BarrierWorkerRuntimeInfoSnapshot;
38use crate::MetaResult;
39use crate::barrier::DatabaseRuntimeInfoSnapshot;
40use crate::barrier::checkpoint::{
41 BatchRefreshJobCheckpointControl, BatchRefreshLogicalFragments, BatchRefreshRenderResult,
42};
43use crate::barrier::context::{GlobalBarrierWorkerContext, GlobalBarrierWorkerContextImpl};
44use crate::barrier::progress::TrackingJob;
45use crate::barrier::rpc::to_partial_graph_id;
46use crate::controller::fragment::{InflightActorInfo, InflightFragmentInfo};
47use crate::controller::scale::{
48 FragmentRenderMap, LoadedFragment, LoadedFragmentContext, RenderedGraph,
49 render_actor_assignments,
50};
51use crate::controller::utils::StreamingJobExtraInfo;
52use crate::manager::ActiveStreamingWorkerNodes;
53use crate::model::{ActorId, FragmentDownstreamRelation, FragmentId, StreamActor};
54use crate::rpc::ddl_controller::refill_upstream_sink_union_in_table;
55use crate::stream::cdc::reload_cdc_table_snapshot_splits;
56use crate::stream::{
57 SourceChange, StreamFragmentGraph, UpstreamSinkInfo, cleanup_dropped_streaming_jobs,
58};
59
60#[derive(Debug)]
61pub(crate) struct UpstreamSinkRecoveryInfo {
62 target_fragment_id: FragmentId,
63 upstream_infos: Vec<UpstreamSinkInfo>,
64}
65
66#[derive(Debug)]
67pub struct LoadedRecoveryContext {
68 pub fragment_context: LoadedFragmentContext,
69 pub job_extra_info: HashMap<JobId, StreamingJobExtraInfo>,
70 pub upstream_sink_recovery: HashMap<JobId, UpstreamSinkRecoveryInfo>,
71 pub fragment_relations: FragmentDownstreamRelation,
72}
73
74impl LoadedRecoveryContext {
75 fn empty(fragment_context: LoadedFragmentContext) -> Self {
76 Self {
77 fragment_context,
78 job_extra_info: HashMap::new(),
79 upstream_sink_recovery: HashMap::new(),
80 fragment_relations: FragmentDownstreamRelation::default(),
81 }
82 }
83}
84
85pub struct RenderedDatabaseRuntimeInfo {
86 pub job_infos: HashMap<JobId, HashMap<FragmentId, InflightFragmentInfo>>,
87 pub stream_actors: HashMap<ActorId, StreamActor>,
88 pub source_splits: HashMap<ActorId, Vec<SplitImpl>>,
89 pub batch_refresh: HashMap<JobId, BatchRefreshRenderResult>,
91}
92
93pub fn render_runtime_info(
94 actor_id_generator: &AtomicU32,
95 worker_nodes: &ActiveStreamingWorkerNodes,
96 recovery_context: &LoadedRecoveryContext,
97 database_id: DatabaseId,
98) -> MetaResult<Option<RenderedDatabaseRuntimeInfo>> {
99 let Some(mut per_database_context) =
100 recovery_context.fragment_context.for_database(database_id)
101 else {
102 return Ok(None);
103 };
104
105 assert!(!per_database_context.is_empty());
106
107 let batch_refresh_job_ids: HashSet<JobId> = per_database_context
110 .job_map
111 .iter()
112 .filter(|(_, model)| model.refresh_interval_sec.is_some())
113 .map(|(job_id, _)| *job_id)
114 .collect();
115
116 let mut batch_refresh_logical = HashMap::new();
117 if !batch_refresh_job_ids.is_empty() {
118 let batch_refresh_fragment_ids: HashSet<FragmentId> = batch_refresh_job_ids
119 .iter()
120 .flat_map(|job_id| {
121 per_database_context
122 .job_fragments
123 .get(job_id)
124 .unwrap()
125 .keys()
126 .copied()
127 })
128 .collect();
129
130 for &job_id in &batch_refresh_job_ids {
131 let fragments = per_database_context.job_fragments.remove(&job_id).unwrap();
132 let downstreams = fragments
133 .keys()
134 .filter_map(|fid| {
135 recovery_context
136 .fragment_relations
137 .get(fid)
138 .map(|r| (*fid, r.clone()))
139 })
140 .collect();
141 batch_refresh_logical.insert(
142 job_id,
143 BatchRefreshLogicalFragments {
144 fragments,
145 downstreams,
146 },
147 );
148 }
149
150 per_database_context.ensembles.retain(|ensemble| {
152 !ensemble
153 .component_fragments()
154 .all(|fid| batch_refresh_fragment_ids.contains(&fid))
155 });
156 }
157
158 let mut batch_refresh = HashMap::new();
160 for (job_id, logical) in batch_refresh_logical {
161 let extra = recovery_context
162 .job_extra_info
163 .get(&job_id)
164 .expect("should have extra info");
165 let streaming_job_model = per_database_context
166 .job_map
167 .get(&job_id)
168 .expect("should have streaming job model");
169 let database_model = &per_database_context.database_map[&database_id];
170 let partial_graph_id = to_partial_graph_id(database_id, Some(job_id));
171
172 let render_result = BatchRefreshJobCheckpointControl::render_actors_and_build_job_info(
173 &logical.fragments,
174 &logical.downstreams,
175 &extra.job_definition,
176 actor_id_generator,
177 worker_nodes.current(),
178 &database_model.resource_group,
179 streaming_job_model,
180 partial_graph_id,
181 )?;
182
183 batch_refresh.insert(job_id, render_result);
184 }
185
186 if per_database_context.ensembles.is_empty() {
188 return Ok(Some(RenderedDatabaseRuntimeInfo {
189 job_infos: HashMap::new(),
190 stream_actors: HashMap::new(),
191 source_splits: HashMap::new(),
192 batch_refresh,
193 }));
194 }
195
196 let RenderedGraph { mut fragments, .. } = render_actor_assignments(
197 actor_id_generator,
198 worker_nodes.current(),
199 &per_database_context,
200 )?;
201
202 let single_database = match fragments.remove(&database_id) {
203 Some(info) => info,
204 None => return Ok(None),
205 };
206
207 let mut database_map = HashMap::from([(database_id, single_database)]);
208 recovery_table_with_upstream_sinks(
209 &mut database_map,
210 &recovery_context.upstream_sink_recovery,
211 )?;
212 let stream_actors = build_stream_actors(&database_map, &recovery_context.job_extra_info)?;
213
214 let job_infos = database_map
215 .remove(&database_id)
216 .expect("database entry must exist");
217
218 let mut source_splits = HashMap::new();
219 for fragment_infos in job_infos.values() {
220 for fragment in fragment_infos.values() {
221 for (actor_id, info) in &fragment.actors {
222 source_splits.insert(*actor_id, info.splits.clone());
223 }
224 }
225 }
226
227 Ok(Some(RenderedDatabaseRuntimeInfo {
228 job_infos,
229 stream_actors,
230 source_splits,
231 batch_refresh,
232 }))
233}
234
235fn recovery_table_with_upstream_sinks(
240 inflight_jobs: &mut FragmentRenderMap,
241 upstream_sink_recovery: &HashMap<JobId, UpstreamSinkRecoveryInfo>,
242) -> MetaResult<()> {
243 if upstream_sink_recovery.is_empty() {
244 return Ok(());
245 }
246
247 let mut seen_jobs = HashSet::new();
248
249 for jobs in inflight_jobs.values_mut() {
250 for (job_id, fragments) in jobs {
251 if !seen_jobs.insert(*job_id) {
252 return Err(anyhow::anyhow!("Duplicate job id found: {}", job_id).into());
253 }
254
255 if let Some(recovery) = upstream_sink_recovery.get(job_id) {
256 if let Some(target_fragment) = fragments.get_mut(&recovery.target_fragment_id) {
257 refill_upstream_sink_union_in_table(
258 &mut target_fragment.nodes,
259 &recovery.upstream_infos,
260 );
261 } else {
262 return Err(anyhow::anyhow!(
263 "target fragment {} not found for upstream sink recovery of job {}",
264 recovery.target_fragment_id,
265 job_id
266 )
267 .into());
268 }
269 }
270 }
271 }
272
273 Ok(())
274}
275
276fn build_stream_actors(
282 all_info: &FragmentRenderMap,
283 job_extra_info: &HashMap<JobId, StreamingJobExtraInfo>,
284) -> MetaResult<HashMap<ActorId, StreamActor>> {
285 let mut stream_actors = HashMap::new();
286
287 for (job_id, streaming_info) in all_info.values().flatten() {
288 let extra_info = job_extra_info
289 .get(job_id)
290 .cloned()
291 .ok_or_else(|| anyhow!("no streaming job info for {}", job_id))?;
292 let expr_context = extra_info.stream_context().to_expr_context();
293 let job_definition = extra_info.job_definition;
294 let config_override = extra_info.config_override;
295
296 for (fragment_id, fragment_infos) in streaming_info {
297 for (actor_id, InflightActorInfo { vnode_bitmap, .. }) in &fragment_infos.actors {
298 stream_actors.insert(
299 *actor_id,
300 StreamActor {
301 actor_id: *actor_id,
302 fragment_id: *fragment_id,
303 vnode_bitmap: vnode_bitmap.clone(),
304 mview_definition: job_definition.clone(),
305 expr_context: Some(expr_context.clone()),
306 config_override: config_override.clone(),
307 },
308 );
309 }
310 }
311 }
312 Ok(stream_actors)
313}
314
315impl GlobalBarrierWorkerContextImpl {
316 fn resolve_job_committed_epoch(
317 job_id: JobId,
318 fragments: &HashMap<FragmentId, LoadedFragment>,
319 state_table_committed_epochs: &HashMap<TableId, u64>,
320 ) -> MetaResult<u64> {
321 let mut table_id_iter = fragments
322 .values()
323 .flat_map(|fragment| fragment.state_table_ids.iter().copied());
324 let Some(first_table_id) = table_id_iter.next() else {
325 bail!("job {} has no state table", job_id);
326 };
327 let committed_epoch = *state_table_committed_epochs
328 .get(&first_table_id)
329 .ok_or_else(|| anyhow!("cannot get committed epoch on table {}.", first_table_id))?;
330 for table_id in table_id_iter {
331 let table_committed_epoch = *state_table_committed_epochs
332 .get(&table_id)
333 .ok_or_else(|| anyhow!("cannot get committed epoch on table {}.", table_id))?;
334 if committed_epoch != table_committed_epoch {
335 bail!(
336 "table {} has committed epoch {} different to other table {} with committed epoch {} in job {}",
337 first_table_id,
338 committed_epoch,
339 table_id,
340 table_committed_epoch,
341 job_id
342 );
343 }
344 }
345
346 Ok(committed_epoch)
347 }
348
349 async fn finish_completed_batch_refresh_background_jobs(
350 &self,
351 recovery_context: &LoadedRecoveryContext,
352 state_table_committed_epochs: &HashMap<TableId, u64>,
353 creating_jobs: &mut HashSet<JobId>,
354 ) -> MetaResult<()> {
355 let creating_job_ids = creating_jobs.iter().copied().collect_vec();
356 for job_id in creating_job_ids {
357 let Some(job) = recovery_context.fragment_context.job_map.get(&job_id) else {
358 continue;
359 };
360 if job.refresh_interval_sec.is_none() {
361 continue;
362 }
363 let Some(fragments) = recovery_context.fragment_context.job_fragments.get(&job_id)
364 else {
365 continue;
366 };
367
368 let committed_epoch =
369 Self::resolve_job_committed_epoch(job_id, fragments, state_table_committed_epochs)?;
370 let snapshot_backfill_info = StreamFragmentGraph::collect_snapshot_backfill_info_impl(
371 fragments
372 .values()
373 .map(|fragment| (&fragment.nodes, fragment.fragment_type_mask)),
374 )?
375 .0
376 .ok_or_else(|| anyhow!("batch refresh job {} has no snapshot backfill info", job_id))?;
377 let snapshot_epoch = snapshot_backfill_info
378 .upstream_mv_table_id_to_backfill_epoch
379 .values()
380 .find_map(|e| *e)
381 .unwrap_or(committed_epoch);
382 if committed_epoch < snapshot_epoch {
383 continue;
384 }
385
386 info!(
387 %job_id,
388 committed_epoch,
389 snapshot_epoch,
390 "finish completed batch refresh background job during recovery"
391 );
392 self.finish_creating_job(TrackingJob::recovered_from_fragment_nodes(
393 job_id,
394 fragments
395 .iter()
396 .map(|(fragment_id, fragment)| (*fragment_id, &fragment.nodes)),
397 ))
398 .await?;
399 creating_jobs.remove(&job_id);
400 }
401
402 Ok(())
403 }
404
405 async fn apply_pre_applied_drop_cancel(
406 &self,
407 database_id: Option<DatabaseId>,
408 ) -> MetaResult<bool> {
409 let drop_cancel = self.scheduled_barriers.pre_apply_drop_cancel(database_id);
410 let has_drop_streaming_jobs = !drop_cancel.streaming_job_ids.is_empty();
411 cleanup_dropped_streaming_jobs(
412 &self.refresh_manager,
413 &self.hummock_manager,
414 &self.metadata_manager,
415 drop_cancel.streaming_job_ids,
416 drop_cancel.dropped_state_table_ids,
417 "drop_streaming_jobs",
418 )
419 .await?;
420 Ok(has_drop_streaming_jobs)
421 }
422
423 async fn clean_dirty_streaming_jobs(&self, database_id: Option<DatabaseId>) -> MetaResult<()> {
426 self.metadata_manager
427 .catalog_controller
428 .clean_dirty_subscription(database_id)
429 .await?;
430
431 let cleaned_dirty_jobs = self
432 .metadata_manager
433 .catalog_controller
434 .clean_dirty_creating_jobs(database_id)
435 .await?;
436 for sink_id in &cleaned_dirty_jobs.sink_ids {
437 self.iceberg_compaction_manager
438 .clear_iceberg_maintenance_by_sink_id(*sink_id);
439 }
440 if database_id.is_some() {
441 cleanup_dropped_streaming_jobs(
444 &self.refresh_manager,
445 &self.hummock_manager,
446 &self.metadata_manager,
447 cleaned_dirty_jobs.streaming_job_ids,
448 cleaned_dirty_jobs.dropped_table_ids,
449 "clean_dirty_creating_jobs",
450 )
451 .await?;
452 }
453 self.metadata_manager
454 .reset_all_refresh_jobs_to_idle()
455 .await?;
456
457 self.source_manager
459 .apply_source_change(SourceChange::DropSource {
460 dropped_source_ids: cleaned_dirty_jobs.source_ids,
461 })
462 .await;
463
464 Ok(())
465 }
466
467 async fn reset_sink_coordinator(&self, database_id: Option<DatabaseId>) -> MetaResult<()> {
468 if let Some(database_id) = database_id {
469 let sink_ids = self
470 .metadata_manager
471 .catalog_controller
472 .list_sink_ids(Some(database_id))
473 .await?;
474 self.sink_manager
475 .stop_sink_coordinator(sink_ids.clone())
476 .await;
477 self.iceberg_pk_index_sink_manager
478 .unregister_sinks(sink_ids);
479 } else {
480 self.sink_manager.reset().await;
481 self.iceberg_pk_index_sink_manager.reset();
482 }
483 Ok(())
484 }
485
486 async fn reregister_iceberg_pk_index_sinks(
488 &self,
489 database_id: Option<DatabaseId>,
490 ) -> MetaResult<()> {
491 let pb_sinks = self
492 .metadata_manager
493 .catalog_controller
494 .list_sinks()
495 .await?;
496 let mut futs = FuturesUnordered::new();
497 for pb_sink in pb_sinks {
498 if database_id.is_some_and(|db_id| pb_sink.database_id != db_id)
499 || !crate::manager::iceberg_pk_index_sink::is_iceberg_pk_index_sink(
500 &pb_sink.properties,
501 )
502 {
503 continue;
504 }
505 let config = crate::manager::iceberg_pk_index_sink::build_iceberg_config(&pb_sink)
506 .with_context(|| {
507 format!(
508 "build iceberg config while re-registering v3 sink {}",
509 pb_sink.id
510 )
511 })?;
512 let state_table_id = self
513 .metadata_manager
514 .catalog_controller
515 .get_sink_state_table_ids(pb_sink.id)
516 .await?
517 .into_iter()
518 .next()
519 .ok_or_else(|| {
520 anyhow!(
521 "no state table found while re-registering iceberg v3 sink {}",
522 pb_sink.id
523 )
524 })?;
525 let recovered_epoch = self
526 .hummock_manager
527 .on_current_version(|version| version.table_committed_epoch(state_table_id))
528 .await
529 .ok_or_else(|| {
530 anyhow!(
531 "cannot get committed epoch for iceberg v3 sink {} state table {}",
532 pb_sink.id,
533 state_table_id
534 )
535 })?;
536 let manager = &self.iceberg_pk_index_sink_manager;
537 futs.push(async move {
538 let partial_graph_id = to_partial_graph_id(pb_sink.database_id, None);
539 let result = manager
540 .register_sink(pb_sink.id, partial_graph_id, config)
541 .await;
542 if result.is_ok() {
543 manager.advance_committed_epochs([(partial_graph_id, recovered_epoch)]);
544 }
545 (pb_sink.id, result)
546 });
547 }
548
549 while let Some((id, res)) = futs.next().await {
550 if let Err(e) = res {
551 let msg = format!("register iceberg v3 sink {} during recovery", id);
552 return Err(e.context(msg).into());
553 }
554 }
555 Ok(())
556 }
557
558 async fn abort_dirty_pending_sink_state(
559 &self,
560 database_id: Option<DatabaseId>,
561 ) -> MetaResult<()> {
562 let pending_sinks: HashSet<SinkId> = self
563 .metadata_manager
564 .catalog_controller
565 .list_all_pending_sinks(database_id)
566 .await?;
567
568 if pending_sinks.is_empty() {
569 return Ok(());
570 }
571
572 let sink_with_state_tables: HashMap<SinkId, Vec<TableId>> = self
573 .metadata_manager
574 .catalog_controller
575 .fetch_sink_with_state_table_ids(pending_sinks)
576 .await?;
577
578 let mut sink_committed_epoch: HashMap<SinkId, u64> = HashMap::new();
579
580 for (sink_id, table_ids) in sink_with_state_tables {
581 let Some(table_id) = table_ids.first() else {
582 return Err(anyhow!("no state table id in sink: {}", sink_id).into());
583 };
584
585 self.hummock_manager
586 .on_current_version(|version| -> MetaResult<()> {
587 if let Some(committed_epoch) = version.table_committed_epoch(*table_id) {
588 assert!(
589 sink_committed_epoch
590 .insert(sink_id, committed_epoch)
591 .is_none()
592 );
593 Ok(())
594 } else {
595 Err(anyhow!("cannot get committed epoch on table {}.", table_id).into())
596 }
597 })
598 .await?;
599 }
600
601 self.metadata_manager
602 .catalog_controller
603 .abort_pending_sink_epochs(sink_committed_epoch)
604 .await?;
605
606 Ok(())
607 }
608
609 async fn purge_state_table_from_hummock(
610 &self,
611 all_state_table_ids: &HashSet<TableId>,
612 ) -> MetaResult<()> {
613 self.hummock_manager.purge(all_state_table_ids).await?;
614 Ok(())
615 }
616
617 async fn list_creating_jobs(
618 &self,
619 database_id: Option<DatabaseId>,
620 ) -> MetaResult<HashSet<JobId>> {
621 let mgr = &self.metadata_manager;
622 Ok(mgr
623 .catalog_controller
624 .list_creating_jobs(false, database_id)
625 .await?
626 .into_iter()
627 .map(|(job_id, _, _, _, _)| job_id)
628 .collect())
629 }
630
631 async fn load_recovery_context(
632 &self,
633 database_id: Option<DatabaseId>,
634 ) -> MetaResult<LoadedRecoveryContext> {
635 let inner = self
636 .metadata_manager
637 .catalog_controller
638 .get_inner_read_guard()
639 .await;
640 let txn = inner.db.begin().await?;
641
642 let fragment_context = self
643 .metadata_manager
644 .catalog_controller
645 .load_fragment_context_in_txn(&txn, database_id)
646 .await
647 .inspect_err(|err| {
648 warn!(error = %err.as_report(), "load fragment context failed");
649 })?;
650
651 if fragment_context.is_empty() {
652 return Ok(LoadedRecoveryContext::empty(fragment_context));
653 }
654
655 let job_ids = fragment_context.job_map.keys().copied().collect_vec();
656 let job_extra_info = self
657 .metadata_manager
658 .catalog_controller
659 .get_streaming_job_extra_info_in_txn(&txn, job_ids)
660 .await?;
661
662 let mut upstream_targets = HashMap::new();
663 for fragment in fragment_context
664 .job_fragments
665 .values()
666 .flat_map(|fragments| fragments.values())
667 {
668 let mut has_upstream_union = false;
669 visit_stream_node_cont(&fragment.nodes, |node| {
670 if let Some(PbNodeBody::UpstreamSinkUnion(_)) = node.node_body {
671 has_upstream_union = true;
672 false
673 } else {
674 true
675 }
676 });
677
678 if has_upstream_union
679 && let Some(previous) =
680 upstream_targets.insert(fragment.job_id, fragment.fragment_id)
681 {
682 bail!(
683 "multiple upstream sink union fragments found for job {}, fragment {}, kept {}",
684 fragment.job_id,
685 fragment.fragment_id,
686 previous
687 );
688 }
689 }
690
691 let mut upstream_sink_recovery = HashMap::new();
692 if !upstream_targets.is_empty() {
693 let tables = self
694 .metadata_manager
695 .catalog_controller
696 .get_user_created_table_by_ids_in_txn(&txn, upstream_targets.keys().copied())
697 .await?;
698
699 for table in tables {
700 let job_id = table.id.as_job_id();
701 let Some(target_fragment_id) = upstream_targets.get(&job_id) else {
702 tracing::debug!(
704 job_id = %job_id,
705 "upstream sink union target fragment not found for table"
706 );
707 continue;
708 };
709
710 let upstream_infos = self
711 .metadata_manager
712 .catalog_controller
713 .get_all_upstream_sink_infos_in_txn(&txn, &table, *target_fragment_id as _)
714 .await?;
715
716 upstream_sink_recovery.insert(
717 job_id,
718 UpstreamSinkRecoveryInfo {
719 target_fragment_id: *target_fragment_id,
720 upstream_infos,
721 },
722 );
723 }
724 }
725
726 let fragment_relations = self
727 .metadata_manager
728 .catalog_controller
729 .get_fragment_downstream_relations_in_txn(
730 &txn,
731 fragment_context
732 .job_fragments
733 .values()
734 .flat_map(|fragments| fragments.keys().copied())
735 .collect_vec(),
736 )
737 .await?;
738
739 Ok(LoadedRecoveryContext {
740 fragment_context,
741 job_extra_info,
742 upstream_sink_recovery,
743 fragment_relations,
744 })
745 }
746
747 #[expect(clippy::type_complexity)]
748 fn resolve_hummock_version_epochs(
749 creating_jobs: impl Iterator<Item = (JobId, &HashMap<FragmentId, LoadedFragment>)>,
750 version: &HummockVersion,
751 table_change_log: &TableChangeLogs,
752 ) -> MetaResult<(
753 HashMap<TableId, u64>,
754 HashMap<TableId, Vec<(Vec<u64>, u64)>>,
755 )> {
756 let table_committed_epoch: HashMap<_, _> = version
757 .state_table_info
758 .info()
759 .iter()
760 .map(|(table_id, info)| (*table_id, info.committed_epoch))
761 .collect();
762 let get_table_committed_epoch = |table_id| -> anyhow::Result<u64> {
763 Ok(*table_committed_epoch
764 .get(&table_id)
765 .ok_or_else(|| anyhow!("cannot get committed epoch on table {}.", table_id))?)
766 };
767 let mut min_downstream_committed_epochs = HashMap::new();
768 for (job_id, fragments) in creating_jobs {
769 let job_committed_epoch =
770 Self::resolve_job_committed_epoch(job_id, fragments, &table_committed_epoch)?;
771 if let (Some(snapshot_backfill_info), _) =
772 StreamFragmentGraph::collect_snapshot_backfill_info_impl(
773 fragments
774 .values()
775 .map(|fragment| (&fragment.nodes, fragment.fragment_type_mask)),
776 )?
777 {
778 for (upstream_table, snapshot_epoch) in
779 snapshot_backfill_info.upstream_mv_table_id_to_backfill_epoch
780 {
781 let snapshot_epoch = snapshot_epoch.ok_or_else(|| {
782 anyhow!(
783 "recovered snapshot backfill job {} has not filled snapshot epoch to upstream {}",
784 job_id, upstream_table
785 )
786 })?;
787 let pinned_epoch = max(snapshot_epoch, job_committed_epoch);
788 match min_downstream_committed_epochs.entry(upstream_table) {
789 Entry::Occupied(entry) => {
790 let prev_min_epoch = entry.into_mut();
791 *prev_min_epoch = min(*prev_min_epoch, pinned_epoch);
792 }
793 Entry::Vacant(entry) => {
794 entry.insert(pinned_epoch);
795 }
796 }
797 }
798 }
799 }
800 let mut log_epochs = HashMap::new();
801 for (upstream_table_id, downstream_committed_epoch) in min_downstream_committed_epochs {
802 let upstream_committed_epoch = get_table_committed_epoch(upstream_table_id)?;
803 match upstream_committed_epoch.cmp(&downstream_committed_epoch) {
804 Ordering::Less => {
805 bail!(
806 "downstream epoch {} later than upstream epoch {} of table {}",
807 downstream_committed_epoch,
808 upstream_committed_epoch,
809 upstream_table_id
810 );
811 }
812 Ordering::Equal => {
813 continue;
814 }
815 Ordering::Greater => {
816 if let Some(table_change_log) = table_change_log.get(&upstream_table_id) {
817 let epochs = table_change_log
818 .filter_epoch((downstream_committed_epoch, upstream_committed_epoch))
819 .map(|epoch_log| {
820 (
821 epoch_log.non_checkpoint_epochs.clone(),
822 epoch_log.checkpoint_epoch,
823 )
824 })
825 .collect_vec();
826 let first_epochs = epochs.first();
827 if let Some((_, first_checkpoint_epoch)) = &first_epochs
828 && *first_checkpoint_epoch == downstream_committed_epoch
829 {
830 } else {
831 bail!(
832 "resolved first log epoch {:?} on table {} not matched with downstream committed epoch {}",
833 epochs,
834 upstream_table_id,
835 downstream_committed_epoch
836 );
837 }
838 log_epochs
839 .try_insert(upstream_table_id, epochs)
840 .expect("non-duplicated");
841 } else {
842 bail!(
843 "upstream table {} on epoch {} has lagged downstream on epoch {} but no table change log",
844 upstream_table_id,
845 upstream_committed_epoch,
846 downstream_committed_epoch
847 );
848 }
849 }
850 }
851 }
852 Ok((table_committed_epoch, log_epochs))
853 }
854
855 pub(super) async fn reload_runtime_info_impl(
856 &self,
857 ) -> MetaResult<BarrierWorkerRuntimeInfoSnapshot> {
858 {
859 {
860 {
861 self.clean_dirty_streaming_jobs(None)
862 .await
863 .context("clean dirty streaming jobs")?;
864
865 self.reset_sink_coordinator(None)
866 .await
867 .context("reset sink coordinator")?;
868 self.abort_dirty_pending_sink_state(None)
869 .await
870 .context("abort dirty pending sink state")?;
871
872 self.reregister_iceberg_pk_index_sinks(None)
876 .await
877 .context("re-register iceberg v3 sinks after recovery")?;
878
879 tracing::info!("recovering creating job progress");
881 let mut initial_creating_jobs = self
882 .list_creating_jobs(None)
883 .await
884 .context("recover creating job progress should not fail")?;
885
886 tracing::info!("recovered creating job progress");
887
888 let _ = self.apply_pre_applied_drop_cancel(None).await?;
890 self.metadata_manager
891 .catalog_controller
892 .cleanup_dropped_tables()
893 .await;
894
895 let active_streaming_nodes =
896 ActiveStreamingWorkerNodes::new_snapshot(self.metadata_manager.clone())
897 .await?;
898
899 let creating_streaming_jobs =
900 initial_creating_jobs.iter().cloned().collect_vec();
901
902 tracing::info!(
903 "creating streaming jobs: {:?} total {}",
904 creating_streaming_jobs,
905 creating_streaming_jobs.len()
906 );
907
908 let unreschedulable_jobs = {
909 let mut unreschedulable_jobs = HashSet::new();
910
911 for job_id in creating_streaming_jobs {
912 let scan_types = self
913 .metadata_manager
914 .get_job_backfill_scan_types(job_id)
915 .await?;
916
917 if scan_types
918 .values()
919 .any(|scan_type| !scan_type.is_reschedulable(false))
920 {
921 unreschedulable_jobs.insert(job_id);
922 }
923 }
924
925 unreschedulable_jobs
926 };
927
928 if !unreschedulable_jobs.is_empty() {
929 info!("unreschedulable creating jobs: {:?}", unreschedulable_jobs);
930 }
931
932 if !unreschedulable_jobs.is_empty() {
936 bail!(
937 "Recovery for unreschedulable creating jobs is not yet implemented. \
938 This path is triggered when the following jobs have at least one scan type that is not reschedulable: {:?}.",
939 unreschedulable_jobs
940 );
941 }
942
943 let mut recovery_context = self.load_recovery_context(None).await?;
944 if self.apply_pre_applied_drop_cancel(None).await? {
945 recovery_context = self.load_recovery_context(None).await?;
946 }
947
948 self.purge_state_table_from_hummock(
949 &recovery_context
950 .fragment_context
951 .job_fragments
952 .values()
953 .flat_map(|fragments| fragments.values())
954 .flat_map(|fragment| fragment.state_table_ids.iter().copied())
955 .collect(),
956 )
957 .await
958 .context("purge state table from hummock")?;
959
960 let (state_table_committed_epochs, state_table_log_epochs) = self
961 .hummock_manager
962 .on_current_version_and_table_change_log(|version, table_change_log| {
963 Self::resolve_hummock_version_epochs(
964 recovery_context
965 .fragment_context
966 .job_fragments
967 .iter()
968 .filter_map(|(job_id, job)| {
969 initial_creating_jobs
970 .contains(job_id)
971 .then_some((*job_id, job))
972 }),
973 version,
974 table_change_log,
975 )
976 })
977 .await?;
978
979 self.finish_completed_batch_refresh_background_jobs(
980 &recovery_context,
981 &state_table_committed_epochs,
982 &mut initial_creating_jobs,
983 )
984 .await?;
985
986 let mv_depended_subscriptions = self
987 .metadata_manager
988 .get_mv_depended_subscriptions(None)
989 .await?;
990
991 let creating_jobs = {
993 let mut refreshed_creating_jobs = self
994 .list_creating_jobs(None)
995 .await
996 .context("recover creating job progress should not fail")?;
997 recovery_context
998 .fragment_context
999 .job_map
1000 .keys()
1001 .filter_map(|job_id| {
1002 refreshed_creating_jobs.remove(job_id).then_some(*job_id)
1003 })
1004 .collect()
1005 };
1006
1007 let database_infos = self
1008 .metadata_manager
1009 .catalog_controller
1010 .list_databases()
1011 .await?;
1012
1013 let cdc_table_snapshot_splits =
1014 reload_cdc_table_snapshot_splits(&self.env.meta_store_ref().conn, None)
1015 .await?;
1016
1017 Ok(BarrierWorkerRuntimeInfoSnapshot {
1018 active_streaming_nodes,
1019 recovery_context,
1020 state_table_committed_epochs,
1021 state_table_log_epochs,
1022 mv_depended_subscriptions,
1023 creating_jobs,
1024 hummock_version_stats: self.hummock_manager.get_version_stats().await,
1025 database_infos,
1026 cdc_table_snapshot_splits,
1027 })
1028 }
1029 }
1030 }
1031 }
1032
1033 pub(super) async fn reload_database_runtime_info_impl(
1034 &self,
1035 database_id: DatabaseId,
1036 ) -> MetaResult<DatabaseRuntimeInfoSnapshot> {
1037 self.clean_dirty_streaming_jobs(Some(database_id))
1038 .await
1039 .context("clean dirty streaming jobs")?;
1040
1041 self.reset_sink_coordinator(Some(database_id))
1042 .await
1043 .context("reset sink coordinator")?;
1044 self.abort_dirty_pending_sink_state(Some(database_id))
1045 .await
1046 .context("abort dirty pending sink state")?;
1047 self.reregister_iceberg_pk_index_sinks(Some(database_id))
1048 .await
1049 .context("re-register iceberg v3 sinks after recovery")?;
1050
1051 tracing::info!(?database_id, "recovering creating job progress of database");
1053
1054 let mut creating_jobs = self
1055 .list_creating_jobs(Some(database_id))
1056 .await
1057 .context("recover creating job progress of database should not fail")?;
1058 tracing::info!(?database_id, "recovered creating job progress");
1059
1060 let _ = self
1062 .apply_pre_applied_drop_cancel(Some(database_id))
1063 .await?;
1064
1065 let recovery_context = self.load_recovery_context(Some(database_id)).await?;
1066
1067 let missing_creating_jobs = creating_jobs
1068 .iter()
1069 .filter(|job_id| {
1070 !recovery_context
1071 .fragment_context
1072 .job_map
1073 .contains_key(*job_id)
1074 })
1075 .copied()
1076 .collect_vec();
1077 if !missing_creating_jobs.is_empty() {
1078 warn!(
1079 database_id = %database_id,
1080 missing_job_ids = ?missing_creating_jobs,
1081 "creating jobs missing in rendered info"
1082 );
1083 }
1084
1085 let (state_table_committed_epochs, state_table_log_epochs) = self
1086 .hummock_manager
1087 .on_current_version_and_table_change_log(|version, table_change_log| {
1088 Self::resolve_hummock_version_epochs(
1089 creating_jobs.iter().filter_map(|job_id| {
1090 recovery_context
1091 .fragment_context
1092 .job_fragments
1093 .get(job_id)
1094 .map(|job| (*job_id, job))
1095 }),
1096 version,
1097 table_change_log,
1098 )
1099 })
1100 .await?;
1101
1102 self.finish_completed_batch_refresh_background_jobs(
1103 &recovery_context,
1104 &state_table_committed_epochs,
1105 &mut creating_jobs,
1106 )
1107 .await?;
1108
1109 let mv_depended_subscriptions = self
1110 .metadata_manager
1111 .get_mv_depended_subscriptions(Some(database_id))
1112 .await?;
1113
1114 let cdc_table_snapshot_splits =
1115 reload_cdc_table_snapshot_splits(&self.env.meta_store_ref().conn, Some(database_id))
1116 .await?;
1117
1118 self.refresh_manager
1119 .remove_trackers_by_database(database_id);
1120
1121 Ok(DatabaseRuntimeInfoSnapshot {
1122 recovery_context,
1123 state_table_committed_epochs,
1124 state_table_log_epochs,
1125 mv_depended_subscriptions,
1126 creating_jobs,
1127 cdc_table_snapshot_splits,
1128 })
1129 }
1130}
1131
1132#[cfg(test)]
1133mod tests {
1134 use std::collections::HashMap;
1135
1136 use risingwave_common::catalog::FragmentTypeMask;
1137 use risingwave_common::id::WorkerId;
1138 use risingwave_meta_model::DispatcherType;
1139 use risingwave_meta_model::fragment::DistributionType;
1140 use risingwave_pb::stream_plan::stream_node::PbNodeBody;
1141 use risingwave_pb::stream_plan::{
1142 PbDispatchOutputMapping, PbStreamNode, UpstreamSinkUnionNode as PbUpstreamSinkUnionNode,
1143 };
1144
1145 use super::*;
1146 use crate::controller::fragment::InflightActorInfo;
1147 use crate::model::DownstreamFragmentRelation;
1148 use crate::stream::UpstreamSinkInfo;
1149
1150 #[test]
1151 fn test_recovery_table_with_upstream_sinks_updates_union_node() {
1152 let database_id = DatabaseId::new(1);
1153 let job_id = JobId::new(10);
1154 let fragment_id = FragmentId::new(100);
1155 let sink_fragment_id = FragmentId::new(200);
1156
1157 let mut inflight_jobs: FragmentRenderMap = HashMap::new();
1158 let fragment = InflightFragmentInfo {
1159 fragment_id,
1160 distribution_type: DistributionType::Hash,
1161 fragment_type_mask: FragmentTypeMask::empty(),
1162 vnode_count: 1,
1163 nodes: PbStreamNode {
1164 node_body: Some(PbNodeBody::UpstreamSinkUnion(Box::new(
1165 PbUpstreamSinkUnionNode {
1166 init_upstreams: vec![],
1167 },
1168 ))),
1169 ..Default::default()
1170 },
1171 actors: HashMap::new(),
1172 state_table_ids: HashSet::new(),
1173 };
1174
1175 inflight_jobs
1176 .entry(database_id)
1177 .or_default()
1178 .entry(job_id)
1179 .or_default()
1180 .insert(fragment_id, fragment);
1181
1182 let upstream_sink_recovery = HashMap::from([(
1183 job_id,
1184 UpstreamSinkRecoveryInfo {
1185 target_fragment_id: fragment_id,
1186 upstream_infos: vec![UpstreamSinkInfo {
1187 sink_id: SinkId::new(1),
1188 sink_fragment_id,
1189 sink_output_fields: vec![],
1190 sink_original_target_columns: vec![],
1191 project_exprs: vec![],
1192 new_sink_downstream: DownstreamFragmentRelation {
1193 downstream_fragment_id: FragmentId::new(300),
1194 dispatcher_type: DispatcherType::Hash,
1195 dist_key_indices: vec![],
1196 output_mapping: PbDispatchOutputMapping::default(),
1197 },
1198 }],
1199 },
1200 )]);
1201
1202 recovery_table_with_upstream_sinks(&mut inflight_jobs, &upstream_sink_recovery).unwrap();
1203
1204 let updated = inflight_jobs
1205 .get(&database_id)
1206 .unwrap()
1207 .get(&job_id)
1208 .unwrap()
1209 .get(&fragment_id)
1210 .unwrap();
1211
1212 let PbNodeBody::UpstreamSinkUnion(updated_union) =
1213 updated.nodes.node_body.as_ref().unwrap()
1214 else {
1215 panic!("expected upstream sink union node");
1216 };
1217
1218 assert_eq!(updated_union.init_upstreams.len(), 1);
1219 assert_eq!(
1220 updated_union.init_upstreams[0].upstream_fragment_id,
1221 sink_fragment_id.as_raw_id()
1222 );
1223 }
1224
1225 #[test]
1226 fn test_build_stream_actors_uses_preloaded_extra_info() {
1227 let database_id = DatabaseId::new(2);
1228 let job_id = JobId::new(20);
1229 let fragment_id = FragmentId::new(120);
1230 let actor_id = ActorId::new(500);
1231
1232 let mut inflight_jobs: FragmentRenderMap = HashMap::new();
1233 inflight_jobs
1234 .entry(database_id)
1235 .or_default()
1236 .entry(job_id)
1237 .or_default()
1238 .insert(
1239 fragment_id,
1240 InflightFragmentInfo {
1241 fragment_id,
1242 distribution_type: DistributionType::Hash,
1243 fragment_type_mask: FragmentTypeMask::empty(),
1244 vnode_count: 1,
1245 nodes: PbStreamNode::default(),
1246 actors: HashMap::from([(
1247 actor_id,
1248 InflightActorInfo {
1249 worker_id: WorkerId::new(1),
1250 vnode_bitmap: None,
1251 splits: vec![],
1252 },
1253 )]),
1254 state_table_ids: HashSet::new(),
1255 },
1256 );
1257
1258 let job_extra_info = HashMap::from([(
1259 job_id,
1260 StreamingJobExtraInfo {
1261 timezone: Some("UTC".to_owned()),
1262 config_override: "cfg".into(),
1263 job_definition: "definition".to_owned(),
1264 backfill_orders: None,
1265 refresh_interval_sec: None,
1266 },
1267 )]);
1268
1269 let stream_actors = build_stream_actors(&inflight_jobs, &job_extra_info).unwrap();
1270
1271 let actor = stream_actors.get(&actor_id).unwrap();
1272 assert_eq!(actor.actor_id, actor_id);
1273 assert_eq!(actor.fragment_id, fragment_id);
1274 assert_eq!(actor.mview_definition, "definition");
1275 assert_eq!(&*actor.config_override, "cfg");
1276 let expr_ctx = actor.expr_context.as_ref().unwrap();
1277 assert_eq!(expr_ctx.time_zone, "UTC");
1278 }
1279}