Skip to main content

risingwave_meta/barrier/context/
recovery.rs

1// Copyright 2024 RisingWave Labs
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use 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    /// Batch refresh jobs rendered during `render_runtime_info`.
90    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    // Extract batch refresh jobs before rendering via `render_actor_assignments`.
108    // They will be rendered independently using the unified `render_actors_and_build_job_info`.
109    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        // Remove ensembles that only contain batch refresh fragments.
151        per_database_context.ensembles.retain(|ensemble| {
152            !ensemble
153                .component_fragments()
154                .all(|fid| batch_refresh_fragment_ids.contains(&fid))
155        });
156    }
157
158    // Render actors for each batch refresh job.
159    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 all fragments were batch refresh, no normal rendering needed.
187    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
235/// For normal DDL operations, the `UpstreamSinkUnion` operator is modified dynamically, and does not persist the
236/// newly added or deleted upstreams in meta-store. Therefore, when restoring jobs, we need to restore the
237/// information required by the operator based on the current state of the upstream (sink) and downstream (table) of
238/// the operator. All necessary metadata must be preloaded before rendering.
239fn 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
276/// Assembles `StreamActor` instances from rendered fragment info and job context.
277///
278/// This function combines the actor assignments from `FragmentRenderMap` with
279/// runtime context (timezone, config, definition) from `StreamingJobExtraInfo`
280/// to produce the final `StreamActor` structures needed for recovery.
281fn 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    /// Clean catalogs for jobs in the initial state and creating jobs whose progress cannot be
424    /// recovered.
425    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            // Per-database recovery does not run the global Hummock purge below. Unregister the
442            // dirty jobs cleaned in this database through the normal dropped-table cleanup path.
443            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        // unregister cleaned sources.
458        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    /// Re-register iceberg pk-index sink commit coordinators after recovery wipes them.
487    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                    // This should not happen unless catalog changes or legacy metadata are involved.
703                    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                    // We must abort dirty pending sink state before registering iceberg pk-index sinks,
873                    // otherwise recover_pending will take speculative (epoch > committed epoch) pending sink state
874                    // as valid and cause duplicated iceberg commit.
875                    self.reregister_iceberg_pk_index_sinks(None)
876                        .await
877                        .context("re-register iceberg v3 sinks after recovery")?;
878
879                    // Creating job progress needs to be recovered.
880                    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                    // This is a quick path to accelerate the process of dropping and canceling streaming jobs.
889                    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                    // Resolve actor info for recovery. If there's no actor to recover, most of the
933                    // following steps will be no-op, while the compute nodes will still be reset.
934                    // TODO(error-handling): attach context to the errors and log them together, instead of inspecting everywhere.
935                    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                    // Refresh creating job progress for the final snapshot to reflect any catalog changes.
992                    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        // Creating job progress needs to be recovered.
1052        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        // This is a quick path to accelerate the process of dropping and canceling streaming jobs.
1061        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}