1mod split_assignment;
16mod worker;
17use std::borrow::BorrowMut;
18use std::cmp::Ordering;
19use std::collections::hash_map::Entry;
20use std::collections::{BTreeMap, BTreeSet, BinaryHeap, HashMap, HashSet};
21use std::sync::Arc;
22use std::time::Duration;
23
24use anyhow::Context;
25use risingwave_common::catalog::DatabaseId;
26use risingwave_common::id::ObjectId;
27use risingwave_common::metrics::{
28 LabelGuardedHistogram, LabelGuardedIntCounter, LabelGuardedIntGauge,
29};
30use risingwave_common::panic_if_debug;
31use risingwave_connector::WithOptionsSecResolved;
32use risingwave_connector::error::ConnectorResult;
33use risingwave_connector::source::{
34 AnySplitEnumerator, ConnectorProperties, SourceEnumeratorContext, SourceEnumeratorInfo,
35 SplitId, SplitImpl, SplitMetaData,
36};
37use risingwave_meta_model::SourceId;
38use risingwave_pb::catalog::Source;
39use risingwave_pb::source::{ConnectorSplit, ConnectorSplits};
40pub use split_assignment::{SplitDiffOptions, SplitState, align_splits, reassign_splits};
41use thiserror_ext::AsReport;
42use tokio::sync::mpsc::{UnboundedReceiver, UnboundedSender};
43use tokio::sync::{Mutex, MutexGuard, oneshot};
44use tokio::task::JoinHandle;
45use tokio::time::MissedTickBehavior;
46use tokio::{select, time};
47pub use worker::create_source_worker;
48use worker::{ConnectorSourceWorkerHandle, create_source_worker_async};
49
50use crate::barrier::{BarrierScheduler, Command, ReplaceStreamJobPlan};
51use crate::manager::{MetaSrvEnv, MetadataManager};
52use crate::model::{ActorId, FragmentId, StreamJobFragments};
53use crate::rpc::metrics::MetaMetrics;
54use crate::{MetaError, MetaResult};
55
56pub type SourceManagerRef = Arc<SourceManager>;
57pub type SplitAssignment = HashMap<FragmentId, HashMap<ActorId, Vec<SplitImpl>>>;
59
60pub type SourceSplitAssignment = HashMap<SourceId, DiscoveredSplits>;
66
67#[derive(Debug, Clone)]
72pub enum DiscoveredSplits {
73 Fixed(BTreeMap<Arc<str>, SplitImpl>),
75 Adaptive(SplitImpl),
78}
79
80#[derive(Debug, Clone)]
86pub enum ReplaceJobSplitPlan {
87 Discovered(SourceSplitAssignment),
91 AlignFromPrevious,
97}
98
99pub type ConnectorPropsChange = HashMap<ObjectId, HashMap<String, String>>;
101
102const DEFAULT_SOURCE_TICK_TIMEOUT: Duration = Duration::from_secs(10);
103
104pub struct SourceManager {
107 pub paused: Mutex<()>,
108 barrier_scheduler: BarrierScheduler,
109 core: Mutex<SourceManagerCore>,
110 pub metrics: Arc<MetaMetrics>,
111}
112pub struct SourceManagerCore {
113 metadata_manager: MetadataManager,
114
115 managed_sources: HashMap<SourceId, ConnectorSourceWorkerHandle>,
117 source_fragments: HashMap<SourceId, BTreeSet<FragmentId>>,
119 backfill_fragments: HashMap<SourceId, BTreeSet<(FragmentId, FragmentId)>>,
121
122 env: MetaSrvEnv,
123}
124
125pub struct SourceManagerRunningInfo {
126 pub source_fragments: HashMap<SourceId, BTreeSet<FragmentId>>,
127 pub backfill_fragments: HashMap<SourceId, BTreeSet<(FragmentId, FragmentId)>>,
128}
129
130impl SourceManagerCore {
131 fn new(
132 metadata_manager: MetadataManager,
133 managed_sources: HashMap<SourceId, ConnectorSourceWorkerHandle>,
134 source_fragments: HashMap<SourceId, BTreeSet<FragmentId>>,
135 backfill_fragments: HashMap<SourceId, BTreeSet<(FragmentId, FragmentId)>>,
136 env: MetaSrvEnv,
137 ) -> Self {
138 Self {
139 metadata_manager,
140 managed_sources,
141 source_fragments,
142 backfill_fragments,
143 env,
144 }
145 }
146
147 pub fn apply_source_change(&mut self, source_change: SourceChange) {
149 let mut added_source_fragments = Default::default();
150 let mut added_backfill_fragments = Default::default();
151 let mut finished_backfill_fragments = Default::default();
152 let mut fragment_replacements = Default::default();
153 let mut dropped_source_fragments = Default::default();
154 let mut dropped_source_ids = Default::default();
155 let mut recreate_source_id_map_new_props: Vec<(SourceId, HashMap<String, String>)> =
156 Default::default();
157
158 match source_change {
159 SourceChange::CreateJob {
160 added_source_fragments: added_source_fragments_,
161 added_backfill_fragments: added_backfill_fragments_,
162 } => {
163 added_source_fragments = added_source_fragments_;
164 added_backfill_fragments = added_backfill_fragments_;
165 }
166 SourceChange::CreateJobFinished {
167 finished_backfill_fragments: finished_backfill_fragments_,
168 } => {
169 finished_backfill_fragments = finished_backfill_fragments_;
170 }
171
172 SourceChange::DropMv {
173 dropped_source_fragments: dropped_source_fragments_,
174 } => {
175 dropped_source_fragments = dropped_source_fragments_;
176 }
177 SourceChange::ReplaceJob {
178 dropped_source_fragments: dropped_source_fragments_,
179 added_source_fragments: added_source_fragments_,
180 fragment_replacements: fragment_replacements_,
181 } => {
182 dropped_source_fragments = dropped_source_fragments_;
183 added_source_fragments = added_source_fragments_;
184 fragment_replacements = fragment_replacements_;
185 }
186 SourceChange::DropSource {
187 dropped_source_ids: dropped_source_ids_,
188 } => {
189 dropped_source_ids = dropped_source_ids_;
190 }
191
192 SourceChange::UpdateSourceProps {
193 source_id_map_new_props,
194 } => {
195 for (source_id, new_props) in source_id_map_new_props {
196 recreate_source_id_map_new_props.push((source_id, new_props));
197 }
198 }
199 }
200
201 for source_id in dropped_source_ids {
202 let dropped_fragments = self.source_fragments.remove(&source_id);
203
204 if let Some(handle) = self.managed_sources.remove(&source_id) {
205 handle.terminate(dropped_fragments);
206 }
207 if let Some(_fragments) = self.backfill_fragments.remove(&source_id) {
208 }
215 }
216
217 for (source_id, fragments) in added_source_fragments {
218 self.source_fragments
219 .entry(source_id)
220 .or_default()
221 .extend(fragments);
222 }
223
224 for (source_id, fragments) in added_backfill_fragments {
225 self.backfill_fragments
226 .entry(source_id)
227 .or_default()
228 .extend(fragments);
229 }
230
231 for (source_id, fragments) in finished_backfill_fragments {
232 let handle = self.managed_sources.get(&source_id).unwrap_or_else(|| {
233 panic!(
234 "source {} not found when adding backfill fragments {:?}",
235 source_id, fragments
236 );
237 });
238 handle.finish_backfill(fragments.iter().map(|(id, _up_id)| *id).collect());
239 }
240
241 for (source_id, fragment_ids) in dropped_source_fragments {
242 self.drop_source_fragments(Some(source_id), fragment_ids);
243 }
244
245 for (old_fragment_id, new_fragment_id) in fragment_replacements {
246 self.drop_source_fragments(None, BTreeSet::from([old_fragment_id]));
248
249 for fragment_ids in self.backfill_fragments.values_mut() {
250 let mut new_backfill_fragment_ids = fragment_ids.clone();
251 for (fragment_id, upstream_fragment_id) in fragment_ids.iter() {
252 assert_ne!(
253 fragment_id, upstream_fragment_id,
254 "backfill fragment should not be replaced"
255 );
256 if *upstream_fragment_id == old_fragment_id {
257 new_backfill_fragment_ids.remove(&(*fragment_id, *upstream_fragment_id));
258 new_backfill_fragment_ids.insert((*fragment_id, new_fragment_id));
259 }
260 }
261 *fragment_ids = new_backfill_fragment_ids;
262 }
263 }
264
265 for (source_id, new_props) in recreate_source_id_map_new_props {
266 if let Some(handle) = self.managed_sources.get_mut(&source_id) {
267 let props_wrapper =
270 WithOptionsSecResolved::without_secrets(new_props.into_iter().collect());
271 let props = ConnectorProperties::extract(props_wrapper, false).unwrap(); handle.update_props(props);
273 tracing::info!("update source {source_id} properties in source manager");
274 } else {
275 tracing::info!("job id {source_id} is not registered in source manager");
276 }
277 }
278 }
279
280 fn drop_source_fragments(
281 &mut self,
282 source_id: Option<SourceId>,
283 dropped_fragment_ids: BTreeSet<FragmentId>,
284 ) {
285 if let Some(source_id) = source_id {
286 if let Entry::Occupied(mut entry) = self.source_fragments.entry(source_id) {
287 let mut dropped_ids = vec![];
288 let managed_fragment_ids = entry.get_mut();
289 for fragment_id in &dropped_fragment_ids {
290 managed_fragment_ids.remove(fragment_id);
291 dropped_ids.push(*fragment_id);
292 }
293 if let Some(handle) = self.managed_sources.get(&source_id) {
294 handle.drop_fragments(dropped_ids);
295 } else {
296 panic_if_debug!(
297 "source {source_id} not found when dropping fragment {dropped_ids:?}",
298 );
299 }
300 if managed_fragment_ids.is_empty() {
301 entry.remove();
302 }
303 }
304 } else {
305 for (source_id, fragment_ids) in &mut self.source_fragments {
306 let mut dropped_ids = vec![];
307 for fragment_id in &dropped_fragment_ids {
308 if fragment_ids.remove(fragment_id) {
309 dropped_ids.push(*fragment_id);
310 }
311 }
312 if !dropped_ids.is_empty() {
313 if let Some(handle) = self.managed_sources.get(source_id) {
314 handle.drop_fragments(dropped_ids);
315 } else {
316 panic_if_debug!(
317 "source {source_id} not found when dropping fragment {dropped_ids:?}",
318 );
319 }
320 }
321 }
322 }
323 }
324}
325
326impl SourceManager {
327 const DEFAULT_SOURCE_TICK_INTERVAL: Duration = Duration::from_secs(10);
328
329 pub async fn new(
330 barrier_scheduler: BarrierScheduler,
331 metadata_manager: MetadataManager,
332 metrics: Arc<MetaMetrics>,
333 env: MetaSrvEnv,
334 ) -> MetaResult<Self> {
335 let mut managed_sources = HashMap::new();
336 {
337 let sources = metadata_manager.list_sources().await?;
338 for source in sources {
339 if source
340 .info
341 .as_ref()
342 .is_some_and(|info| info.external_table.is_some())
343 {
344 continue;
345 }
346 create_source_worker_async(
347 source,
348 &mut managed_sources,
349 metrics.clone(),
350 env.await_tree_reg().clone(),
351 )?
352 }
353 }
354
355 let source_fragments = metadata_manager
356 .catalog_controller
357 .load_source_fragment_ids()
358 .await?
359 .into_iter()
360 .map(|(source_id, fragment_ids)| {
361 (
362 source_id as SourceId,
363 fragment_ids.into_iter().map(|id| id as _).collect(),
364 )
365 })
366 .collect();
367 let backfill_fragments = metadata_manager
368 .catalog_controller
369 .load_backfill_fragment_ids()
370 .await?;
371
372 let core = Mutex::new(SourceManagerCore::new(
373 metadata_manager,
374 managed_sources,
375 source_fragments,
376 backfill_fragments,
377 env,
378 ));
379
380 Ok(Self {
381 barrier_scheduler,
382 core,
383 paused: Mutex::new(()),
384 metrics,
385 })
386 }
387
388 pub async fn validate_source_once(
389 &self,
390 source_id: SourceId,
391 new_source_props: WithOptionsSecResolved,
392 ) -> MetaResult<()> {
393 let props = ConnectorProperties::extract(new_source_props, false).unwrap();
394
395 {
396 let mut enumerator = props
397 .create_split_enumerator(Arc::new(SourceEnumeratorContext {
398 metrics: self.metrics.source_enumerator_metrics.clone(),
399 info: SourceEnumeratorInfo { source_id },
400 }))
401 .await
402 .context("failed to create SplitEnumerator")?;
403
404 validate_enumerator_once(&mut *enumerator).await?;
405 }
406 Ok(())
407 }
408
409 #[await_tree::instrument]
411 pub async fn handle_replace_job(
412 &self,
413 dropped_job_fragments: &StreamJobFragments,
414 added_source_fragments: HashMap<SourceId, BTreeSet<FragmentId>>,
415 replace_plan: &ReplaceStreamJobPlan,
416 ) {
417 let dropped_source_fragments = dropped_job_fragments.stream_source_fragments();
419
420 self.apply_source_change(SourceChange::ReplaceJob {
421 dropped_source_fragments,
422 added_source_fragments,
423 fragment_replacements: replace_plan.fragment_replacements(),
424 })
425 .await;
426 }
427
428 #[await_tree::instrument("apply_source_change({source_change})")]
431 pub async fn apply_source_change(&self, source_change: SourceChange) {
432 let need_force_tick = matches!(source_change, SourceChange::UpdateSourceProps { .. });
433 let updated_source_ids = if let SourceChange::UpdateSourceProps {
434 ref source_id_map_new_props,
435 } = source_change
436 {
437 source_id_map_new_props.keys().cloned().collect::<Vec<_>>()
438 } else {
439 Vec::new()
440 };
441
442 {
443 let mut core = self.core.lock().await;
444 core.apply_source_change(source_change);
445 }
446
447 if need_force_tick {
449 self.force_tick_updated_sources(updated_source_ids).await;
450 }
451 }
452
453 #[await_tree::instrument("register_source({})", source.name)]
455 pub async fn register_source(&self, source: &Source) -> MetaResult<()> {
456 tracing::debug!("register_source: {}", source.get_id());
457 let mut core = self.core.lock().await;
458 let source_id = source.get_id();
459 if core.managed_sources.contains_key(&source_id) {
460 tracing::warn!("source {} already registered", source_id);
461 return Ok(());
462 }
463
464 let handle = create_source_worker(
465 source,
466 self.metrics.clone(),
467 core.env.await_tree_reg().clone(),
468 )
469 .await
470 .context("failed to create source worker")?;
471
472 core.managed_sources.insert(source_id, handle);
473
474 Ok(())
475 }
476
477 pub async fn register_source_with_handle(
479 &self,
480 source_id: SourceId,
481 handle: ConnectorSourceWorkerHandle,
482 ) {
483 let mut core = self.core.lock().await;
484 if core.managed_sources.contains_key(&source_id) {
485 tracing::warn!("source {} already registered", source_id);
486 return;
487 }
488
489 core.managed_sources.insert(source_id, handle);
490 }
491
492 pub async fn get_running_info(&self) -> SourceManagerRunningInfo {
493 let core = self.core.lock().await;
494
495 SourceManagerRunningInfo {
496 source_fragments: core.source_fragments.clone(),
497 backfill_fragments: core.backfill_fragments.clone(),
498 }
499 }
500
501 async fn tick(&self) -> MetaResult<()> {
510 let split_states = {
511 let core_guard = self.core.lock().await;
512 core_guard.reassign_splits().await?
513 };
514
515 for (database_id, split_state) in split_states {
516 if !split_state.split_assignment.is_empty() {
517 let command = Command::SourceChangeSplit(split_state);
518 tracing::info!(command = ?command, "pushing down split assignment command");
519 self.barrier_scheduler
520 .run_command(database_id, command)
521 .await?;
522 }
523 }
524
525 Ok(())
526 }
527
528 pub async fn run(&self) -> MetaResult<()> {
529 let mut ticker = time::interval(Self::DEFAULT_SOURCE_TICK_INTERVAL);
530 ticker.set_missed_tick_behavior(MissedTickBehavior::Skip);
531 loop {
532 ticker.tick().await;
533 let _pause_guard = self.paused.lock().await;
534 if let Err(e) = self.tick().await {
535 tracing::error!(
536 error = %e.as_report(),
537 "source manager tick failed",
538 );
539 }
540 }
541 }
542
543 pub async fn pause_tick(&self) -> MutexGuard<'_, ()> {
545 tracing::debug!("pausing tick lock in source manager");
546 self.paused.lock().await
547 }
548
549 async fn force_tick_updated_sources(&self, updated_source_ids: Vec<SourceId>) {
551 let core = self.core.lock().await;
552 for source_id in updated_source_ids {
553 if let Some(handle) = core.managed_sources.get(&source_id) {
554 tracing::info!("forcing tick for updated source {}", source_id);
555 if let Err(e) = handle.force_tick().await {
556 tracing::warn!(
557 error = %e.as_report(),
558 "failed to force tick for source {} after properties update",
559 source_id
560 );
561 }
562 } else {
563 tracing::warn!(
564 "source {} not found when trying to force tick after update",
565 source_id
566 );
567 }
568 }
569 }
570
571 pub async fn reset_source_splits(&self, source_id: SourceId) -> MetaResult<()> {
575 tracing::warn!(
576 %source_id,
577 "UNSAFE: Resetting source splits - clearing cached state and triggering re-discovery"
578 );
579
580 let core = self.core.lock().await;
581 if let Some(handle) = core.managed_sources.get(&source_id) {
582 let prev_splits = handle.splits.take();
584 tracing::info!(
585 %source_id,
586 prev_splits = ?prev_splits.as_ref().map(|s| s.len()),
587 "Clearing cached splits"
588 );
589
590 tracing::info!(
592 %source_id,
593 "Triggering split re-discovery via force_tick"
594 );
595 handle.force_tick().await.with_context(|| {
596 format!(
597 "failed to force tick for source {} after split reset",
598 source_id
599 )
600 })?;
601
602 tracing::info!(
603 %source_id,
604 "Split reset completed - new splits will be assigned on next tick"
605 );
606 Ok(())
607 } else {
608 Err(anyhow::anyhow!("source {} not found in source manager", source_id).into())
609 }
610 }
611
612 pub async fn validate_inject_source_offsets(
619 &self,
620 source_id: SourceId,
621 split_offsets: &HashMap<String, String>,
622 ) -> MetaResult<Vec<String>> {
623 let (fragment_ids, env) = {
624 let core = self.core.lock().await;
625
626 let _ = core.managed_sources.get(&source_id).ok_or_else(|| {
628 MetaError::invalid_parameter(format!(
629 "source {} not found in source manager",
630 source_id
631 ))
632 })?;
633
634 let mut ids = Vec::new();
635 if let Some(src_frags) = core.source_fragments.get(&source_id) {
636 ids.extend(src_frags.iter().copied());
637 }
638 if let Some(backfill_frags) = core.backfill_fragments.get(&source_id) {
639 ids.extend(
640 backfill_frags
641 .iter()
642 .flat_map(|(id, upstream)| [*id, *upstream]),
643 );
644 }
645 (ids, core.env.clone())
646 };
647
648 if fragment_ids.is_empty() {
649 return Err(MetaError::invalid_parameter(format!(
650 "source {} has no running fragments",
651 source_id
652 )));
653 }
654
655 let guard = env.shared_actor_infos().read_guard();
656 let mut assigned_split_ids = HashSet::new();
657 for fragment_id in fragment_ids {
658 if let Some(fragment) = guard.get_fragment(fragment_id) {
659 for actor in fragment.actors.values() {
660 for split in &actor.splits {
661 assigned_split_ids.insert(split.id().to_string());
662 }
663 }
664 }
665 }
666
667 let mut invalid_splits = Vec::new();
669 for split_id in split_offsets.keys() {
670 if !assigned_split_ids.contains(split_id) {
671 invalid_splits.push(split_id.clone());
672 }
673 }
674
675 if !invalid_splits.is_empty() {
676 return Err(MetaError::invalid_parameter(format!(
677 "invalid split IDs for source {}: {:?}. Valid splits are: {:?}",
678 source_id,
679 invalid_splits,
680 assigned_split_ids.iter().collect::<Vec<_>>()
681 )));
682 }
683
684 tracing::info!(
685 source_id = %source_id,
686 num_splits = split_offsets.len(),
687 "Validated inject source offsets request"
688 );
689
690 Ok(split_offsets.keys().cloned().collect())
691 }
692}
693
694async fn validate_enumerator_once(enumerator: &mut dyn AnySplitEnumerator) -> MetaResult<()> {
695 let _ = tokio::time::timeout(DEFAULT_SOURCE_TICK_TIMEOUT, enumerator.list_splits())
696 .await
697 .context("failed to list splits")??;
698 Ok(())
699}
700
701#[derive(strum::Display, Debug)]
702pub enum SourceChange {
703 CreateJob {
706 added_source_fragments: HashMap<SourceId, BTreeSet<FragmentId>>,
707 added_backfill_fragments: HashMap<SourceId, BTreeSet<(FragmentId, FragmentId)>>,
709 },
710 UpdateSourceProps {
711 source_id_map_new_props: HashMap<SourceId, HashMap<String, String>>,
714 },
715 CreateJobFinished {
719 finished_backfill_fragments: HashMap<SourceId, BTreeSet<(FragmentId, FragmentId)>>,
721 },
722 DropSource { dropped_source_ids: Vec<SourceId> },
724 DropMv {
725 dropped_source_fragments: HashMap<SourceId, BTreeSet<FragmentId>>,
727 },
728 ReplaceJob {
729 dropped_source_fragments: HashMap<SourceId, BTreeSet<FragmentId>>,
730 added_source_fragments: HashMap<SourceId, BTreeSet<FragmentId>>,
731 fragment_replacements: HashMap<FragmentId, FragmentId>,
732 },
733}
734
735pub fn build_actor_connector_splits(
736 splits: &HashMap<ActorId, Vec<SplitImpl>>,
737) -> HashMap<ActorId, ConnectorSplits> {
738 splits
739 .iter()
740 .map(|(&actor_id, splits)| {
741 (
742 actor_id,
743 ConnectorSplits {
744 splits: splits.iter().map(ConnectorSplit::from).collect(),
745 },
746 )
747 })
748 .collect()
749}
750
751pub fn build_actor_split_impls(
752 actor_splits: &HashMap<ActorId, ConnectorSplits>,
753) -> HashMap<ActorId, Vec<SplitImpl>> {
754 actor_splits
755 .iter()
756 .map(|(actor_id, ConnectorSplits { splits })| {
757 (
758 *actor_id,
759 splits
760 .iter()
761 .map(|split| SplitImpl::try_from(split).unwrap())
762 .collect(),
763 )
764 })
765 .collect()
766}