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