risingwave_stream/task/barrier_manager/
mod.rs1pub 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
37pub(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#[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 fn send_event(&self, event: LocalBarrierEvent) {
139 let _ = self.barrier_event_sender.send(event);
141 }
142
143 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 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}