1use std::collections::hash_map::Entry;
16use std::collections::{HashMap, HashSet};
17use std::fmt::{Display, Formatter};
18
19use itertools::Itertools;
20use risingwave_common::bitmap::Bitmap;
21use risingwave_common::catalog::{DatabaseId, TableId};
22use risingwave_common::hash::{ActorMapping, VnodeCountCompat};
23use risingwave_common::id::{JobId, SinkId, SourceId};
24use risingwave_common::must_match;
25use risingwave_common::types::Timestamptz;
26use risingwave_common::util::epoch::Epoch;
27use risingwave_connector::source::{CdcTableSnapshotSplitRaw, SplitImpl};
28use risingwave_hummock_sdk::change_log::build_table_change_log_delta;
29use risingwave_hummock_sdk::vector_index::VectorIndexDelta;
30use risingwave_meta_model::{DispatcherType, WorkerId, fragment_relation, streaming_job};
31use risingwave_pb::catalog::CreateType;
32use risingwave_pb::common::PbActorInfo;
33use risingwave_pb::hummock::vector_index_delta::PbVectorIndexInit;
34use risingwave_pb::plan_common::{ColumnCatalog as PbColumnCatalog, PbField};
35use risingwave_pb::source::{
36 ConnectorSplit, ConnectorSplits, PbCdcTableSnapshotSplits,
37 PbCdcTableSnapshotSplitsWithGeneration,
38};
39use risingwave_pb::stream_plan::add_mutation::PbNewUpstreamSink;
40use risingwave_pb::stream_plan::barrier::BarrierKind as PbBarrierKind;
41use risingwave_pb::stream_plan::barrier_mutation::Mutation;
42use risingwave_pb::stream_plan::connector_props_change_mutation::ConnectorPropsInfo;
43use risingwave_pb::stream_plan::sink_schema_change::Op as PbSinkSchemaChangeOp;
44use risingwave_pb::stream_plan::throttle_mutation::ThrottleConfig;
45use risingwave_pb::stream_plan::update_mutation::{DispatcherUpdate, MergeUpdate};
46use risingwave_pb::stream_plan::{
47 AddMutation, ConnectorPropsChangeMutation, Dispatcher, Dispatchers, DropSubscriptionsMutation,
48 ListFinishMutation, LoadFinishMutation, PauseMutation, PbSinkAddColumnsOp, PbSinkDropColumnsOp,
49 PbSinkSchemaChange, PbStreamNode, PbUpstreamSinkInfo, ResumeMutation,
50 SourceChangeSplitMutation, StartFragmentBackfillMutation, StopMutation,
51 SubscriptionUpstreamInfo, ThrottleMutation, UpdateMutation,
52};
53use risingwave_pb::stream_service::BarrierCompleteResponse;
54use tracing::warn;
55
56use super::info::InflightDatabaseInfo;
57use crate::barrier::backfill_order_control::get_nodes_with_backfill_dependencies;
58use crate::barrier::complete_task::CompleteBarrierTask;
59use crate::barrier::edge_builder::FragmentEdgeBuildResult;
60use crate::barrier::info::BarrierInfo;
61use crate::barrier::partial_graph::PartialGraphBarrierInfo;
62use crate::barrier::rpc::{ControlStreamManager, to_partial_graph_id};
63use crate::barrier::utils::{collect_new_vector_index_info, collect_resp_info};
64use crate::controller::fragment::{InflightActorInfo, InflightFragmentInfo};
65use crate::controller::scale::LoadedFragmentContext;
66use crate::controller::utils::StreamingJobExtraInfo;
67use crate::hummock::NewTableFragmentInfo;
68use crate::manager::{StreamingJob, StreamingJobType};
69use crate::model::{
70 ActorId, ActorUpstreams, DispatcherId, FragmentActorDispatchers, FragmentDownstreamRelation,
71 FragmentId, FragmentReplaceUpstream, StreamActor, StreamActorWithDispatchers,
72 StreamJobActorsToCreate, StreamJobFragments, StreamJobFragmentsToCreate, SubscriptionId,
73};
74use crate::stream::{
75 AutoRefreshSchemaSinkContext, ConnectorPropsChange, ExtendedFragmentBackfillOrder,
76 ReplaceJobSplitPlan, SourceSplitAssignment, SplitAssignment, SplitState, UpstreamSinkInfo,
77 build_actor_connector_splits,
78};
79use crate::{MetaError, MetaResult};
80
81#[derive(Debug, Clone)]
84pub struct Reschedule {
85 pub added_actors: HashMap<WorkerId, Vec<ActorId>>,
87
88 pub removed_actors: HashSet<ActorId>,
90
91 pub vnode_bitmap_updates: HashMap<ActorId, Bitmap>,
93
94 pub upstream_fragment_dispatcher_ids: Vec<(FragmentId, DispatcherId)>,
96 pub upstream_dispatcher_mapping: Option<ActorMapping>,
101
102 pub downstream_fragment_ids: Vec<FragmentId>,
104
105 pub actor_splits: HashMap<ActorId, Vec<SplitImpl>>,
109
110 pub newly_created_actors: HashMap<ActorId, (StreamActorWithDispatchers, WorkerId)>,
111}
112
113#[derive(Debug, Clone)]
114pub struct ReschedulePlan {
115 pub reschedules: HashMap<FragmentId, Reschedule>,
116 pub fragment_actors: HashMap<FragmentId, HashSet<ActorId>>,
119}
120
121#[derive(Debug, Clone)]
123pub struct RescheduleContext {
124 pub loaded: LoadedFragmentContext,
125 pub job_extra_info: HashMap<JobId, StreamingJobExtraInfo>,
126 pub upstream_fragments: HashMap<FragmentId, HashMap<FragmentId, DispatcherType>>,
127 pub downstream_fragments: HashMap<FragmentId, HashMap<FragmentId, DispatcherType>>,
128 pub downstream_relations: HashMap<(FragmentId, FragmentId), fragment_relation::Model>,
129}
130
131impl RescheduleContext {
132 pub fn empty() -> Self {
133 Self {
134 loaded: LoadedFragmentContext::default(),
135 job_extra_info: HashMap::new(),
136 upstream_fragments: HashMap::new(),
137 downstream_fragments: HashMap::new(),
138 downstream_relations: HashMap::new(),
139 }
140 }
141
142 pub fn is_empty(&self) -> bool {
143 self.loaded.is_empty()
144 }
145
146 pub fn for_database(&self, database_id: DatabaseId) -> Option<Self> {
147 let loaded = self.loaded.for_database(database_id)?;
148 let job_ids: HashSet<JobId> = loaded.job_map.keys().copied().collect();
149 let fragment_ids: HashSet<FragmentId> = loaded
152 .job_fragments
153 .values()
154 .flat_map(|fragments| fragments.keys().copied())
155 .collect();
156
157 let job_extra_info = self
158 .job_extra_info
159 .iter()
160 .filter(|(job_id, _)| job_ids.contains(*job_id))
161 .map(|(job_id, info)| (*job_id, info.clone()))
162 .collect();
163
164 let upstream_fragments = self
165 .upstream_fragments
166 .iter()
167 .filter(|(fragment_id, _)| fragment_ids.contains(*fragment_id))
168 .map(|(fragment_id, upstreams)| (*fragment_id, upstreams.clone()))
169 .collect();
170
171 let downstream_fragments = self
172 .downstream_fragments
173 .iter()
174 .filter(|(fragment_id, _)| fragment_ids.contains(*fragment_id))
175 .map(|(fragment_id, downstreams)| (*fragment_id, downstreams.clone()))
176 .collect();
177
178 let downstream_relations = self
179 .downstream_relations
180 .iter()
181 .filter(|((source_fragment_id, _), _)| fragment_ids.contains(source_fragment_id))
185 .map(|(key, relation)| (*key, relation.clone()))
186 .collect();
187
188 Some(Self {
189 loaded,
190 job_extra_info,
191 upstream_fragments,
192 downstream_fragments,
193 downstream_relations,
194 })
195 }
196
197 pub fn into_database_contexts(self) -> HashMap<DatabaseId, Self> {
200 let Self {
201 loaded,
202 job_extra_info,
203 upstream_fragments,
204 downstream_fragments,
205 downstream_relations,
206 } = self;
207
208 let mut contexts: HashMap<_, _> = loaded
209 .into_database_contexts()
210 .into_iter()
211 .map(|(database_id, loaded)| {
212 (
213 database_id,
214 Self {
215 loaded,
216 job_extra_info: HashMap::new(),
217 upstream_fragments: HashMap::new(),
218 downstream_fragments: HashMap::new(),
219 downstream_relations: HashMap::new(),
220 },
221 )
222 })
223 .collect();
224
225 if contexts.is_empty() {
226 return contexts;
227 }
228
229 let mut job_databases = HashMap::new();
230 let mut fragment_databases = HashMap::new();
231 for (&database_id, context) in &contexts {
232 for job_id in context.loaded.job_map.keys().copied() {
233 job_databases.insert(job_id, database_id);
234 }
235 for fragment_id in context
236 .loaded
237 .job_fragments
238 .values()
239 .flat_map(|fragments| fragments.keys().copied())
240 {
241 fragment_databases.insert(fragment_id, database_id);
242 }
243 }
244
245 for (job_id, info) in job_extra_info {
246 if let Some(database_id) = job_databases.get(&job_id).copied() {
247 contexts
248 .get_mut(&database_id)
249 .expect("database context should exist for job")
250 .job_extra_info
251 .insert(job_id, info);
252 }
253 }
254
255 for (fragment_id, upstreams) in upstream_fragments {
256 if let Some(database_id) = fragment_databases.get(&fragment_id).copied() {
257 contexts
258 .get_mut(&database_id)
259 .expect("database context should exist for fragment")
260 .upstream_fragments
261 .insert(fragment_id, upstreams);
262 }
263 }
264
265 for (fragment_id, downstreams) in downstream_fragments {
266 if let Some(database_id) = fragment_databases.get(&fragment_id).copied() {
267 contexts
268 .get_mut(&database_id)
269 .expect("database context should exist for fragment")
270 .downstream_fragments
271 .insert(fragment_id, downstreams);
272 }
273 }
274
275 for ((source_fragment_id, target_fragment_id), relation) in downstream_relations {
276 if let Some(database_id) = fragment_databases.get(&source_fragment_id).copied() {
279 contexts
280 .get_mut(&database_id)
281 .expect("database context should exist for relation source")
282 .downstream_relations
283 .insert((source_fragment_id, target_fragment_id), relation);
284 }
285 }
286
287 contexts
288 }
289}
290
291#[derive(Debug, Clone)]
298pub struct ReplaceStreamJobPlan {
299 pub old_fragments: StreamJobFragments,
300 pub new_fragments: StreamJobFragmentsToCreate,
301 pub database_resource_group: String,
303 pub replace_upstream: FragmentReplaceUpstream,
306 pub upstream_fragment_downstreams: FragmentDownstreamRelation,
307 pub split_plan: ReplaceJobSplitPlan,
310 pub streaming_job: StreamingJob,
312 pub streaming_job_model: streaming_job::Model,
314 pub tmp_id: JobId,
316 pub to_drop_state_table_ids: Vec<TableId>,
318 pub auto_refresh_schema_sinks: Option<Vec<AutoRefreshSchemaSinkContext>>,
319}
320
321impl ReplaceStreamJobPlan {
322 pub fn fragment_replacements(&self) -> HashMap<FragmentId, FragmentId> {
324 let mut fragment_replacements = HashMap::new();
325 for (upstream_fragment_id, new_upstream_fragment_id) in
326 self.replace_upstream.values().flatten()
327 {
328 {
329 let r =
330 fragment_replacements.insert(*upstream_fragment_id, *new_upstream_fragment_id);
331 if let Some(r) = r {
332 assert_eq!(
333 *new_upstream_fragment_id, r,
334 "one fragment is replaced by multiple fragments"
335 );
336 }
337 }
338 }
339 fragment_replacements
340 }
341}
342
343#[derive(educe::Educe, Clone)]
344#[educe(Debug)]
345pub struct CreateStreamingJobCommandInfo {
346 #[educe(Debug(ignore))]
347 pub stream_job_fragments: StreamJobFragmentsToCreate,
348 pub upstream_fragment_downstreams: FragmentDownstreamRelation,
349 pub database_resource_group: String,
351 pub init_split_assignment: SourceSplitAssignment,
353 pub definition: String,
354 pub job_type: StreamingJobType,
355 pub create_type: CreateType,
356 pub streaming_job: StreamingJob,
357 pub fragment_backfill_ordering: ExtendedFragmentBackfillOrder,
358 pub cdc_table_snapshot_splits: Option<Vec<CdcTableSnapshotSplitRaw>>,
359 pub locality_fragment_state_table_mapping: HashMap<FragmentId, Vec<TableId>>,
360 pub is_serverless: bool,
361 pub streaming_job_model: streaming_job::Model,
363 pub replace_sink: Option<SinkId>,
365 pub refresh_interval_sec: Option<u64>,
367}
368
369impl StreamJobFragments {
370 pub(super) fn new_fragment_info<'a>(
373 &'a self,
374 stream_actors: &'a HashMap<FragmentId, Vec<StreamActor>>,
375 actor_location: &'a HashMap<ActorId, WorkerId>,
376 assignment: &'a SplitAssignment,
377 ) -> impl Iterator<Item = (FragmentId, InflightFragmentInfo)> + 'a {
378 self.fragments.values().map(|fragment| {
379 (
380 fragment.fragment_id,
381 InflightFragmentInfo {
382 fragment_id: fragment.fragment_id,
383 distribution_type: fragment.distribution_type.into(),
384 fragment_type_mask: fragment.fragment_type_mask,
385 vnode_count: fragment.vnode_count(),
386 nodes: fragment.nodes.clone(),
387 actors: stream_actors
388 .get(&fragment.fragment_id)
389 .into_iter()
390 .flatten()
391 .map(|actor| {
392 (
393 actor.actor_id,
394 InflightActorInfo {
395 worker_id: actor_location[&actor.actor_id],
396 vnode_bitmap: actor.vnode_bitmap.clone(),
397 splits: assignment
398 .get(&fragment.fragment_id)
399 .and_then(|s| s.get(&actor.actor_id))
400 .cloned()
401 .unwrap_or_default(),
402 },
403 )
404 })
405 .collect(),
406 state_table_ids: fragment.state_table_ids.iter().copied().collect(),
407 },
408 )
409 })
410 }
411}
412
413pub type TableLogEpochs = Vec<(Vec<u64>, u64)>;
414pub type UpstreamTableLogEpochs = HashMap<TableId, TableLogEpochs>;
415pub type SinceTimestampResolvedEpoch = (u64, TableLogEpochs);
416
417#[derive(Debug, Clone)]
418pub struct SnapshotBackfillInfo {
419 pub upstream_mv_table_id_to_backfill_epoch: HashMap<TableId, Option<u64>>,
423}
424
425#[derive(Debug, Clone)]
426pub struct SinceEpochInfo {
427 pub provided_since_epoch: u64,
428 pub resolved: Option<SinceTimestampResolvedEpoch>,
429}
430
431#[derive(Debug, Clone)]
432pub struct BatchRefreshInfo {
433 pub snapshot_backfill_info: SnapshotBackfillInfo,
434 pub refresh_interval_sec: u64,
435}
436
437#[derive(Debug, Clone)]
438pub enum CreateStreamingJobType {
439 Normal,
440 SinkIntoTable(UpstreamSinkInfo),
441 SnapshotBackfill {
442 snapshot_backfill_info: SnapshotBackfillInfo,
443 since_epoch: Option<SinceEpochInfo>,
444 },
445 BatchRefresh(BatchRefreshInfo),
446}
447
448#[derive(Debug)]
453pub enum Command {
454 Flush,
457
458 Pause,
461
462 Resume,
466
467 DropStreamingJobs {
475 streaming_job_ids: HashSet<JobId>,
476 unregistered_state_table_ids: HashSet<TableId>,
478 dropped_sink_fragment_by_targets: HashMap<FragmentId, Vec<FragmentId>>,
480 },
481
482 CreateStreamingJob {
492 info: CreateStreamingJobCommandInfo,
493 job_type: CreateStreamingJobType,
494 cross_db_snapshot_backfill_info: SnapshotBackfillInfo,
495 },
496
497 RescheduleIntent {
499 context: RescheduleContext,
500 reschedule_plan: Option<ReschedulePlan>,
505 },
506
507 ReplaceStreamJob(ReplaceStreamJobPlan),
514
515 SourceChangeSplit(SplitState),
518
519 Throttle {
522 jobs: HashSet<JobId>,
523 config: HashMap<FragmentId, (ThrottleConfig, PbStreamNode)>,
524 },
525
526 CreateSubscription {
529 subscription_id: SubscriptionId,
530 upstream_mv_table_id: TableId,
531 retention_second: u64,
532 },
533
534 DropSubscription {
538 subscription_id: SubscriptionId,
539 upstream_mv_table_id: TableId,
540 },
541
542 AlterSubscriptionRetention {
544 subscription_id: SubscriptionId,
545 upstream_mv_table_id: TableId,
546 retention_second: u64,
547 },
548
549 ConnectorPropsChange(ConnectorPropsChange),
550
551 Refresh {
554 table_id: TableId,
555 associated_source_id: SourceId,
556 },
557 ListFinish {
558 table_id: TableId,
559 associated_source_id: SourceId,
560 },
561 LoadFinish {
562 table_id: TableId,
563 associated_source_id: SourceId,
564 },
565
566 ResetSource {
569 source_id: SourceId,
570 },
571
572 ResumeBackfill {
575 target: ResumeBackfillTarget,
576 },
577
578 InjectSourceOffsets {
582 source_id: SourceId,
583 split_offsets: HashMap<String, String>,
585 },
586}
587
588#[derive(Debug, Clone, Copy)]
589pub enum ResumeBackfillTarget {
590 Job(JobId),
591 Fragment(FragmentId),
592}
593
594impl std::fmt::Display for Command {
596 fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
597 match self {
598 Command::Flush => write!(f, "Flush"),
599 Command::Pause => write!(f, "Pause"),
600 Command::Resume => write!(f, "Resume"),
601 Command::DropStreamingJobs {
602 streaming_job_ids, ..
603 } => {
604 write!(
605 f,
606 "DropStreamingJobs: {}",
607 streaming_job_ids.iter().sorted().join(", ")
608 )
609 }
610 Command::CreateStreamingJob { info, .. } => {
611 write!(f, "CreateStreamingJob: {}", info.streaming_job)
612 }
613 Command::RescheduleIntent {
614 reschedule_plan, ..
615 } => {
616 if reschedule_plan.is_some() {
617 write!(f, "RescheduleIntent(planned)")
618 } else {
619 write!(f, "RescheduleIntent")
620 }
621 }
622 Command::ReplaceStreamJob(plan) => {
623 write!(f, "ReplaceStreamJob: {}", plan.streaming_job)
624 }
625 Command::SourceChangeSplit { .. } => write!(f, "SourceChangeSplit"),
626 Command::Throttle { .. } => write!(f, "Throttle"),
627 Command::CreateSubscription {
628 subscription_id, ..
629 } => write!(f, "CreateSubscription: {subscription_id}"),
630 Command::DropSubscription {
631 subscription_id, ..
632 } => write!(f, "DropSubscription: {subscription_id}"),
633 Command::AlterSubscriptionRetention {
634 subscription_id,
635 retention_second,
636 ..
637 } => write!(
638 f,
639 "AlterSubscriptionRetention: {subscription_id} -> {retention_second}"
640 ),
641 Command::ConnectorPropsChange(_) => write!(f, "ConnectorPropsChange"),
642 Command::Refresh {
643 table_id,
644 associated_source_id,
645 } => write!(
646 f,
647 "Refresh: {} (source: {})",
648 table_id, associated_source_id
649 ),
650 Command::ListFinish {
651 table_id,
652 associated_source_id,
653 } => write!(
654 f,
655 "ListFinish: {} (source: {})",
656 table_id, associated_source_id
657 ),
658 Command::LoadFinish {
659 table_id,
660 associated_source_id,
661 } => write!(
662 f,
663 "LoadFinish: {} (source: {})",
664 table_id, associated_source_id
665 ),
666 Command::ResetSource { source_id } => write!(f, "ResetSource: {source_id}"),
667 Command::ResumeBackfill { target } => match target {
668 ResumeBackfillTarget::Job(job_id) => {
669 write!(f, "ResumeBackfill: job={job_id}")
670 }
671 ResumeBackfillTarget::Fragment(fragment_id) => {
672 write!(f, "ResumeBackfill: fragment={fragment_id}")
673 }
674 },
675 Command::InjectSourceOffsets {
676 source_id,
677 split_offsets,
678 } => write!(
679 f,
680 "InjectSourceOffsets: {} ({} splits)",
681 source_id,
682 split_offsets.len()
683 ),
684 }
685 }
686}
687
688impl Command {
689 pub fn pause() -> Self {
690 Self::Pause
691 }
692
693 pub fn resume() -> Self {
694 Self::Resume
695 }
696
697 pub fn need_checkpoint(&self) -> bool {
698 !matches!(self, Command::Resume)
700 }
701}
702
703#[derive(Debug)]
704pub enum PostCollectCommand {
705 Command(String),
706 DropStreamingJobs,
707 CreateStreamingJob {
708 info: CreateStreamingJobCommandInfo,
709 job_type: CreateStreamingJobType,
710 cross_db_snapshot_backfill_info: SnapshotBackfillInfo,
711 resolved_split_assignment: SplitAssignment,
712 },
713 Reschedule {
714 reschedules: HashMap<FragmentId, Reschedule>,
715 },
716 ReplaceStreamJob {
717 plan: ReplaceStreamJobPlan,
718 resolved_split_assignment: SplitAssignment,
719 },
720 SourceChangeSplit {
721 split_assignment: SplitAssignment,
722 },
723 CreateSubscription {
724 subscription_id: SubscriptionId,
725 },
726 ConnectorPropsChange(ConnectorPropsChange),
727 ResumeBackfill {
728 target: ResumeBackfillTarget,
729 },
730}
731
732impl PostCollectCommand {
733 pub fn barrier() -> Self {
734 PostCollectCommand::Command("barrier".to_owned())
735 }
736
737 pub fn should_checkpoint(&self) -> bool {
738 match self {
739 PostCollectCommand::DropStreamingJobs
740 | PostCollectCommand::CreateStreamingJob { .. }
741 | PostCollectCommand::Reschedule { .. }
742 | PostCollectCommand::ReplaceStreamJob { .. }
743 | PostCollectCommand::SourceChangeSplit { .. }
744 | PostCollectCommand::CreateSubscription { .. }
745 | PostCollectCommand::ConnectorPropsChange(_)
746 | PostCollectCommand::ResumeBackfill { .. } => true,
747 PostCollectCommand::Command(_) => false,
748 }
749 }
750
751 pub fn command_name(&self) -> &str {
752 match self {
753 PostCollectCommand::Command(name) => name.as_str(),
754 PostCollectCommand::DropStreamingJobs => "DropStreamingJobs",
755 PostCollectCommand::CreateStreamingJob { .. } => "CreateStreamingJob",
756 PostCollectCommand::Reschedule { .. } => "Reschedule",
757 PostCollectCommand::ReplaceStreamJob { .. } => "ReplaceStreamJob",
758 PostCollectCommand::SourceChangeSplit { .. } => "SourceChangeSplit",
759 PostCollectCommand::CreateSubscription { .. } => "CreateSubscription",
760 PostCollectCommand::ConnectorPropsChange(_) => "ConnectorPropsChange",
761 PostCollectCommand::ResumeBackfill { .. } => "ResumeBackfill",
762 }
763 }
764}
765
766impl Display for PostCollectCommand {
767 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
768 f.write_str(self.command_name())
769 }
770}
771
772#[derive(Debug, Clone, PartialEq, Eq)]
773pub enum BarrierKind {
774 Initial,
775 Barrier,
776 Checkpoint(Vec<u64>),
778}
779
780impl BarrierKind {
781 pub fn to_protobuf(&self) -> PbBarrierKind {
782 match self {
783 BarrierKind::Initial => PbBarrierKind::Initial,
784 BarrierKind::Barrier => PbBarrierKind::Barrier,
785 BarrierKind::Checkpoint(_) => PbBarrierKind::Checkpoint,
786 }
787 }
788
789 pub fn is_checkpoint(&self) -> bool {
790 matches!(self, BarrierKind::Checkpoint(_))
791 }
792
793 pub fn is_initial(&self) -> bool {
794 matches!(self, BarrierKind::Initial)
795 }
796
797 pub fn as_str_name(&self) -> &'static str {
798 match self {
799 BarrierKind::Initial => "Initial",
800 BarrierKind::Barrier => "Barrier",
801 BarrierKind::Checkpoint(_) => "Checkpoint",
802 }
803 }
804}
805
806fn sink_original_schema_fields(columns: &[PbColumnCatalog]) -> Vec<PbField> {
807 columns
808 .iter()
809 .filter(|col| !col.is_hidden)
810 .map(|col| {
811 let column_desc = col
812 .column_desc
813 .as_ref()
814 .expect("sink column catalog should have a column descriptor");
815 PbField {
816 data_type: Some(
817 column_desc
818 .column_type
819 .as_ref()
820 .expect("sink column descriptor should have a column type")
821 .clone(),
822 ),
823 name: column_desc.name.clone(),
824 }
825 })
826 .collect()
827}
828
829impl BarrierInfo {
830 fn get_truncate_epoch(&self, retention_second: u64) -> Epoch {
831 let Some(truncate_timestamptz) = Timestamptz::from_secs(
832 self.prev_epoch.value().as_timestamptz().timestamp() - retention_second as i64,
833 ) else {
834 warn!(retention_second, prev_epoch = ?self.prev_epoch.value(), "invalid retention second value");
835 return self.prev_epoch.value();
836 };
837 Epoch::from_unix_millis(truncate_timestamptz.timestamp_millis() as u64)
838 }
839}
840
841impl Command {
842 pub(super) fn collect_commit_epoch_info(
843 database_info: &InflightDatabaseInfo,
844 barrier_info: &PartialGraphBarrierInfo,
845 task: &mut CompleteBarrierTask,
846 resps: Vec<BarrierCompleteResponse>,
847 backfill_pinned_log_epoch: HashMap<JobId, (u64, HashSet<TableId>)>,
848 ) {
849 let (
850 sst_to_context,
851 synced_ssts,
852 new_table_watermarks,
853 old_value_ssts,
854 vector_index_adds,
855 truncate_tables,
856 iceberg_pk_index_sink_metadata,
857 ) = collect_resp_info(resps);
858
859 let new_table_fragment_infos = match &barrier_info.post_collect_command {
860 PostCollectCommand::CreateStreamingJob { info, job_type, .. } => {
861 assert!(!matches!(
862 job_type,
863 CreateStreamingJobType::SnapshotBackfill { .. }
864 | CreateStreamingJobType::BatchRefresh(_)
865 ));
866 let table_fragments = &info.stream_job_fragments;
867 let mut table_ids: HashSet<_> =
868 table_fragments.internal_table_ids().into_iter().collect();
869 if let Some(mv_table_id) = table_fragments.mv_table_id() {
870 table_ids.insert(mv_table_id);
871 }
872
873 vec![NewTableFragmentInfo { table_ids }]
874 }
875 _ => vec![],
876 };
877
878 let mut mv_log_store_truncate_epoch = HashMap::new();
879 let mut update_truncate_epoch =
881 |table_id: TableId, truncate_epoch| match mv_log_store_truncate_epoch.entry(table_id) {
882 Entry::Occupied(mut entry) => {
883 let prev_truncate_epoch = entry.get_mut();
884 if truncate_epoch < *prev_truncate_epoch {
885 *prev_truncate_epoch = truncate_epoch;
886 }
887 }
888 Entry::Vacant(entry) => {
889 entry.insert(truncate_epoch);
890 }
891 };
892 for (mv_table_id, max_retention) in database_info.max_subscription_retention() {
893 let truncate_epoch = barrier_info
894 .barrier_info
895 .get_truncate_epoch(max_retention)
896 .0;
897 update_truncate_epoch(mv_table_id, truncate_epoch);
898 }
899 for (_, (backfill_epoch, upstream_mv_table_ids)) in backfill_pinned_log_epoch {
900 for mv_table_id in upstream_mv_table_ids {
901 update_truncate_epoch(mv_table_id, backfill_epoch);
902 }
903 }
904
905 let table_new_change_log = build_table_change_log_delta(
906 old_value_ssts.into_iter(),
907 synced_ssts.iter().map(|sst| &sst.sst_info),
908 must_match!(&barrier_info.barrier_info.kind, BarrierKind::Checkpoint(epochs) => epochs),
909 mv_log_store_truncate_epoch.into_iter(),
910 );
911
912 let epoch = barrier_info.barrier_info.prev_epoch();
913 let info = &mut task.commit_info;
914 for table_id in &barrier_info.table_ids_to_commit {
915 info.tables_to_commit
916 .try_insert(*table_id, epoch)
917 .expect("non duplicate");
918 }
919
920 info.sstables.extend(synced_ssts);
921 info.new_table_watermarks.extend(new_table_watermarks);
922 info.sst_to_context.extend(sst_to_context);
923 info.new_table_fragment_infos
924 .extend(new_table_fragment_infos);
925 info.change_log_delta.extend(table_new_change_log);
926 for (table_id, vector_index_adds) in vector_index_adds {
927 info.vector_index_delta
928 .try_insert(table_id, VectorIndexDelta::Adds(vector_index_adds))
929 .expect("non-duplicate");
930 }
931 if let PostCollectCommand::CreateStreamingJob { info: job_info, .. } =
932 &barrier_info.post_collect_command
933 && let Some(index_table) = collect_new_vector_index_info(job_info)
934 {
935 info.vector_index_delta
936 .try_insert(
937 index_table.id,
938 VectorIndexDelta::Init(PbVectorIndexInit {
939 info: Some(index_table.vector_index_info.unwrap()),
940 }),
941 )
942 .expect("non-duplicate");
943 }
944 info.truncate_tables.extend(truncate_tables);
945 task.iceberg_pk_index_sink_metadata
946 .extend(iceberg_pk_index_sink_metadata);
947 }
948}
949
950impl Command {
951 pub(super) fn pause_to_mutation(is_currently_paused: bool) -> Option<Mutation> {
953 {
954 {
955 if !is_currently_paused {
958 Some(Mutation::Pause(PauseMutation {}))
959 } else {
960 None
961 }
962 }
963 }
964 }
965
966 pub(super) fn resume_to_mutation(is_currently_paused: bool) -> Option<Mutation> {
968 {
969 {
970 if is_currently_paused {
972 Some(Mutation::Resume(ResumeMutation {}))
973 } else {
974 None
975 }
976 }
977 }
978 }
979
980 pub(super) fn source_change_split_to_mutation(split_assignment: &SplitAssignment) -> Mutation {
982 {
983 {
984 let mut diff = HashMap::new();
985
986 for actor_splits in split_assignment.values() {
987 diff.extend(actor_splits.clone());
988 }
989
990 Mutation::Splits(SourceChangeSplitMutation {
991 actor_splits: build_actor_connector_splits(&diff),
992 })
993 }
994 }
995 }
996
997 pub(super) fn throttle_to_mutation(
999 config: &HashMap<FragmentId, (ThrottleConfig, PbStreamNode)>,
1000 ) -> Mutation {
1001 {
1002 {
1003 Mutation::Throttle(ThrottleMutation {
1004 fragment_throttle: config
1005 .iter()
1006 .map(|(fragment_id, (throttle_config, _))| (*fragment_id, *throttle_config))
1007 .collect(),
1008 })
1009 }
1010 }
1011 }
1012
1013 pub(super) fn drop_streaming_jobs_to_mutation(
1015 actors: &Vec<ActorId>,
1016 dropped_sink_fragment_by_targets: &HashMap<FragmentId, Vec<FragmentId>>,
1017 ) -> Mutation {
1018 {
1019 Mutation::Stop(StopMutation {
1020 actors: actors.clone(),
1021 dropped_sink_fragments: dropped_sink_fragment_by_targets
1022 .values()
1023 .flatten()
1024 .cloned()
1025 .collect(),
1026 })
1027 }
1028 }
1029
1030 pub(super) fn create_streaming_job_to_mutation(
1032 info: &CreateStreamingJobCommandInfo,
1033 job_type: &CreateStreamingJobType,
1034 dropped_actors: impl IntoIterator<Item = ActorId>,
1035 is_currently_paused: bool,
1036 edges: &mut FragmentEdgeBuildResult,
1037 control_stream_manager: &ControlStreamManager,
1038 actor_cdc_table_snapshot_splits: Option<HashMap<ActorId, PbCdcTableSnapshotSplits>>,
1039 split_assignment: &SplitAssignment,
1040 stream_actors: &HashMap<FragmentId, Vec<StreamActor>>,
1041 actor_location: &HashMap<ActorId, WorkerId>,
1042 ) -> MetaResult<Mutation> {
1043 {
1044 {
1045 let CreateStreamingJobCommandInfo {
1046 stream_job_fragments,
1047 upstream_fragment_downstreams,
1048 fragment_backfill_ordering,
1049 streaming_job,
1050 ..
1051 } = info;
1052 let database_id = streaming_job.database_id();
1053 let added_actors: Vec<ActorId> = stream_actors
1054 .values()
1055 .flatten()
1056 .map(|actor| actor.actor_id)
1057 .collect();
1058 let dropped_actors = dropped_actors.into_iter().collect();
1059 let actor_splits = split_assignment
1060 .values()
1061 .flat_map(build_actor_connector_splits)
1062 .collect();
1063 let subscriptions_to_add = {
1064 if let CreateStreamingJobType::SnapshotBackfill {
1065 snapshot_backfill_info,
1066 ..
1067 }
1068 | CreateStreamingJobType::BatchRefresh(BatchRefreshInfo {
1069 snapshot_backfill_info,
1070 ..
1071 }) = job_type
1072 {
1073 snapshot_backfill_info
1074 .upstream_mv_table_id_to_backfill_epoch
1075 .keys()
1076 .map(|table_id| SubscriptionUpstreamInfo {
1077 subscriber_id: stream_job_fragments
1078 .stream_job_id()
1079 .as_subscriber_id(),
1080 upstream_mv_table_id: *table_id,
1081 })
1082 .collect()
1083 } else {
1084 Default::default()
1085 }
1086 };
1087 let backfill_nodes_to_pause: Vec<_> =
1088 get_nodes_with_backfill_dependencies(fragment_backfill_ordering)
1089 .into_iter()
1090 .collect();
1091
1092 let new_upstream_sinks =
1093 if let CreateStreamingJobType::SinkIntoTable(UpstreamSinkInfo {
1094 sink_fragment_id,
1095 sink_output_fields,
1096 project_exprs,
1097 new_sink_downstream,
1098 ..
1099 }) = job_type
1100 {
1101 let new_sink_actors = stream_actors
1102 .get(sink_fragment_id)
1103 .unwrap_or_else(|| {
1104 panic!("upstream sink fragment {sink_fragment_id} not exist")
1105 })
1106 .iter()
1107 .map(|actor| {
1108 let worker_id = actor_location[&actor.actor_id];
1109 PbActorInfo {
1110 actor_id: actor.actor_id,
1111 host: Some(control_stream_manager.host_addr(worker_id)),
1112 partial_graph_id: to_partial_graph_id(database_id, None),
1113 }
1114 });
1115 let new_upstream_sink = PbNewUpstreamSink {
1116 info: Some(PbUpstreamSinkInfo {
1117 upstream_fragment_id: *sink_fragment_id,
1118 sink_output_schema: sink_output_fields.clone(),
1119 project_exprs: project_exprs.clone(),
1120 }),
1121 upstream_actors: new_sink_actors.collect(),
1122 };
1123 HashMap::from([(
1124 new_sink_downstream.downstream_fragment_id,
1125 new_upstream_sink,
1126 )])
1127 } else {
1128 HashMap::new()
1129 };
1130
1131 let actor_cdc_table_snapshot_splits = actor_cdc_table_snapshot_splits
1132 .map(|splits| PbCdcTableSnapshotSplitsWithGeneration { splits });
1133
1134 let add_mutation = AddMutation {
1135 actor_dispatchers: edges
1136 .dispatchers
1137 .extract_if(|fragment_id, _| {
1138 upstream_fragment_downstreams.contains_key(fragment_id)
1139 })
1140 .flat_map(|(_, fragment_dispatchers)| fragment_dispatchers.into_iter())
1141 .map(|(actor_id, dispatchers)| (actor_id, Dispatchers { dispatchers }))
1142 .collect(),
1143 added_actors,
1144 actor_splits,
1145 pause: is_currently_paused,
1147 subscriptions_to_add,
1148 backfill_nodes_to_pause,
1149 actor_cdc_table_snapshot_splits,
1150 new_upstream_sinks,
1151 dropped_actors,
1152 sink_log_store_flush: info
1153 .replace_sink
1154 .map(|old_sink_id| vec![old_sink_id])
1155 .unwrap_or_default(),
1156 };
1157
1158 Ok(Mutation::Add(add_mutation))
1159 }
1160 }
1161 }
1162
1163 pub(super) fn replace_stream_job_to_mutation(
1165 ReplaceStreamJobPlan {
1166 old_fragments,
1167 replace_upstream,
1168 upstream_fragment_downstreams,
1169 auto_refresh_schema_sinks,
1170 ..
1171 }: &ReplaceStreamJobPlan,
1172 edges: &mut FragmentEdgeBuildResult,
1173 database_info: &mut InflightDatabaseInfo,
1174 split_assignment: &SplitAssignment,
1175 ) -> MetaResult<Option<Mutation>> {
1176 {
1177 {
1178 let merge_updates = edges
1179 .merge_updates
1180 .extract_if(|fragment_id, _| replace_upstream.contains_key(fragment_id))
1181 .collect();
1182 let dispatchers = edges
1183 .dispatchers
1184 .extract_if(|fragment_id, _| {
1185 upstream_fragment_downstreams.contains_key(fragment_id)
1186 })
1187 .collect();
1188 let actor_cdc_table_snapshot_splits = database_info
1189 .assign_cdc_backfill_splits(old_fragments.stream_job_id)?
1190 .map(|splits| PbCdcTableSnapshotSplitsWithGeneration { splits });
1191 let old_fragments = old_fragments.fragments.keys().copied();
1192 let auto_refresh_sink_fragment_ids = auto_refresh_schema_sinks
1193 .as_ref()
1194 .into_iter()
1195 .flat_map(|sinks| sinks.iter())
1196 .map(|sink| sink.original_fragment.fragment_id);
1197 Ok(Self::generate_update_mutation_for_replace_table(
1198 old_fragments
1199 .chain(auto_refresh_sink_fragment_ids)
1200 .flat_map(|fragment_id| {
1201 database_info.fragment(fragment_id).actors.keys().copied()
1202 }),
1203 merge_updates,
1204 dispatchers,
1205 split_assignment,
1206 actor_cdc_table_snapshot_splits,
1207 auto_refresh_schema_sinks.as_ref(),
1208 ))
1209 }
1210 }
1211 }
1212
1213 pub(super) fn reschedule_to_mutation(
1215 reschedules: &HashMap<FragmentId, Reschedule>,
1216 fragment_actors: &HashMap<FragmentId, HashSet<ActorId>>,
1217 control_stream_manager: &ControlStreamManager,
1218 database_info: &mut InflightDatabaseInfo,
1219 ) -> MetaResult<Option<Mutation>> {
1220 {
1221 {
1222 let database_id = database_info.database_id;
1223 let mut dispatcher_update = HashMap::new();
1224 for reschedule in reschedules.values() {
1225 for &(upstream_fragment_id, dispatcher_id) in
1226 &reschedule.upstream_fragment_dispatcher_ids
1227 {
1228 let upstream_actor_ids = fragment_actors
1230 .get(&upstream_fragment_id)
1231 .expect("should contain");
1232
1233 let upstream_reschedule = reschedules.get(&upstream_fragment_id);
1234
1235 for &actor_id in upstream_actor_ids {
1237 let added_downstream_actor_id = if upstream_reschedule
1238 .map(|reschedule| !reschedule.removed_actors.contains(&actor_id))
1239 .unwrap_or(true)
1240 {
1241 reschedule
1242 .added_actors
1243 .values()
1244 .flatten()
1245 .cloned()
1246 .collect()
1247 } else {
1248 Default::default()
1249 };
1250 dispatcher_update
1252 .try_insert(
1253 (actor_id, dispatcher_id),
1254 DispatcherUpdate {
1255 actor_id,
1256 dispatcher_id,
1257 hash_mapping: reschedule
1258 .upstream_dispatcher_mapping
1259 .as_ref()
1260 .map(|m| m.to_protobuf()),
1261 added_downstream_actor_id,
1262 removed_downstream_actor_id: reschedule
1263 .removed_actors
1264 .iter()
1265 .cloned()
1266 .collect(),
1267 },
1268 )
1269 .unwrap();
1270 }
1271 }
1272 }
1273 let dispatcher_update = dispatcher_update.into_values().collect();
1274
1275 let mut merge_update = HashMap::new();
1276 for (&fragment_id, reschedule) in reschedules {
1277 for &downstream_fragment_id in &reschedule.downstream_fragment_ids {
1278 let downstream_actor_ids = fragment_actors
1280 .get(&downstream_fragment_id)
1281 .expect("should contain");
1282
1283 let downstream_removed_actors: HashSet<_> = reschedules
1287 .get(&downstream_fragment_id)
1288 .map(|downstream_reschedule| {
1289 downstream_reschedule
1290 .removed_actors
1291 .iter()
1292 .copied()
1293 .collect()
1294 })
1295 .unwrap_or_default();
1296
1297 for &actor_id in downstream_actor_ids {
1299 if downstream_removed_actors.contains(&actor_id) {
1300 continue;
1301 }
1302
1303 merge_update
1305 .try_insert(
1306 (actor_id, fragment_id),
1307 MergeUpdate {
1308 actor_id,
1309 upstream_fragment_id: fragment_id,
1310 new_upstream_fragment_id: None,
1311 added_upstream_actors: reschedule
1312 .added_actors
1313 .iter()
1314 .flat_map(|(worker_id, actors)| {
1315 let host =
1316 control_stream_manager.host_addr(*worker_id);
1317 actors.iter().map(move |&actor_id| PbActorInfo {
1318 actor_id,
1319 host: Some(host.clone()),
1320 partial_graph_id: to_partial_graph_id(
1322 database_id,
1323 None,
1324 ),
1325 })
1326 })
1327 .collect(),
1328 removed_upstream_actor_id: reschedule
1329 .removed_actors
1330 .iter()
1331 .cloned()
1332 .collect(),
1333 },
1334 )
1335 .unwrap();
1336 }
1337 }
1338 }
1339 let merge_update = merge_update.into_values().collect();
1340
1341 let mut actor_vnode_bitmap_update = HashMap::new();
1342 for reschedule in reschedules.values() {
1343 for (&actor_id, bitmap) in &reschedule.vnode_bitmap_updates {
1345 let bitmap = bitmap.to_protobuf();
1346 actor_vnode_bitmap_update
1347 .try_insert(actor_id, bitmap)
1348 .unwrap();
1349 }
1350 }
1351 let dropped_actors = reschedules
1352 .values()
1353 .flat_map(|r| r.removed_actors.iter().copied())
1354 .collect();
1355 let mut actor_splits = HashMap::new();
1356 let mut actor_cdc_table_snapshot_splits = HashMap::new();
1357 for (fragment_id, reschedule) in reschedules {
1358 for (actor_id, splits) in &reschedule.actor_splits {
1359 actor_splits.insert(
1360 *actor_id,
1361 ConnectorSplits {
1362 splits: splits.iter().map(ConnectorSplit::from).collect(),
1363 },
1364 );
1365 }
1366
1367 if let Some(assignment) =
1368 database_info.may_assign_fragment_cdc_backfill_splits(*fragment_id)?
1369 {
1370 actor_cdc_table_snapshot_splits.extend(assignment)
1371 }
1372 }
1373
1374 let actor_new_dispatchers = HashMap::new();
1376 let mutation = Mutation::Update(UpdateMutation {
1377 dispatcher_update,
1378 merge_update,
1379 actor_vnode_bitmap_update,
1380 dropped_actors,
1381 actor_splits,
1382 actor_new_dispatchers,
1383 actor_cdc_table_snapshot_splits: Some(PbCdcTableSnapshotSplitsWithGeneration {
1384 splits: actor_cdc_table_snapshot_splits,
1385 }),
1386 sink_schema_change: Default::default(),
1387 subscriptions_to_drop: vec![],
1388 });
1389 tracing::debug!("update mutation: {mutation:?}");
1390 Ok(Some(mutation))
1391 }
1392 }
1393 }
1394
1395 pub(super) fn create_subscription_to_mutation(
1397 upstream_mv_table_id: TableId,
1398 subscription_id: SubscriptionId,
1399 ) -> Mutation {
1400 {
1401 Mutation::Add(AddMutation {
1402 actor_dispatchers: Default::default(),
1403 added_actors: vec![],
1404 actor_splits: Default::default(),
1405 pause: false,
1406 subscriptions_to_add: vec![SubscriptionUpstreamInfo {
1407 upstream_mv_table_id,
1408 subscriber_id: subscription_id.as_subscriber_id(),
1409 }],
1410 backfill_nodes_to_pause: vec![],
1411 actor_cdc_table_snapshot_splits: None,
1412 new_upstream_sinks: Default::default(),
1413 dropped_actors: Default::default(),
1414 sink_log_store_flush: Default::default(),
1415 })
1416 }
1417 }
1418
1419 pub(super) fn drop_subscription_to_mutation(
1421 upstream_mv_table_id: TableId,
1422 subscription_id: SubscriptionId,
1423 ) -> Mutation {
1424 {
1425 Mutation::DropSubscriptions(DropSubscriptionsMutation {
1426 info: vec![SubscriptionUpstreamInfo {
1427 subscriber_id: subscription_id.as_subscriber_id(),
1428 upstream_mv_table_id,
1429 }],
1430 })
1431 }
1432 }
1433
1434 pub(super) fn connector_props_change_to_mutation(config: &ConnectorPropsChange) -> Mutation {
1436 {
1437 {
1438 let mut connector_props_infos = HashMap::default();
1439 for (k, v) in config {
1440 connector_props_infos.insert(
1441 k.as_raw_id(),
1442 ConnectorPropsInfo {
1443 connector_props_info: v.clone(),
1444 },
1445 );
1446 }
1447 Mutation::ConnectorPropsChange(ConnectorPropsChangeMutation {
1448 connector_props_infos,
1449 })
1450 }
1451 }
1452 }
1453
1454 pub(super) fn refresh_to_mutation(
1456 table_id: TableId,
1457 associated_source_id: SourceId,
1458 ) -> Mutation {
1459 Mutation::RefreshStart(risingwave_pb::stream_plan::RefreshStartMutation {
1460 table_id,
1461 associated_source_id,
1462 })
1463 }
1464
1465 pub(super) fn list_finish_to_mutation(associated_source_id: SourceId) -> Mutation {
1467 Mutation::ListFinish(ListFinishMutation {
1468 associated_source_id,
1469 })
1470 }
1471
1472 pub(super) fn load_finish_to_mutation(associated_source_id: SourceId) -> Mutation {
1474 Mutation::LoadFinish(LoadFinishMutation {
1475 associated_source_id,
1476 })
1477 }
1478
1479 pub(super) fn reset_source_to_mutation(source_id: SourceId) -> Mutation {
1481 Mutation::ResetSource(risingwave_pb::stream_plan::ResetSourceMutation {
1482 source_id: source_id.as_raw_id(),
1483 })
1484 }
1485
1486 pub(super) fn resume_backfill_to_mutation(
1488 target: &ResumeBackfillTarget,
1489 database_info: &InflightDatabaseInfo,
1490 ) -> MetaResult<Option<Mutation>> {
1491 {
1492 {
1493 let fragment_ids: HashSet<_> = match target {
1494 ResumeBackfillTarget::Job(job_id) => {
1495 database_info.backfill_fragment_ids_for_job(*job_id)?
1496 }
1497 ResumeBackfillTarget::Fragment(fragment_id) => {
1498 if !database_info.is_backfill_fragment(*fragment_id)? {
1499 return Err(MetaError::invalid_parameter(format!(
1500 "fragment {} is not a backfill node",
1501 fragment_id
1502 )));
1503 }
1504 HashSet::from([*fragment_id])
1505 }
1506 };
1507 if fragment_ids.is_empty() {
1508 warn!(
1509 ?target,
1510 "resume backfill command ignored because no backfill fragments found"
1511 );
1512 Ok(None)
1513 } else {
1514 Ok(Some(Mutation::StartFragmentBackfill(
1515 StartFragmentBackfillMutation {
1516 fragment_ids: fragment_ids.into_iter().collect(),
1517 },
1518 )))
1519 }
1520 }
1521 }
1522 }
1523
1524 pub(super) fn inject_source_offsets_to_mutation(
1526 source_id: SourceId,
1527 split_offsets: &HashMap<String, String>,
1528 ) -> Mutation {
1529 Mutation::InjectSourceOffsets(risingwave_pb::stream_plan::InjectSourceOffsetsMutation {
1530 source_id: source_id.as_raw_id(),
1531 split_offsets: split_offsets.clone(),
1532 })
1533 }
1534
1535 pub(super) fn create_streaming_job_actors_to_create(
1537 info: &CreateStreamingJobCommandInfo,
1538 edges: &mut FragmentEdgeBuildResult,
1539 stream_actors: &HashMap<FragmentId, Vec<StreamActor>>,
1540 actor_location: &HashMap<ActorId, WorkerId>,
1541 ) -> StreamJobActorsToCreate {
1542 {
1543 {
1544 edges.collect_actors_to_create(info.stream_job_fragments.fragments.values().map(
1545 |fragment| {
1546 let actors = stream_actors
1547 .get(&fragment.fragment_id)
1548 .into_iter()
1549 .flatten()
1550 .map(|actor| (actor, actor_location[&actor.actor_id]));
1551 (
1552 fragment.fragment_id,
1553 &fragment.nodes,
1554 actors,
1555 [], )
1557 },
1558 ))
1559 }
1560 }
1561 }
1562
1563 pub(super) fn reschedule_actors_to_create(
1565 reschedules: &HashMap<FragmentId, Reschedule>,
1566 fragment_actors: &HashMap<FragmentId, HashSet<ActorId>>,
1567 database_info: &InflightDatabaseInfo,
1568 control_stream_manager: &ControlStreamManager,
1569 ) -> StreamJobActorsToCreate {
1570 {
1571 {
1572 let mut actor_upstreams = Self::collect_database_partial_graph_actor_upstreams(
1573 reschedules.iter().map(|(fragment_id, reschedule)| {
1574 (
1575 *fragment_id,
1576 reschedule.newly_created_actors.values().map(
1577 |((actor, dispatchers), _)| {
1578 (actor.actor_id, dispatchers.as_slice())
1579 },
1580 ),
1581 )
1582 }),
1583 Some((reschedules, fragment_actors)),
1584 database_info,
1585 control_stream_manager,
1586 );
1587 let mut map: HashMap<WorkerId, HashMap<_, (_, Vec<_>, _)>> = HashMap::new();
1588 for (fragment_id, (actor, dispatchers), worker_id) in
1589 reschedules.iter().flat_map(|(fragment_id, reschedule)| {
1590 reschedule
1591 .newly_created_actors
1592 .values()
1593 .map(|(actors, status)| (*fragment_id, actors, status))
1594 })
1595 {
1596 let upstreams = actor_upstreams.remove(&actor.actor_id).unwrap_or_default();
1597 map.entry(*worker_id)
1598 .or_default()
1599 .entry(fragment_id)
1600 .or_insert_with(|| {
1601 let node = database_info.fragment(fragment_id).nodes.clone();
1602 let subscribers =
1603 database_info.fragment_subscribers(fragment_id).collect();
1604 (node, vec![], subscribers)
1605 })
1606 .1
1607 .push((actor.clone(), upstreams, dispatchers.clone()));
1608 }
1609 map
1610 }
1611 }
1612 }
1613
1614 pub(super) fn replace_stream_job_actors_to_create(
1616 replace_table: &ReplaceStreamJobPlan,
1617 edges: &mut FragmentEdgeBuildResult,
1618 database_info: &InflightDatabaseInfo,
1619 stream_actors: &HashMap<FragmentId, Vec<StreamActor>>,
1620 actor_location: &HashMap<ActorId, WorkerId>,
1621 ) -> StreamJobActorsToCreate {
1622 {
1623 {
1624 let mut actors = edges.collect_actors_to_create(
1625 replace_table
1626 .new_fragments
1627 .fragments
1628 .values()
1629 .map(|fragment| {
1630 let actors = stream_actors
1631 .get(&fragment.fragment_id)
1632 .into_iter()
1633 .flatten()
1634 .map(|actor| (actor, actor_location[&actor.actor_id]));
1635 (
1636 fragment.fragment_id,
1637 &fragment.nodes,
1638 actors,
1639 database_info
1640 .job_subscribers(replace_table.old_fragments.stream_job_id),
1641 )
1642 }),
1643 );
1644
1645 if let Some(sinks) = &replace_table.auto_refresh_schema_sinks {
1647 let sink_actors = edges.collect_actors_to_create(sinks.iter().map(|sink| {
1648 (
1649 sink.new_fragment.fragment_id,
1650 &sink.new_fragment.nodes,
1651 stream_actors
1652 .get(&sink.new_fragment.fragment_id)
1653 .into_iter()
1654 .flatten()
1655 .map(|actor| (actor, actor_location[&actor.actor_id])),
1656 database_info.job_subscribers(sink.original_sink.id.as_job_id()),
1657 )
1658 }));
1659 for (worker_id, fragment_actors) in sink_actors {
1660 actors.entry(worker_id).or_default().extend(fragment_actors);
1661 }
1662 }
1663 actors
1664 }
1665 }
1666 }
1667
1668 fn generate_update_mutation_for_replace_table(
1669 dropped_actors: impl IntoIterator<Item = ActorId>,
1670 merge_updates: HashMap<FragmentId, Vec<MergeUpdate>>,
1671 dispatchers: FragmentActorDispatchers,
1672 split_assignment: &SplitAssignment,
1673 cdc_table_snapshot_split_assignment: Option<PbCdcTableSnapshotSplitsWithGeneration>,
1674 auto_refresh_schema_sinks: Option<&Vec<AutoRefreshSchemaSinkContext>>,
1675 ) -> Option<Mutation> {
1676 let dropped_actors = dropped_actors.into_iter().collect();
1677
1678 let actor_new_dispatchers = dispatchers
1679 .into_values()
1680 .flatten()
1681 .map(|(actor_id, dispatchers)| (actor_id, Dispatchers { dispatchers }))
1682 .collect();
1683
1684 let actor_splits = split_assignment
1685 .values()
1686 .flat_map(build_actor_connector_splits)
1687 .collect();
1688 Some(Mutation::Update(UpdateMutation {
1689 actor_new_dispatchers,
1690 merge_update: merge_updates.into_values().flatten().collect(),
1691 dropped_actors,
1692 actor_splits,
1693 actor_cdc_table_snapshot_splits: cdc_table_snapshot_split_assignment,
1694 sink_schema_change: auto_refresh_schema_sinks
1695 .as_ref()
1696 .into_iter()
1697 .flat_map(|sinks| {
1698 sinks.iter().map(|sink| {
1699 let op = if !sink.removed_column_names.is_empty() {
1700 PbSinkSchemaChangeOp::DropColumns(PbSinkDropColumnsOp {
1701 column_names: sink.removed_column_names.clone(),
1702 })
1703 } else {
1704 PbSinkSchemaChangeOp::AddColumns(PbSinkAddColumnsOp {
1705 fields: sink
1706 .newly_add_fields
1707 .iter()
1708 .map(|field| field.to_prost())
1709 .collect(),
1710 })
1711 };
1712 (
1713 sink.original_sink.id.as_raw_id(),
1714 PbSinkSchemaChange {
1715 original_schema: sink_original_schema_fields(
1716 &sink.original_sink.columns,
1717 ),
1718 op: Some(op),
1719 },
1720 )
1721 })
1722 })
1723 .collect(),
1724 ..Default::default()
1725 }))
1726 }
1727}
1728
1729impl Command {
1730 #[expect(clippy::type_complexity)]
1731 pub(super) fn collect_database_partial_graph_actor_upstreams(
1732 actor_dispatchers: impl Iterator<
1733 Item = (FragmentId, impl Iterator<Item = (ActorId, &[Dispatcher])>),
1734 >,
1735 reschedule_dispatcher_update: Option<(
1736 &HashMap<FragmentId, Reschedule>,
1737 &HashMap<FragmentId, HashSet<ActorId>>,
1738 )>,
1739 database_info: &InflightDatabaseInfo,
1740 control_stream_manager: &ControlStreamManager,
1741 ) -> HashMap<ActorId, ActorUpstreams> {
1742 let mut actor_upstreams: HashMap<ActorId, ActorUpstreams> = HashMap::new();
1743 for (upstream_fragment_id, upstream_actors) in actor_dispatchers {
1744 let upstream_fragment = database_info.fragment(upstream_fragment_id);
1745 for (upstream_actor_id, dispatchers) in upstream_actors {
1746 let upstream_actor_location =
1747 upstream_fragment.actors[&upstream_actor_id].worker_id;
1748 let upstream_actor_host = control_stream_manager.host_addr(upstream_actor_location);
1749 for downstream_actor_id in dispatchers
1750 .iter()
1751 .flat_map(|dispatcher| dispatcher.downstream_actor_id.iter())
1752 {
1753 actor_upstreams
1754 .entry(*downstream_actor_id)
1755 .or_default()
1756 .entry(upstream_fragment_id)
1757 .or_default()
1758 .insert(
1759 upstream_actor_id,
1760 PbActorInfo {
1761 actor_id: upstream_actor_id,
1762 host: Some(upstream_actor_host.clone()),
1763 partial_graph_id: to_partial_graph_id(
1764 database_info.database_id,
1765 None,
1766 ),
1767 },
1768 );
1769 }
1770 }
1771 }
1772 if let Some((reschedules, fragment_actors)) = reschedule_dispatcher_update {
1773 for reschedule in reschedules.values() {
1774 for (upstream_fragment_id, _) in &reschedule.upstream_fragment_dispatcher_ids {
1775 let upstream_fragment = database_info.fragment(*upstream_fragment_id);
1776 let upstream_reschedule = reschedules.get(upstream_fragment_id);
1777 for upstream_actor_id in fragment_actors
1778 .get(upstream_fragment_id)
1779 .expect("should exist")
1780 {
1781 let upstream_actor_location =
1782 upstream_fragment.actors[upstream_actor_id].worker_id;
1783 let upstream_actor_host =
1784 control_stream_manager.host_addr(upstream_actor_location);
1785 if let Some(upstream_reschedule) = upstream_reschedule
1786 && upstream_reschedule
1787 .removed_actors
1788 .contains(upstream_actor_id)
1789 {
1790 continue;
1791 }
1792 for (_, downstream_actor_id) in
1793 reschedule
1794 .added_actors
1795 .iter()
1796 .flat_map(|(worker_id, actors)| {
1797 actors.iter().map(|actor| (*worker_id, *actor))
1798 })
1799 {
1800 actor_upstreams
1801 .entry(downstream_actor_id)
1802 .or_default()
1803 .entry(*upstream_fragment_id)
1804 .or_default()
1805 .insert(
1806 *upstream_actor_id,
1807 PbActorInfo {
1808 actor_id: *upstream_actor_id,
1809 host: Some(upstream_actor_host.clone()),
1810 partial_graph_id: to_partial_graph_id(
1811 database_info.database_id,
1812 None,
1813 ),
1814 },
1815 );
1816 }
1817 }
1818 }
1819 }
1820 }
1821 actor_upstreams
1822 }
1823}
1824
1825#[cfg(test)]
1826mod tests {
1827 use risingwave_pb::data::PbDataType;
1828 use risingwave_pb::data::data_type::PbTypeName;
1829 use risingwave_pb::plan_common::{ColumnCatalog as PbColumnCatalog, ColumnDesc};
1830
1831 use super::sink_original_schema_fields;
1832
1833 fn column(name: &str, type_name: PbTypeName, is_hidden: bool) -> PbColumnCatalog {
1834 PbColumnCatalog {
1835 column_desc: Some(ColumnDesc {
1836 column_type: Some(PbDataType {
1837 type_name: type_name as i32,
1838 ..Default::default()
1839 }),
1840 name: name.to_owned(),
1841 ..Default::default()
1842 }),
1843 is_hidden,
1844 }
1845 }
1846
1847 #[test]
1848 fn test_sink_original_schema_fields_skips_hidden_columns() {
1849 let columns = vec![
1850 column("k", PbTypeName::Int32, false),
1851 column("v", PbTypeName::Varchar, false),
1852 column("_row_id", PbTypeName::Serial, true),
1853 ];
1854
1855 let fields = sink_original_schema_fields(&columns);
1856 let field_names = fields
1857 .iter()
1858 .map(|field| field.name.as_str())
1859 .collect::<Vec<_>>();
1860
1861 assert_eq!(field_names, ["k", "v"]);
1862 }
1863}