Skip to main content

risingwave_stream/task/barrier_manager/
mod.rs

1// Copyright 2025 RisingWave Labs
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15pub mod cdc_progress;
16pub mod progress;
17
18use std::sync::Arc;
19
20pub use progress::CreateMviewProgressReporter;
21use risingwave_common::id::{SourceId, TableId};
22use risingwave_common::util::epoch::EpochPair;
23use risingwave_pb::connector_service::SinkMetadata;
24use risingwave_pb::id::{FragmentId, PartialGraphId, SinkId};
25use risingwave_pb::stream_service::PbIcebergPkIndexSinkRole;
26use tokio::sync::mpsc;
27use tokio::sync::mpsc::{UnboundedReceiver, UnboundedSender, unbounded_channel};
28
29use crate::error::{IntoUnexpectedExit, StreamError};
30use crate::executor::exchange::permit;
31use crate::executor::monitor::StreamingMetrics;
32use crate::executor::{Barrier, BarrierInner};
33use crate::task::barrier_manager::progress::BackfillState;
34use crate::task::cdc_progress::CdcTableBackfillState;
35use crate::task::{ActorId, StreamEnvironment};
36
37/// Events sent from actors via [`LocalBarrierManager`] to [`super::barrier_worker::managed_state::PartialGraphState`].
38///
39/// See [`crate::task`] for architecture overview.
40pub(super) enum LocalBarrierEvent {
41    ReportActorCollected {
42        actor_id: ActorId,
43        epoch: EpochPair,
44    },
45    ReportCreateProgress {
46        epoch: EpochPair,
47        fragment_id: FragmentId,
48        actor: ActorId,
49        state: BackfillState,
50    },
51    ReportSourceListFinished {
52        epoch: EpochPair,
53        actor_id: ActorId,
54        table_id: TableId,
55        associated_source_id: SourceId,
56    },
57    ReportSourceLoadFinished {
58        epoch: EpochPair,
59        actor_id: ActorId,
60        table_id: TableId,
61        associated_source_id: SourceId,
62    },
63    RefreshFinished {
64        epoch: EpochPair,
65        actor_id: ActorId,
66        table_id: TableId,
67        staging_table_id: TableId,
68    },
69    RegisterBarrierSender {
70        actor_id: ActorId,
71        barrier_sender: mpsc::UnboundedSender<Barrier>,
72    },
73    RegisterLocalUpstreamOutput {
74        actor_id: ActorId,
75        upstream_actor_id: ActorId,
76        upstream_partial_graph_id: PartialGraphId,
77        term_id: String,
78        tx: permit::Sender,
79    },
80    ReportCdcTableBackfillProgress {
81        actor_id: ActorId,
82        epoch: EpochPair,
83        state: CdcTableBackfillState,
84    },
85    ReportCdcSourceOffsetUpdated {
86        epoch: EpochPair,
87        actor_id: ActorId,
88        source_id: SourceId,
89    },
90    ReportIcebergPkIndexSinkMetadata {
91        epoch: EpochPair,
92        sink_id: SinkId,
93        actor_id: ActorId,
94        role: PbIcebergPkIndexSinkRole,
95        metadata: Option<SinkMetadata>,
96    },
97}
98
99/// Can send [`LocalBarrierEvent`] to [`super::barrier_worker::managed_state::PartialGraphState::poll_next_event`]
100///
101/// See [`crate::task`] for architecture overview.
102#[derive(Clone)]
103pub struct LocalBarrierManager {
104    barrier_event_sender: UnboundedSender<LocalBarrierEvent>,
105    actor_failure_sender: UnboundedSender<(ActorId, StreamError)>,
106    pub(crate) term_id: String,
107    pub(crate) env: StreamEnvironment,
108}
109
110impl LocalBarrierManager {
111    pub(super) fn new(
112        term_id: String,
113        env: StreamEnvironment,
114    ) -> (
115        Self,
116        UnboundedReceiver<LocalBarrierEvent>,
117        UnboundedReceiver<(ActorId, StreamError)>,
118    ) {
119        let (event_tx, event_rx) = unbounded_channel();
120        let (err_tx, err_rx) = unbounded_channel();
121        (
122            Self {
123                barrier_event_sender: event_tx,
124                actor_failure_sender: err_tx,
125                term_id,
126                env,
127            },
128            event_rx,
129            err_rx,
130        )
131    }
132
133    pub fn for_test() -> Self {
134        Self::new("114514".to_owned(), StreamEnvironment::for_test()).0
135    }
136
137    /// Event is handled by [`super::barrier_worker::managed_state::PartialGraphState::poll_next_event`]
138    fn send_event(&self, event: LocalBarrierEvent) {
139        // ignore error, because the current barrier manager maybe a stale one
140        let _ = self.barrier_event_sender.send(event);
141    }
142
143    /// When a [`crate::executor::StreamConsumer`] (typically [`crate::executor::DispatchExecutor`]) get a barrier, it should report
144    /// and collect this barrier with its own `actor_id` using this function.
145    pub fn collect<M>(&self, actor_id: ActorId, barrier: &BarrierInner<M>) {
146        self.send_event(LocalBarrierEvent::ReportActorCollected {
147            actor_id,
148            epoch: barrier.epoch,
149        })
150    }
151
152    /// When a actor exit unexpectedly, it should report this event using this function, so meta
153    /// will notice actor's exit while collecting.
154    pub fn notify_failure(&self, actor_id: ActorId, err: StreamError) {
155        let _ = self
156            .actor_failure_sender
157            .send((actor_id, err.into_unexpected_exit(actor_id)));
158    }
159
160    pub fn subscribe_barrier(&self, actor_id: ActorId) -> UnboundedReceiver<Barrier> {
161        let (tx, rx) = mpsc::unbounded_channel();
162        self.send_event(LocalBarrierEvent::RegisterBarrierSender {
163            actor_id,
164            barrier_sender: tx,
165        });
166        rx
167    }
168
169    pub fn register_local_upstream_output(
170        &self,
171        actor_id: ActorId,
172        upstream_actor_id: ActorId,
173        upstream_fragment_id: FragmentId,
174        upstream_partial_graph_id: PartialGraphId,
175        metrics: Arc<StreamingMetrics>,
176    ) -> permit::Receiver {
177        let upstream_fragment_id_str = upstream_fragment_id.to_string();
178        let fragment_channel_buffered_bytes = metrics
179            .fragment_channel_buffered_bytes
180            .with_guarded_label_values(&[&upstream_fragment_id_str]);
181        let (tx, rx) = permit::channel_from_config_with_metrics(
182            self.env.global_config(),
183            permit::ChannelMetrics {
184                sender_actor_channel_buffered_bytes: fragment_channel_buffered_bytes.clone(),
185                receiver_actor_channel_buffered_bytes: fragment_channel_buffered_bytes,
186            },
187        );
188        self.send_event(LocalBarrierEvent::RegisterLocalUpstreamOutput {
189            actor_id,
190            upstream_actor_id,
191            upstream_partial_graph_id,
192            term_id: self.term_id.clone(),
193            tx,
194        });
195        rx
196    }
197
198    pub fn report_source_list_finished(
199        &self,
200        epoch: EpochPair,
201        actor_id: ActorId,
202        table_id: TableId,
203        associated_source_id: SourceId,
204    ) {
205        self.send_event(LocalBarrierEvent::ReportSourceListFinished {
206            epoch,
207            actor_id,
208            table_id,
209            associated_source_id,
210        });
211    }
212
213    pub fn report_source_load_finished(
214        &self,
215        epoch: EpochPair,
216        actor_id: ActorId,
217        table_id: TableId,
218        associated_source_id: SourceId,
219    ) {
220        self.send_event(LocalBarrierEvent::ReportSourceLoadFinished {
221            epoch,
222            actor_id,
223            table_id,
224            associated_source_id,
225        });
226    }
227
228    pub fn report_refresh_finished(
229        &self,
230        epoch: EpochPair,
231        actor_id: ActorId,
232        table_id: TableId,
233        staging_table_id: TableId,
234    ) {
235        self.send_event(LocalBarrierEvent::RefreshFinished {
236            epoch,
237            actor_id,
238            table_id,
239            staging_table_id,
240        });
241    }
242
243    pub fn report_cdc_source_offset_updated(
244        &self,
245        epoch: EpochPair,
246        actor_id: ActorId,
247        source_id: SourceId,
248    ) {
249        self.send_event(LocalBarrierEvent::ReportCdcSourceOffsetUpdated {
250            epoch,
251            actor_id,
252            source_id,
253        });
254    }
255
256    pub fn report_iceberg_pk_index_sink_metadata(
257        &self,
258        epoch: EpochPair,
259        sink_id: risingwave_common::id::SinkId,
260        actor_id: ActorId,
261        role: PbIcebergPkIndexSinkRole,
262        metadata: Option<SinkMetadata>,
263    ) {
264        self.send_event(LocalBarrierEvent::ReportIcebergPkIndexSinkMetadata {
265            epoch,
266            sink_id,
267            actor_id,
268            role,
269            metadata,
270        });
271    }
272}
273
274#[cfg(test)]
275impl LocalBarrierManager {
276    pub(super) fn spawn_for_test()
277    -> crate::task::barrier_worker::EventSender<crate::task::barrier_worker::LocalActorOperation>
278    {
279        use std::sync::Arc;
280        use std::sync::atomic::AtomicU64;
281
282        use crate::executor::monitor::StreamingMetrics;
283        use crate::task::barrier_worker::{EventSender, LocalBarrierWorker};
284
285        let (tx, rx) = unbounded_channel();
286        let _join_handle = LocalBarrierWorker::spawn(
287            StreamEnvironment::for_test(),
288            Arc::new(StreamingMetrics::unused()),
289            None,
290            Arc::new(AtomicU64::new(0)),
291            rx,
292        );
293        EventSender(tx)
294    }
295}