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::{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    /// Batch refresh jobs rendered during `render_runtime_info`.
90    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    // Extract batch refresh jobs before rendering via `render_actor_assignments`.
109    // They will be rendered independently using the unified `render_actors_and_build_job_info`.
110    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        // Remove ensembles that only contain batch refresh fragments.
152        per_database_context.ensembles.retain(|ensemble| {
153            !ensemble
154                .component_fragments()
155                .all(|fid| batch_refresh_fragment_ids.contains(&fid))
156        });
157    }
158
159    // Render actors for each batch refresh job.
160    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 all fragments were batch refresh, no normal rendering needed.
189    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
237/// For normal DDL operations, the `UpstreamSinkUnion` operator is modified dynamically, and does not persist the
238/// newly added or deleted upstreams in meta-store. Therefore, when restoring jobs, we need to restore the
239/// information required by the operator based on the current state of the upstream (sink) and downstream (table) of
240/// the operator. All necessary metadata must be preloaded before rendering.
241fn 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
278/// Assembles `StreamActor` instances from rendered fragment info and job context.
279///
280/// This function combines the actor assignments from `FragmentRenderMap` with
281/// runtime context (timezone, config, definition) from `StreamingJobExtraInfo`
282/// to produce the final `StreamActor` structures needed for recovery.
283fn 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    /// Clean catalogs for jobs in the initial state and creating jobs whose progress cannot be
426    /// recovered.
427    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            // Per-database recovery does not run the global Hummock purge below. Unregister the
444            // dirty jobs cleaned in this database through the normal dropped-table cleanup path.
445            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        // unregister cleaned sources.
456        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    /// Re-register iceberg pk-index sink commit coordinators after recovery wipes them.
466    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                    // This should not happen unless catalog changes or legacy metadata are involved.
682                    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                    // We must abort dirty pending sink state before registering iceberg pk-index sinks,
849                    // otherwise recover_pending will take speculative (epoch > committed epoch) pending sink state
850                    // as valid and cause duplicated iceberg commit.
851                    self.reregister_iceberg_pk_index_sinks(None)
852                        .await
853                        .context("re-register iceberg v3 sinks after recovery")?;
854
855                    // Creating job progress needs to be recovered.
856                    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                    // This is a quick path to accelerate the process of dropping and canceling streaming jobs.
865                    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                    // Resolve actor info for recovery. If there's no actor to recover, most of the
910                    // following steps will be no-op, while the compute nodes will still be reset.
911                    // TODO(error-handling): attach context to the errors and log them together, instead of inspecting everywhere.
912                    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                    // Refresh creating job progress for the final snapshot to reflect any catalog changes.
969                    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        // Creating job progress needs to be recovered.
1026        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        // This is a quick path to accelerate the process of dropping and canceling streaming jobs.
1035        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}