1use std::collections::hash_map::Entry;
16use std::collections::{HashMap, HashSet};
17use std::mem::replace;
18use std::sync::Arc;
19
20use itertools::Itertools;
21use parking_lot::RawRwLock;
22use parking_lot::lock_api::RwLockReadGuard;
23use risingwave_common::bitmap::Bitmap;
24use risingwave_common::catalog::{DatabaseId, FragmentTypeFlag, FragmentTypeMask, TableId};
25use risingwave_common::id::JobId;
26use risingwave_common::util::epoch::EpochPair;
27use risingwave_common::util::stream_graph_visitor::visit_stream_node_mut;
28use risingwave_connector::source::{SplitImpl, SplitMetaData};
29use risingwave_meta_model::WorkerId;
30use risingwave_meta_model::fragment::DistributionType;
31use risingwave_pb::ddl_service::PbBackfillType;
32use risingwave_pb::hummock::HummockVersionStats;
33use risingwave_pb::id::SubscriberId;
34use risingwave_pb::meta::PbFragmentWorkerSlotMapping;
35use risingwave_pb::meta::subscribe_response::Operation;
36use risingwave_pb::source::PbCdcTableSnapshotSplits;
37use risingwave_pb::stream_plan::PbUpstreamSinkInfo;
38use risingwave_pb::stream_plan::barrier_mutation::Mutation;
39use risingwave_pb::stream_plan::stream_node::NodeBody;
40use risingwave_pb::stream_service::BarrierCompleteResponse;
41use tracing::{info, warn};
42
43use crate::barrier::cdc_progress::{CdcProgress, CdcTableBackfillTracker};
44use crate::barrier::command::{
45 CreateStreamingJobCommandInfo, PostCollectCommand, ReplaceStreamJobPlan, ThrottleConfigMap,
46 extract_throttle_config,
47};
48use crate::barrier::edge_builder::{FragmentEdgeBuildResult, FragmentEdgeBuilder};
49use crate::barrier::progress::{CreateMviewProgressTracker, StagingCommitInfo};
50use crate::barrier::rpc::{ControlStreamManager, to_partial_graph_id};
51use crate::barrier::{
52 BackfillProgress, BarrierKind, CreateStreamingJobType, FragmentBackfillProgress, TracedEpoch,
53};
54use crate::controller::fragment::{InflightActorInfo, InflightFragmentInfo};
55use crate::controller::utils::rebuild_fragment_mapping;
56use crate::manager::NotificationManagerRef;
57use crate::model::{
58 ActorId, ActorNewNoShuffle, BackfillUpstreamType, FragmentId, StreamActor, StreamJobFragments,
59};
60use crate::stream::UpstreamSinkInfo;
61use crate::{MetaError, MetaResult};
62
63#[derive(Debug, Clone)]
64pub struct SharedActorInfo {
65 pub worker_id: WorkerId,
66 pub vnode_bitmap: Option<Bitmap>,
67 pub splits: Vec<SplitImpl>,
68}
69
70impl From<&InflightActorInfo> for SharedActorInfo {
71 fn from(value: &InflightActorInfo) -> Self {
72 Self {
73 worker_id: value.worker_id,
74 vnode_bitmap: value.vnode_bitmap.clone(),
75 splits: value.splits.clone(),
76 }
77 }
78}
79
80#[derive(Debug, Clone)]
81pub struct SharedFragmentInfo {
82 pub fragment_id: FragmentId,
83 pub job_id: JobId,
84 pub distribution_type: DistributionType,
85 pub actors: HashMap<ActorId, SharedActorInfo>,
86 pub vnode_count: usize,
87 pub fragment_type_mask: FragmentTypeMask,
88 pub state_table_ids: HashSet<TableId>,
89}
90
91impl From<(&InflightFragmentInfo, JobId)> for SharedFragmentInfo {
92 fn from(pair: (&InflightFragmentInfo, JobId)) -> Self {
93 let (info, job_id) = pair;
94
95 let InflightFragmentInfo {
96 fragment_id,
97 distribution_type,
98 fragment_type_mask,
99 actors,
100 vnode_count,
101 state_table_ids,
102 ..
103 } = info;
104
105 Self {
106 fragment_id: *fragment_id,
107 job_id,
108 distribution_type: *distribution_type,
109 fragment_type_mask: *fragment_type_mask,
110 actors: actors
111 .iter()
112 .map(|(actor_id, actor)| (*actor_id, actor.into()))
113 .collect(),
114 vnode_count: *vnode_count,
115 state_table_ids: state_table_ids.clone(),
116 }
117 }
118}
119
120#[derive(Default, Debug)]
121pub struct SharedActorInfosInner {
122 info: HashMap<DatabaseId, HashMap<FragmentId, SharedFragmentInfo>>,
123}
124
125impl SharedActorInfosInner {
126 pub fn get_fragment(&self, fragment_id: FragmentId) -> Option<&SharedFragmentInfo> {
127 self.info
128 .values()
129 .find_map(|database| database.get(&fragment_id))
130 }
131
132 pub fn iter_over_fragments(&self) -> impl Iterator<Item = (&FragmentId, &SharedFragmentInfo)> {
133 self.info.values().flatten()
134 }
135}
136
137#[derive(Clone, educe::Educe)]
138#[educe(Debug)]
139pub struct SharedActorInfos {
140 inner: Arc<parking_lot::RwLock<SharedActorInfosInner>>,
141 #[educe(Debug(ignore))]
142 notification_manager: NotificationManagerRef,
143}
144
145impl SharedActorInfos {
146 pub fn read_guard(&self) -> RwLockReadGuard<'_, RawRwLock, SharedActorInfosInner> {
147 self.inner.read()
148 }
149
150 pub fn list_assignments(&self) -> HashMap<ActorId, Vec<SplitImpl>> {
151 let core = self.inner.read();
152 core.iter_over_fragments()
153 .flat_map(|(_, fragment)| {
154 fragment
155 .actors
156 .iter()
157 .map(|(actor_id, info)| (*actor_id, info.splits.clone()))
158 })
159 .collect()
160 }
161
162 pub fn migrate_splits_for_source_actors(
168 &self,
169 fragment_id: FragmentId,
170 prev_actor_ids: &[ActorId],
171 curr_actor_ids: &[ActorId],
172 ) -> MetaResult<HashMap<ActorId, Vec<SplitImpl>>> {
173 let guard = self.read_guard();
174
175 let prev_splits = prev_actor_ids
176 .iter()
177 .flat_map(|actor_id| {
178 guard
180 .get_fragment(fragment_id)
181 .and_then(|info| info.actors.get(actor_id))
182 .map(|actor| actor.splits.clone())
183 .unwrap_or_default()
184 })
185 .map(|split| (split.id(), split))
186 .collect();
187
188 let empty_actor_splits = curr_actor_ids
189 .iter()
190 .map(|actor_id| (*actor_id, vec![]))
191 .collect();
192
193 let diff = crate::stream::source_manager::reassign_splits(
194 fragment_id,
195 empty_actor_splits,
196 &prev_splits,
197 std::default::Default::default(),
199 )
200 .unwrap_or_default();
201
202 Ok(diff)
203 }
204}
205
206impl SharedActorInfos {
207 pub(crate) fn new(notification_manager: NotificationManagerRef) -> Self {
208 Self {
209 inner: Arc::new(Default::default()),
210 notification_manager,
211 }
212 }
213
214 pub(super) fn remove_database(&self, database_id: DatabaseId) {
215 if let Some(database) = self.inner.write().info.remove(&database_id) {
216 let mapping = database
217 .into_values()
218 .map(|fragment| rebuild_fragment_mapping(&fragment))
219 .collect_vec();
220 if !mapping.is_empty() {
221 self.notification_manager
222 .notify_streaming_fragment_mapping(Operation::Delete, mapping);
223 }
224 }
225 }
226
227 pub(super) fn retain_databases(&self, database_ids: impl IntoIterator<Item = DatabaseId>) {
228 let database_ids: HashSet<_> = database_ids.into_iter().collect();
229
230 let mut mapping = Vec::new();
231 for fragment in self
232 .inner
233 .write()
234 .info
235 .extract_if(|database_id, _| !database_ids.contains(database_id))
236 .flat_map(|(_, fragments)| fragments.into_values())
237 {
238 mapping.push(rebuild_fragment_mapping(&fragment));
239 }
240 if !mapping.is_empty() {
241 self.notification_manager
242 .notify_streaming_fragment_mapping(Operation::Delete, mapping);
243 }
244 }
245
246 pub(super) fn recover_database(
247 &self,
248 database_id: DatabaseId,
249 fragments: impl Iterator<Item = (&InflightFragmentInfo, JobId)>,
250 ) {
251 let mut remaining_fragments: HashMap<_, _> = fragments
252 .map(|info @ (fragment, _)| (fragment.fragment_id, info))
253 .collect();
254 let mut writer = self.start_writer(database_id);
256 let database = writer.write_guard.info.entry(database_id).or_default();
257 for (_, fragment) in database.extract_if(|fragment_id, fragment_infos| {
258 if let Some(info) = remaining_fragments.remove(fragment_id) {
259 let info = info.into();
260 writer
261 .updated_fragment_mapping
262 .get_or_insert_default()
263 .push(rebuild_fragment_mapping(&info));
264 *fragment_infos = info;
265 false
266 } else {
267 true
268 }
269 }) {
270 writer
271 .deleted_fragment_mapping
272 .get_or_insert_default()
273 .push(rebuild_fragment_mapping(&fragment));
274 }
275 for (fragment_id, info) in remaining_fragments {
276 let info = info.into();
277 writer
278 .added_fragment_mapping
279 .get_or_insert_default()
280 .push(rebuild_fragment_mapping(&info));
281 database.insert(fragment_id, info);
282 }
283 writer.finish();
284 }
285
286 pub(super) fn upsert(
287 &self,
288 database_id: DatabaseId,
289 infos: impl IntoIterator<Item = (&InflightFragmentInfo, JobId)>,
290 ) {
291 let mut writer = self.start_writer(database_id);
292 writer.upsert(infos);
293 writer.finish();
294 }
295
296 pub(super) fn start_writer(&self, database_id: DatabaseId) -> SharedActorInfoWriter<'_> {
297 SharedActorInfoWriter {
298 database_id,
299 write_guard: self.inner.write(),
300 notification_manager: &self.notification_manager,
301 added_fragment_mapping: None,
302 updated_fragment_mapping: None,
303 deleted_fragment_mapping: None,
304 }
305 }
306}
307
308pub(super) struct SharedActorInfoWriter<'a> {
309 database_id: DatabaseId,
310 write_guard: parking_lot::RwLockWriteGuard<'a, SharedActorInfosInner>,
311 notification_manager: &'a NotificationManagerRef,
312 added_fragment_mapping: Option<Vec<PbFragmentWorkerSlotMapping>>,
313 updated_fragment_mapping: Option<Vec<PbFragmentWorkerSlotMapping>>,
314 deleted_fragment_mapping: Option<Vec<PbFragmentWorkerSlotMapping>>,
315}
316
317impl SharedActorInfoWriter<'_> {
318 pub(super) fn upsert(
319 &mut self,
320 infos: impl IntoIterator<Item = (&InflightFragmentInfo, JobId)>,
321 ) {
322 let database = self.write_guard.info.entry(self.database_id).or_default();
323 for info @ (fragment, _) in infos {
324 match database.entry(fragment.fragment_id) {
325 Entry::Occupied(mut entry) => {
326 let info = info.into();
327 self.updated_fragment_mapping
328 .get_or_insert_default()
329 .push(rebuild_fragment_mapping(&info));
330 entry.insert(info);
331 }
332 Entry::Vacant(entry) => {
333 let info = info.into();
334 self.added_fragment_mapping
335 .get_or_insert_default()
336 .push(rebuild_fragment_mapping(&info));
337 entry.insert(info);
338 }
339 }
340 }
341 }
342
343 pub(super) fn remove(&mut self, info: &InflightFragmentInfo) {
344 if let Some(database) = self.write_guard.info.get_mut(&self.database_id)
345 && let Some(fragment) = database.remove(&info.fragment_id)
346 {
347 self.deleted_fragment_mapping
348 .get_or_insert_default()
349 .push(rebuild_fragment_mapping(&fragment));
350 }
351 }
352
353 pub(super) fn finish(self) {
354 if let Some(mapping) = self.added_fragment_mapping {
355 self.notification_manager
356 .notify_streaming_fragment_mapping(Operation::Add, mapping);
357 }
358 if let Some(mapping) = self.updated_fragment_mapping {
359 self.notification_manager
360 .notify_streaming_fragment_mapping(Operation::Update, mapping);
361 }
362 if let Some(mapping) = self.deleted_fragment_mapping {
363 self.notification_manager
364 .notify_streaming_fragment_mapping(Operation::Delete, mapping);
365 }
366 }
367}
368
369#[derive(Debug, Clone)]
370pub(super) struct BarrierInfo {
371 pub prev_epoch: TracedEpoch,
372 pub curr_epoch: TracedEpoch,
373 pub kind: BarrierKind,
374}
375
376impl BarrierInfo {
377 pub(super) fn prev_epoch(&self) -> u64 {
378 self.prev_epoch.value().0
379 }
380
381 pub(super) fn curr_epoch(&self) -> u64 {
382 self.curr_epoch.value().0
383 }
384
385 pub(super) fn epoch(&self) -> EpochPair {
386 EpochPair {
387 curr: self.curr_epoch(),
388 prev: self.prev_epoch(),
389 }
390 }
391}
392
393#[derive(Clone, Debug)]
394pub enum SubscriberType {
395 Subscription(u64),
396 SnapshotBackfill,
397}
398
399#[derive(Debug)]
400pub(super) enum CreateStreamingJobStatus {
401 Init,
402 Creating { tracker: CreateMviewProgressTracker },
403 Created,
404}
405
406#[derive(Debug)]
407pub(super) struct InflightStreamingJobInfo {
408 pub job_id: JobId,
409 pub fragment_infos: HashMap<FragmentId, InflightFragmentInfo>,
410 pub subscribers: HashMap<SubscriberId, SubscriberType>,
411 pub status: CreateStreamingJobStatus,
412 pub cdc_table_backfill_tracker: Option<CdcTableBackfillTracker>,
413}
414
415impl InflightStreamingJobInfo {
416 pub fn fragment_infos(&self) -> impl Iterator<Item = &InflightFragmentInfo> + '_ {
417 self.fragment_infos.values()
418 }
419
420 pub fn snapshot_backfill_actor_ids(
421 fragment_infos: &HashMap<FragmentId, InflightFragmentInfo>,
422 ) -> impl Iterator<Item = ActorId> + '_ {
423 fragment_infos
424 .values()
425 .filter(|fragment| {
426 fragment
427 .fragment_type_mask
428 .contains(FragmentTypeFlag::SnapshotBackfillStreamScan)
429 })
430 .flat_map(|fragment| fragment.actors.keys().copied())
431 }
432
433 pub fn tracking_progress_actor_ids(
434 fragment_infos: &HashMap<FragmentId, InflightFragmentInfo>,
435 ) -> Vec<(ActorId, BackfillUpstreamType)> {
436 StreamJobFragments::tracking_progress_actor_ids_impl(fragment_infos.values().map(
437 |fragment| {
438 (
439 fragment.fragment_type_mask,
440 &fragment.nodes,
441 fragment.actors.keys().copied(),
442 )
443 },
444 ))
445 }
446}
447
448impl<'a> IntoIterator for &'a InflightStreamingJobInfo {
449 type Item = &'a InflightFragmentInfo;
450
451 type IntoIter = impl Iterator<Item = &'a InflightFragmentInfo> + 'a;
452
453 fn into_iter(self) -> Self::IntoIter {
454 self.fragment_infos()
455 }
456}
457
458#[derive(Debug)]
459pub struct InflightDatabaseInfo {
460 pub(super) database_id: DatabaseId,
461 jobs: HashMap<JobId, InflightStreamingJobInfo>,
462 fragment_location: HashMap<FragmentId, JobId>,
463 pub(super) shared_actor_infos: SharedActorInfos,
464}
465
466impl InflightDatabaseInfo {
467 pub(super) fn job_ids(&self) -> impl Iterator<Item = JobId> + '_ {
468 self.jobs.keys().copied()
469 }
470
471 pub fn fragment_infos(&self) -> impl Iterator<Item = &InflightFragmentInfo> + '_ {
472 self.jobs.values().flat_map(|job| job.fragment_infos())
473 }
474
475 pub fn contains_job(&self, job_id: JobId) -> bool {
476 self.jobs.contains_key(&job_id)
477 }
478
479 pub fn job_fragment_infos(
481 &self,
482 job_id: JobId,
483 ) -> impl Iterator<Item = &InflightFragmentInfo> + '_ {
484 self.jobs
485 .get(&job_id)
486 .into_iter()
487 .flat_map(|job| job.fragment_infos())
488 }
489
490 pub(super) fn job_id_by_fragment(&self, fragment_id: FragmentId) -> Option<JobId> {
491 self.fragment_location.get(&fragment_id).copied()
492 }
493
494 pub fn fragment(&self, fragment_id: FragmentId) -> &InflightFragmentInfo {
495 let job_id = self.fragment_location[&fragment_id];
496 self.jobs
497 .get(&job_id)
498 .expect("should exist")
499 .fragment_infos
500 .get(&fragment_id)
501 .expect("should exist")
502 }
503
504 pub(super) fn backfill_fragment_ids_for_job(
505 &self,
506 job_id: JobId,
507 ) -> MetaResult<HashSet<FragmentId>> {
508 let job = self
509 .jobs
510 .get(&job_id)
511 .ok_or_else(|| MetaError::invalid_parameter(format!("job {} not found", job_id)))?;
512 Ok(job
513 .fragment_infos
514 .iter()
515 .filter_map(|(fragment_id, fragment)| {
516 fragment
517 .fragment_type_mask
518 .contains_any([
519 FragmentTypeFlag::StreamScan,
520 FragmentTypeFlag::SourceScan,
521 FragmentTypeFlag::LocalityProvider,
522 ])
523 .then_some(*fragment_id)
524 })
525 .collect())
526 }
527
528 pub(super) fn is_backfill_fragment(&self, fragment_id: FragmentId) -> MetaResult<bool> {
529 let job_id = self.fragment_location.get(&fragment_id).ok_or_else(|| {
530 MetaError::invalid_parameter(format!("fragment {} not found", fragment_id))
531 })?;
532 let fragment = self
533 .jobs
534 .get(job_id)
535 .expect("should exist")
536 .fragment_infos
537 .get(&fragment_id)
538 .expect("should exist");
539 Ok(fragment.fragment_type_mask.contains_any([
540 FragmentTypeFlag::StreamScan,
541 FragmentTypeFlag::SourceScan,
542 FragmentTypeFlag::LocalityProvider,
543 ]))
544 }
545
546 pub fn gen_backfill_progress(&self) -> impl Iterator<Item = (JobId, BackfillProgress)> + '_ {
547 self.jobs
548 .iter()
549 .filter_map(|(job_id, job)| match &job.status {
550 CreateStreamingJobStatus::Init => None,
551 CreateStreamingJobStatus::Creating { tracker } => {
552 let progress = tracker.gen_backfill_progress();
553 Some((
554 *job_id,
555 BackfillProgress {
556 progress,
557 backfill_type: PbBackfillType::NormalBackfill,
558 },
559 ))
560 }
561 CreateStreamingJobStatus::Created => None,
562 })
563 }
564
565 pub fn gen_cdc_progress(&self) -> impl Iterator<Item = (JobId, CdcProgress)> + '_ {
566 self.jobs.iter().filter_map(|(job_id, job)| {
567 job.cdc_table_backfill_tracker
568 .as_ref()
569 .map(|tracker| (*job_id, tracker.gen_cdc_progress()))
570 })
571 }
572
573 pub fn gen_fragment_backfill_progress(&self) -> Vec<FragmentBackfillProgress> {
574 let mut result = Vec::new();
575 for job in self.jobs.values() {
576 let CreateStreamingJobStatus::Creating { tracker } = &job.status else {
577 continue;
578 };
579 let fragment_progress = tracker.collect_fragment_progress(&job.fragment_infos, true);
580 result.extend(fragment_progress);
581 }
582 result
583 }
584
585 pub(super) fn may_assign_fragment_cdc_backfill_splits(
586 &mut self,
587 fragment_id: FragmentId,
588 ) -> MetaResult<Option<HashMap<ActorId, PbCdcTableSnapshotSplits>>> {
589 let job_id = self.fragment_location[&fragment_id];
590 let job = self.jobs.get_mut(&job_id).expect("should exist");
591 if let Some(tracker) = &mut job.cdc_table_backfill_tracker {
592 let cdc_scan_fragment_id = tracker.cdc_scan_fragment_id();
593 if cdc_scan_fragment_id != fragment_id {
594 return Ok(None);
595 }
596 let actors = job.fragment_infos[&cdc_scan_fragment_id]
597 .actors
598 .keys()
599 .copied()
600 .collect();
601 tracker.reassign_splits(actors).map(Some)
602 } else {
603 Ok(None)
604 }
605 }
606
607 pub(super) fn assign_cdc_backfill_splits(
608 &mut self,
609 job_id: JobId,
610 ) -> MetaResult<Option<HashMap<ActorId, PbCdcTableSnapshotSplits>>> {
611 let job = self.jobs.get_mut(&job_id).expect("should exist");
612 if let Some(tracker) = &mut job.cdc_table_backfill_tracker {
613 let cdc_scan_fragment_id = tracker.cdc_scan_fragment_id();
614 let actors = job.fragment_infos[&cdc_scan_fragment_id]
615 .actors
616 .keys()
617 .copied()
618 .collect();
619 tracker.reassign_splits(actors).map(Some)
620 } else {
621 Ok(None)
622 }
623 }
624
625 pub(super) fn apply_collected_command(
626 &mut self,
627 command: &PostCollectCommand,
628 resps: &HashMap<WorkerId, BarrierCompleteResponse>,
629 version_stats: &HummockVersionStats,
630 ) {
631 if let PostCollectCommand::CreateStreamingJob { info, job_type, .. } = command {
632 match job_type {
633 CreateStreamingJobType::Normal | CreateStreamingJobType::SinkIntoTable(_) => {
634 let job_id = info.streaming_job.id();
635 if let Some(job_info) = self.jobs.get_mut(&job_id) {
636 let CreateStreamingJobStatus::Init = replace(
637 &mut job_info.status,
638 CreateStreamingJobStatus::Creating {
639 tracker: CreateMviewProgressTracker::new(
640 info,
641 version_stats,
642 &job_info.fragment_infos,
643 ),
644 },
645 ) else {
646 unreachable!("should be init before collect the first barrier")
647 };
648 } else {
649 info!(%job_id, "newly create job get cancelled before first barrier is collected")
650 }
651 }
652 CreateStreamingJobType::Independent { .. } => {
653 }
655 }
656 }
657 if let PostCollectCommand::Reschedule { reschedules, .. } = command {
658 debug_assert!(
660 reschedules
661 .values()
662 .all(|reschedule| reschedule.vnode_bitmap_updates.is_empty()),
663 "Reschedule should not carry vnode bitmap updates when actors are rebuilt"
664 );
665
666 let related_job_ids = reschedules
668 .keys()
669 .filter_map(|fragment_id| self.fragment_location.get(fragment_id))
670 .cloned()
671 .collect::<HashSet<_>>();
672 for job_id in related_job_ids {
673 if let Some(job) = self.jobs.get_mut(&job_id)
674 && let CreateStreamingJobStatus::Creating { tracker, .. } = &mut job.status
675 {
676 tracker.refresh_after_reschedule(&job.fragment_infos, version_stats);
677 }
678 }
679 }
680 for progress in resps.values().flat_map(|resp| &resp.create_mview_progress) {
681 let Some(job_id) = self.fragment_location.get(&progress.fragment_id) else {
682 warn!(
683 "update the progress of an non-existent creating streaming job: {progress:?}, which could be cancelled"
684 );
685 continue;
686 };
687 let tracker = match &mut self.jobs.get_mut(job_id).expect("should exist").status {
688 CreateStreamingJobStatus::Init => {
689 continue;
690 }
691 CreateStreamingJobStatus::Creating { tracker, .. } => tracker,
692 CreateStreamingJobStatus::Created => {
693 if !progress.done {
694 warn!("update the progress of an created streaming job: {progress:?}");
695 }
696 continue;
697 }
698 };
699 tracker.apply_progress(progress, version_stats);
700 }
701 for progress in resps
702 .values()
703 .flat_map(|resp| &resp.cdc_table_backfill_progress)
704 {
705 let Some(job_id) = self.fragment_location.get(&progress.fragment_id) else {
706 warn!(
707 "update the cdc progress of an non-existent creating streaming job: {progress:?}, which could be cancelled"
708 );
709 continue;
710 };
711 let Some(tracker) = &mut self
712 .jobs
713 .get_mut(job_id)
714 .expect("should exist")
715 .cdc_table_backfill_tracker
716 else {
717 warn!("update the cdc progress of an created streaming job: {progress:?}");
718 continue;
719 };
720 tracker.update_split_progress(progress);
721 }
722 for cdc_offset_updated in resps
724 .values()
725 .flat_map(|resp| &resp.cdc_source_offset_updated)
726 {
727 let source_id = cdc_offset_updated.source_id;
728 let job_id = source_id.as_share_source_job_id();
729 if let Some(job) = self.jobs.get_mut(&job_id) {
730 if let CreateStreamingJobStatus::Creating { tracker, .. } = &mut job.status {
731 tracker.mark_cdc_source_finished();
732 }
733 } else {
734 warn!(
735 "update cdc source offset for non-existent creating streaming job: source_id={}, job_id={}",
736 cdc_offset_updated.source_id, job_id
737 );
738 }
739 }
740 }
741
742 fn iter_creating_job_tracker(&self) -> impl Iterator<Item = &CreateMviewProgressTracker> {
743 self.jobs.values().filter_map(|job| match &job.status {
744 CreateStreamingJobStatus::Init => None,
745 CreateStreamingJobStatus::Creating { tracker, .. } => Some(tracker),
746 CreateStreamingJobStatus::Created => None,
747 })
748 }
749
750 fn iter_mut_creating_job_tracker(
751 &mut self,
752 ) -> impl Iterator<Item = &mut CreateMviewProgressTracker> {
753 self.jobs
754 .values_mut()
755 .filter_map(|job| match &mut job.status {
756 CreateStreamingJobStatus::Init => None,
757 CreateStreamingJobStatus::Creating { tracker, .. } => Some(tracker),
758 CreateStreamingJobStatus::Created => None,
759 })
760 }
761
762 pub(super) fn has_pending_finished_jobs(&self) -> bool {
763 self.iter_creating_job_tracker()
764 .any(|tracker| tracker.is_finished())
765 }
766
767 pub(super) fn take_pending_backfill_nodes(&mut self) -> Vec<FragmentId> {
768 self.iter_mut_creating_job_tracker()
769 .flat_map(|tracker| tracker.take_pending_backfill_nodes())
770 .collect()
771 }
772
773 pub(super) fn take_staging_commit_info(&mut self) -> StagingCommitInfo {
774 let mut finished_jobs = vec![];
775 let mut table_ids_to_truncate = vec![];
776 let mut finished_cdc_table_backfill = vec![];
777 for (job_id, job) in &mut self.jobs {
778 if let CreateStreamingJobStatus::Creating { tracker, .. } = &mut job.status {
779 let (is_finished, truncate_table_ids) = tracker.collect_staging_commit_info();
780 table_ids_to_truncate.extend(truncate_table_ids);
781 if is_finished {
782 let CreateStreamingJobStatus::Creating { tracker, .. } =
783 replace(&mut job.status, CreateStreamingJobStatus::Created)
784 else {
785 unreachable!()
786 };
787 finished_jobs.push(tracker.into_tracking_job());
788 }
789 }
790 if let Some(tracker) = &mut job.cdc_table_backfill_tracker
791 && tracker.take_pre_completed()
792 {
793 finished_cdc_table_backfill.push(*job_id);
794 }
795 }
796 StagingCommitInfo {
797 finished_jobs,
798 table_ids_to_truncate,
799 finished_cdc_table_backfill,
800 }
801 }
802
803 pub fn fragment_subscribers(
804 &self,
805 fragment_id: FragmentId,
806 ) -> impl Iterator<Item = SubscriberId> + '_ {
807 let job_id = self.fragment_location[&fragment_id];
808 self.jobs[&job_id].subscribers.keys().copied()
809 }
810
811 pub fn job_subscribers(&self, job_id: JobId) -> impl Iterator<Item = SubscriberId> + '_ {
812 self.jobs[&job_id].subscribers.keys().copied()
813 }
814
815 pub fn subscribed_tables(&self) -> impl Iterator<Item = TableId> + '_ {
816 self.jobs.iter().filter_map(|(job_id, info)| {
817 info.subscribers
818 .values()
819 .any(|subscriber| matches!(subscriber, SubscriberType::Subscription(_)))
820 .then_some(job_id.as_mv_table_id())
821 })
822 }
823
824 pub fn register_subscriber(
825 &mut self,
826 job_id: JobId,
827 subscriber_id: SubscriberId,
828 subscriber: SubscriberType,
829 ) {
830 self.jobs
831 .get_mut(&job_id)
832 .expect("should exist")
833 .subscribers
834 .try_insert(subscriber_id, subscriber)
835 .expect("non duplicate");
836 }
837
838 pub fn unregister_subscriber(
839 &mut self,
840 job_id: JobId,
841 subscriber_id: SubscriberId,
842 ) -> Option<SubscriberType> {
843 self.jobs
844 .get_mut(&job_id)
845 .expect("should exist")
846 .subscribers
847 .remove(&subscriber_id)
848 }
849
850 pub fn update_subscription_retention(
851 &mut self,
852 job_id: JobId,
853 subscriber_id: SubscriberId,
854 retention_second: u64,
855 ) {
856 let job = self.jobs.get_mut(&job_id).expect("should exist");
857 match job.subscribers.get_mut(&subscriber_id) {
858 Some(SubscriberType::Subscription(current_retention)) => {
859 *current_retention = retention_second;
860 }
861 Some(SubscriberType::SnapshotBackfill) => {
862 warn!(
863 %job_id,
864 %subscriber_id,
865 "cannot update retention for snapshot backfill subscriber"
866 );
867 }
868 None => {
869 warn!(%job_id, %subscriber_id, "subscription subscriber not found");
870 }
871 }
872 }
873
874 fn fragment_mut(&mut self, fragment_id: FragmentId) -> (&mut InflightFragmentInfo, JobId) {
875 let job_id = self.fragment_location[&fragment_id];
876 let fragment = self
877 .jobs
878 .get_mut(&job_id)
879 .expect("should exist")
880 .fragment_infos
881 .get_mut(&fragment_id)
882 .expect("should exist");
883 (fragment, job_id)
884 }
885
886 fn empty_inner(database_id: DatabaseId, shared_actor_infos: SharedActorInfos) -> Self {
887 Self {
888 database_id,
889 jobs: Default::default(),
890 fragment_location: Default::default(),
891 shared_actor_infos,
892 }
893 }
894
895 pub fn empty(database_id: DatabaseId, shared_actor_infos: SharedActorInfos) -> Self {
896 shared_actor_infos.remove_database(database_id);
898 Self::empty_inner(database_id, shared_actor_infos)
899 }
900
901 pub fn recover(
902 database_id: DatabaseId,
903 jobs: impl Iterator<Item = InflightStreamingJobInfo>,
904 shared_actor_infos: SharedActorInfos,
905 ) -> Self {
906 let mut info = Self::empty_inner(database_id, shared_actor_infos);
907 for job in jobs {
908 info.add_existing(job);
909 }
910 info
911 }
912
913 pub fn is_empty(&self) -> bool {
914 self.jobs.is_empty()
915 }
916
917 pub fn add_existing(&mut self, job: InflightStreamingJobInfo) {
918 let InflightStreamingJobInfo {
919 job_id,
920 fragment_infos,
921 subscribers,
922 status,
923 cdc_table_backfill_tracker,
924 } = job;
925 self.jobs
926 .try_insert(
927 job_id,
928 InflightStreamingJobInfo {
929 job_id,
930 subscribers,
931 fragment_infos: Default::default(), status,
933 cdc_table_backfill_tracker,
934 },
935 )
936 .expect("non-duplicate");
937 self.pre_apply_new_fragments(
938 fragment_infos
939 .into_iter()
940 .map(|(fragment_id, info)| (fragment_id, job_id, info)),
941 );
942 }
943
944 pub(crate) fn pre_apply_new_job(
946 &mut self,
947 job_id: JobId,
948 cdc_table_backfill_tracker: Option<CdcTableBackfillTracker>,
949 ) {
950 {
951 self.jobs
952 .try_insert(
953 job_id,
954 InflightStreamingJobInfo {
955 job_id,
956 fragment_infos: Default::default(),
957 subscribers: Default::default(), status: CreateStreamingJobStatus::Init,
959 cdc_table_backfill_tracker,
960 },
961 )
962 .expect("non-duplicate");
963 }
964 }
965
966 pub(crate) fn pre_apply_new_fragments(
968 &mut self,
969 fragments: impl IntoIterator<Item = (FragmentId, JobId, InflightFragmentInfo)>,
970 ) {
971 {
972 let shared_infos = self.shared_actor_infos.clone();
973 let mut shared_actor_writer = shared_infos.start_writer(self.database_id);
974 for (fragment_id, job_id, info) in fragments {
975 {
976 {
977 let fragment_infos = self.jobs.get_mut(&job_id).expect("should exist");
978 shared_actor_writer.upsert([(&info, job_id)]);
979 fragment_infos
980 .fragment_infos
981 .try_insert(fragment_id, info)
982 .expect("non duplicate");
983 self.fragment_location
984 .try_insert(fragment_id, job_id)
985 .expect("non duplicate");
986 }
987 }
988 }
989 shared_actor_writer.finish();
990 }
991 }
992
993 pub(crate) fn pre_apply_reschedule(
996 &mut self,
997 fragment_id: FragmentId,
998 new_actors: HashMap<ActorId, InflightActorInfo>,
999 actor_update_vnode_bitmap: HashMap<ActorId, Bitmap>,
1000 actor_splits: HashMap<ActorId, Vec<SplitImpl>>,
1001 ) {
1002 {
1003 {
1004 {
1005 {
1006 let (info, _) = self.fragment_mut(fragment_id);
1007 let actors = &mut info.actors;
1008 for (actor_id, new_vnodes) in actor_update_vnode_bitmap {
1009 actors
1010 .get_mut(&actor_id)
1011 .expect("should exist")
1012 .vnode_bitmap = Some(new_vnodes);
1013 }
1014 for (actor_id, actor) in new_actors {
1015 actors
1016 .try_insert(actor_id as _, actor)
1017 .expect("non-duplicate");
1018 }
1019 for (actor_id, splits) in actor_splits {
1020 actors.get_mut(&actor_id).expect("should exist").splits = splits;
1021 }
1022 }
1024 }
1025 }
1026 }
1027 }
1028
1029 pub(crate) fn pre_apply_replace_node_upstream(
1031 &mut self,
1032 fragment_id: FragmentId,
1033 replace_map: &HashMap<FragmentId, FragmentId>,
1034 ) {
1035 {
1036 {
1037 {
1038 {
1039 let mut remaining_fragment_ids: HashSet<_> =
1040 replace_map.keys().cloned().collect();
1041 let (info, _) = self.fragment_mut(fragment_id);
1042 visit_stream_node_mut(&mut info.nodes, |node| {
1043 if let NodeBody::Merge(m) = node
1044 && let Some(new_upstream_fragment_id) =
1045 replace_map.get(&m.upstream_fragment_id)
1046 {
1047 if !remaining_fragment_ids.remove(&m.upstream_fragment_id) {
1048 if cfg!(debug_assertions) {
1049 panic!(
1050 "duplicate upstream fragment: {:?} {:?}",
1051 m, replace_map
1052 );
1053 } else {
1054 warn!(?m, ?replace_map, "duplicate upstream fragment");
1055 }
1056 }
1057 m.upstream_fragment_id = *new_upstream_fragment_id;
1058 }
1059 });
1060 if cfg!(debug_assertions) {
1061 assert!(
1062 remaining_fragment_ids.is_empty(),
1063 "non-existing fragment to replace: {:?} {:?} {:?}",
1064 remaining_fragment_ids,
1065 info.nodes,
1066 replace_map
1067 );
1068 } else {
1069 warn!(?remaining_fragment_ids, node = ?info.nodes, ?replace_map, "non-existing fragment to replace");
1070 }
1071 }
1072 }
1073 }
1074 }
1075 }
1076
1077 pub(crate) fn pre_apply_add_node_upstream(
1079 &mut self,
1080 fragment_id: FragmentId,
1081 new_upstream_info: &PbUpstreamSinkInfo,
1082 ) {
1083 {
1084 {
1085 {
1086 {
1087 let (info, _) = self.fragment_mut(fragment_id);
1088 let mut injected = false;
1089 visit_stream_node_mut(&mut info.nodes, |node| {
1090 if let NodeBody::UpstreamSinkUnion(u) = node {
1091 if cfg!(debug_assertions) {
1092 let current_upstream_fragment_ids = u
1093 .init_upstreams
1094 .iter()
1095 .map(|upstream| upstream.upstream_fragment_id)
1096 .collect::<HashSet<_>>();
1097 if current_upstream_fragment_ids
1098 .contains(&new_upstream_info.upstream_fragment_id)
1099 {
1100 panic!(
1101 "duplicate upstream fragment: {:?} {:?}",
1102 u, new_upstream_info
1103 );
1104 }
1105 }
1106 u.init_upstreams.push(new_upstream_info.clone());
1107 injected = true;
1108 }
1109 });
1110 assert!(injected, "should inject upstream into UpstreamSinkUnion");
1111 }
1112 }
1113 }
1114 }
1115 }
1116
1117 pub(crate) fn pre_apply_drop_node_upstream(
1119 &mut self,
1120 fragment_id: FragmentId,
1121 drop_upstream_fragment_ids: &[FragmentId],
1122 ) {
1123 if !self.fragment_location.contains_key(&fragment_id) {
1124 warn!(
1125 target_fragment_id = %fragment_id,
1126 drop_upstream_fragment_ids = ?drop_upstream_fragment_ids,
1127 "skip dropping upstream sink fragments for non-existing target fragment"
1128 );
1129 return;
1130 }
1131 {
1132 {
1133 {
1134 {
1135 let (info, _) = self.fragment_mut(fragment_id);
1136 let mut removed = false;
1137 visit_stream_node_mut(&mut info.nodes, |node| {
1138 if let NodeBody::UpstreamSinkUnion(u) = node {
1139 if cfg!(debug_assertions) {
1140 let current_upstream_fragment_ids = u
1141 .init_upstreams
1142 .iter()
1143 .map(|upstream| upstream.upstream_fragment_id)
1144 .collect::<HashSet<FragmentId>>();
1145 for drop_fragment_id in drop_upstream_fragment_ids {
1146 if !current_upstream_fragment_ids.contains(drop_fragment_id)
1147 {
1148 panic!(
1149 "non-existing upstream fragment to drop: {:?} {:?} {:?}",
1150 u, drop_upstream_fragment_ids, drop_fragment_id
1151 );
1152 }
1153 }
1154 }
1155 u.init_upstreams.retain(|upstream| {
1156 !drop_upstream_fragment_ids
1157 .contains(&upstream.upstream_fragment_id)
1158 });
1159 removed = true;
1160 }
1161 });
1162 assert!(removed, "should remove upstream from UpstreamSinkUnion");
1163 }
1164 }
1165 }
1166 }
1167 }
1168
1169 pub(crate) fn pre_apply_throttle(
1171 &mut self,
1172 config: &mut ThrottleConfigMap,
1173 ) -> Option<Mutation> {
1174 extract_throttle_config(config, |fragment_id, stream_node| {
1175 if !self.fragment_location.contains_key(&fragment_id) {
1176 return false;
1177 }
1178 self.fragment_mut(fragment_id).0.nodes = stream_node.clone();
1179 true
1180 })
1181 }
1182
1183 pub(crate) fn pre_apply_split_assignments(
1185 &mut self,
1186 assignments: impl IntoIterator<Item = (FragmentId, HashMap<ActorId, Vec<SplitImpl>>)>,
1187 ) {
1188 {
1189 let shared_infos = self.shared_actor_infos.clone();
1190 let mut shared_actor_writer = shared_infos.start_writer(self.database_id);
1191 {
1192 {
1193 for (fragment_id, actor_splits) in assignments {
1194 let (info, job_id) = self.fragment_mut(fragment_id);
1195 let actors = &mut info.actors;
1196 for (actor_id, splits) in actor_splits {
1197 actors.get_mut(&actor_id).expect("should exist").splits = splits;
1198 }
1199 shared_actor_writer.upsert([(&*info, job_id)]);
1200 }
1201 }
1202 }
1203 shared_actor_writer.finish();
1204 }
1205 }
1206
1207 pub(super) fn build_edge(
1208 &self,
1209 info: Option<(&CreateStreamingJobCommandInfo, bool)>,
1210 replace_job: Option<&ReplaceStreamJobPlan>,
1211 new_upstream_sink: Option<&UpstreamSinkInfo>,
1212 control_stream_manager: &ControlStreamManager,
1213 stream_actors: &HashMap<FragmentId, Vec<StreamActor>>,
1214 actor_location: &HashMap<ActorId, WorkerId>,
1215 ) -> MetaResult<(FragmentEdgeBuildResult, ActorNewNoShuffle)> {
1216 let existing_fragment_ids = info
1224 .into_iter()
1225 .flat_map(|(info, _)| info.upstream_fragment_downstreams.keys())
1226 .chain(replace_job.into_iter().flat_map(|replace_job| {
1227 replace_job
1228 .upstream_fragment_downstreams
1229 .keys()
1230 .filter(|fragment_id| {
1231 info.map(|(info, _)| {
1232 !info
1233 .stream_job_fragments
1234 .fragments
1235 .contains_key(*fragment_id)
1236 })
1237 .unwrap_or(true)
1238 })
1239 .chain(replace_job.replace_upstream.keys())
1240 }))
1241 .chain(
1242 new_upstream_sink
1243 .into_iter()
1244 .map(|ctx| &ctx.new_sink_downstream.downstream_fragment_id),
1245 )
1246 .cloned();
1247 let new_fragments = info
1249 .into_iter()
1250 .flat_map(|(info, is_snapshot_backfill)| {
1251 let partial_graph_id = to_partial_graph_id(
1252 self.database_id,
1253 is_snapshot_backfill.then_some(info.streaming_job.id()),
1254 );
1255 info.stream_job_fragments
1256 .fragments
1257 .values()
1258 .map(move |fragment| (partial_graph_id, fragment))
1259 })
1260 .chain(replace_job.into_iter().flat_map(|replace_job| {
1261 replace_job
1262 .new_fragments
1263 .fragments
1264 .values()
1265 .chain(
1266 replace_job
1267 .auto_refresh_schema_sinks
1268 .as_ref()
1269 .into_iter()
1270 .flat_map(move |sinks| sinks.iter().map(|sink| &sink.new_fragment)),
1271 )
1272 .map(|fragment| {
1273 (
1274 to_partial_graph_id(self.database_id, None),
1276 fragment,
1277 )
1278 })
1279 }));
1280
1281 let database_partial_graph_id = to_partial_graph_id(self.database_id, None);
1282 let mut builder = FragmentEdgeBuilder::new()
1283 .add_existing_fragments(
1284 existing_fragment_ids.map(|fragment_id| self.fragment(fragment_id)),
1285 database_partial_graph_id,
1286 control_stream_manager,
1287 )
1288 .add_new_logical_fragments(
1289 new_fragments,
1290 stream_actors,
1291 actor_location,
1292 control_stream_manager,
1293 )
1294 .finish_fragments();
1295 if let Some((info, _)) = info {
1296 builder = builder
1297 .add_relations(&info.upstream_fragment_downstreams)?
1298 .add_relations(&info.stream_job_fragments.downstreams)?;
1299 }
1300 if let Some(replace_job) = replace_job {
1301 builder = builder
1302 .add_relations(&replace_job.upstream_fragment_downstreams)?
1303 .add_relations(&replace_job.new_fragments.downstreams)?;
1304 }
1305 if let Some(new_upstream_sink) = new_upstream_sink {
1306 let sink_fragment_id = new_upstream_sink.sink_fragment_id;
1307 let new_sink_downstream = &new_upstream_sink.new_sink_downstream;
1308 builder = builder.add_edge(sink_fragment_id, new_sink_downstream)?;
1309 }
1310 if let Some(replace_job) = replace_job {
1311 for (fragment_id, fragment_replacement) in &replace_job.replace_upstream {
1312 for (original_upstream_fragment_id, new_upstream_fragment_id) in
1313 fragment_replacement
1314 {
1315 builder = builder.replace_upstream(
1316 *fragment_id,
1317 *original_upstream_fragment_id,
1318 *new_upstream_fragment_id,
1319 );
1320 }
1321 }
1322 }
1323 Ok(builder.build())
1324 }
1325
1326 pub(crate) fn post_apply_reschedules(
1328 &mut self,
1329 reschedules: impl IntoIterator<Item = (FragmentId, HashSet<ActorId>)>,
1330 ) {
1331 let inner = self.shared_actor_infos.clone();
1332 let mut shared_actor_writer = inner.start_writer(self.database_id);
1333 {
1334 {
1335 {
1336 for (fragment_id, to_remove) in reschedules {
1337 let job_id = self.fragment_location[&fragment_id];
1338 let info = self
1339 .jobs
1340 .get_mut(&job_id)
1341 .expect("should exist")
1342 .fragment_infos
1343 .get_mut(&fragment_id)
1344 .expect("should exist");
1345 for actor_id in to_remove {
1346 assert!(info.actors.remove(&actor_id).is_some());
1347 }
1348 shared_actor_writer.upsert([(&*info, job_id)]);
1349 }
1350 }
1351 }
1352 }
1353 shared_actor_writer.finish();
1354 }
1355
1356 pub(crate) fn post_apply_remove_fragments(
1358 &mut self,
1359 fragment_ids: impl IntoIterator<Item = FragmentId>,
1360 ) {
1361 let inner = self.shared_actor_infos.clone();
1362 let mut shared_actor_writer = inner.start_writer(self.database_id);
1363 {
1364 {
1365 {
1366 for fragment_id in fragment_ids {
1367 let job_id = self
1368 .fragment_location
1369 .remove(&fragment_id)
1370 .expect("should exist");
1371 let job = self.jobs.get_mut(&job_id).expect("should exist");
1372 let fragment = job
1373 .fragment_infos
1374 .remove(&fragment_id)
1375 .expect("should exist");
1376 shared_actor_writer.remove(&fragment);
1377 if job.fragment_infos.is_empty() {
1378 self.jobs.remove(&job_id).expect("should exist");
1379 }
1380 }
1381 }
1382 }
1383 }
1384 shared_actor_writer.finish();
1385 }
1386
1387 pub(crate) fn post_apply_remove_job(
1388 &mut self,
1389 job_id: JobId,
1390 ) -> Option<InflightStreamingJobInfo> {
1391 let job = self.jobs.remove(&job_id)?;
1392 let inner = self.shared_actor_infos.clone();
1393 let mut shared_actor_writer = inner.start_writer(self.database_id);
1394 for (fragment_id, fragment) in &job.fragment_infos {
1395 self.fragment_location
1396 .remove(fragment_id)
1397 .expect("should exist");
1398 shared_actor_writer.remove(fragment);
1399 }
1400 shared_actor_writer.finish();
1401 Some(job)
1402 }
1403}
1404
1405impl InflightFragmentInfo {
1406 pub(crate) fn actor_ids_to_collect(
1408 infos: impl IntoIterator<Item = &Self>,
1409 ) -> HashMap<WorkerId, HashSet<ActorId>> {
1410 let mut ret: HashMap<_, HashSet<_>> = HashMap::new();
1411 for (actor_id, actor) in infos.into_iter().flat_map(|info| info.actors.iter()) {
1412 assert!(
1413 ret.entry(actor.worker_id)
1414 .or_default()
1415 .insert(*actor_id as _)
1416 )
1417 }
1418 ret
1419 }
1420
1421 pub fn existing_table_ids<'a>(
1422 infos: impl IntoIterator<Item = &'a Self> + 'a,
1423 ) -> impl Iterator<Item = TableId> + 'a {
1424 infos
1425 .into_iter()
1426 .flat_map(|info| info.state_table_ids.iter().cloned())
1427 }
1428
1429 pub fn workers<'a>(
1430 infos: impl IntoIterator<Item = &'a Self> + 'a,
1431 ) -> impl Iterator<Item = WorkerId> + 'a {
1432 infos
1433 .into_iter()
1434 .flat_map(|fragment| fragment.actors.values().map(|actor| actor.worker_id))
1435 }
1436
1437 pub fn contains_worker<'a>(
1438 infos: impl IntoIterator<Item = &'a Self> + 'a,
1439 worker_id: WorkerId,
1440 ) -> bool {
1441 Self::workers(infos).any(|existing_worker_id| existing_worker_id == worker_id)
1442 }
1443}
1444
1445impl InflightDatabaseInfo {
1446 pub fn contains_worker(&self, worker_id: WorkerId) -> bool {
1447 InflightFragmentInfo::contains_worker(self.fragment_infos(), worker_id)
1448 }
1449
1450 pub fn existing_table_ids(&self) -> impl Iterator<Item = TableId> + '_ {
1451 InflightFragmentInfo::existing_table_ids(self.fragment_infos())
1452 }
1453}