risingwave_stream/executor/
barrier_align.rs1use std::sync::Arc;
16use std::time::Instant;
17
18use anyhow::Context;
19use await_tree::InstrumentAwait;
20use enum_as_inner::EnumAsInner;
21use futures::StreamExt;
22use futures::future::{Either, select};
23use futures_async_stream::try_stream;
24use risingwave_common::bail;
25
26use super::error::StreamExecutorError;
27use super::{Barrier, BoxedMessageStream, Message, StreamChunk, StreamExecutorResult, Watermark};
28use crate::executor::monitor::StreamingMetrics;
29use crate::task::{ActorId, FragmentId};
30
31pub type AlignedMessageStreamItem = StreamExecutorResult<AlignedMessage>;
32pub trait AlignedMessageStream = futures::Stream<Item = AlignedMessageStreamItem> + Send;
33
34#[cfg_attr(any(test, feature = "test"), derive(PartialEq))]
35#[derive(Debug, EnumAsInner)]
36pub enum AlignedMessage {
37 Barrier(Barrier),
38 WatermarkLeft(Watermark),
39 WatermarkRight(Watermark),
40 Left(StreamChunk),
41 Right(StreamChunk),
42}
43
44#[try_stream(ok = AlignedMessage, error = StreamExecutorError)]
45pub async fn barrier_align(
46 mut left: BoxedMessageStream,
47 mut right: BoxedMessageStream,
48 actor_id: ActorId,
49 fragment_id: FragmentId,
50 metrics: Arc<StreamingMetrics>,
51 executor_name: &str,
52) {
53 let actor_id = actor_id.to_string();
54 let fragment_id = fragment_id.to_string();
55 let left_barrier_align_duration = metrics.barrier_align_duration.with_guarded_label_values(&[
56 actor_id.as_str(),
57 fragment_id.as_str(),
58 "left",
59 executor_name,
60 ]);
61 let right_barrier_align_duration = metrics.barrier_align_duration.with_guarded_label_values(&[
62 actor_id.as_str(),
63 fragment_id.as_str(),
64 "right",
65 executor_name,
66 ]);
67 loop {
68 let prefer_left: bool = rand::random();
69 let select_result = if prefer_left {
70 select(left.next(), right.next()).await
71 } else {
72 match select(right.next(), left.next()).await {
73 Either::Left(x) => Either::Right(x),
74 Either::Right(x) => Either::Left(x),
75 }
76 };
77 match select_result {
78 Either::Left((None, _)) => {
79 while let Some(msg) = right.next().await {
81 match msg? {
82 Message::Watermark(watermark) => {
83 yield AlignedMessage::WatermarkRight(watermark)
84 }
85 Message::Chunk(chunk) => yield AlignedMessage::Right(chunk),
86 Message::Barrier(_) => {
87 bail!("right barrier received while left stream end");
88 }
89 }
90 }
91 break;
92 }
93 Either::Right((None, _)) => {
94 while let Some(msg) = left.next().await {
96 match msg? {
97 Message::Watermark(watermark) => {
98 yield AlignedMessage::WatermarkLeft(watermark)
99 }
100 Message::Chunk(chunk) => yield AlignedMessage::Left(chunk),
101 Message::Barrier(_) => {
102 bail!("left barrier received while right stream end");
103 }
104 }
105 }
106 break;
107 }
108 Either::Left((Some(msg), _)) => match msg? {
109 Message::Watermark(watermark) => yield AlignedMessage::WatermarkLeft(watermark),
110 Message::Chunk(chunk) => yield AlignedMessage::Left(chunk),
111 Message::Barrier(barrier) => loop {
112 let start_time = Instant::now();
113 match right
115 .next()
116 .instrument_await(await_tree::span!(
117 "barrier_align_wait_right epoch={}",
118 barrier.epoch.curr
119 ))
120 .await
121 .context("failed to poll right message, stream closed unexpectedly")??
122 {
123 Message::Watermark(watermark) => {
124 yield AlignedMessage::WatermarkRight(watermark)
125 }
126 Message::Chunk(chunk) => yield AlignedMessage::Right(chunk),
127 Message::Barrier(barrier) => {
128 yield AlignedMessage::Barrier(barrier);
129 right_barrier_align_duration
130 .inc_by(start_time.elapsed().as_nanos() as u64);
131 break;
132 }
133 }
134 },
135 },
136 Either::Right((Some(msg), _)) => match msg? {
137 Message::Watermark(watermark) => yield AlignedMessage::WatermarkRight(watermark),
138 Message::Chunk(chunk) => yield AlignedMessage::Right(chunk),
139 Message::Barrier(barrier) => loop {
140 let start_time = Instant::now();
141 match left
143 .next()
144 .instrument_await(await_tree::span!(
145 "barrier_align_wait_left epoch={}",
146 barrier.epoch.curr
147 ))
148 .await
149 .context("failed to poll left message, stream closed unexpectedly")??
150 {
151 Message::Watermark(watermark) => {
152 yield AlignedMessage::WatermarkLeft(watermark)
153 }
154 Message::Chunk(chunk) => yield AlignedMessage::Left(chunk),
155 Message::Barrier(barrier) => {
156 yield AlignedMessage::Barrier(barrier);
157 left_barrier_align_duration
158 .inc_by(start_time.elapsed().as_nanos() as u64);
159 break;
160 }
161 }
162 },
163 },
164 }
165 }
166}
167
168#[cfg(test)]
169mod tests {
170 use std::time::Duration;
171
172 use async_stream::try_stream;
173 use futures::{Stream, TryStreamExt};
174 use risingwave_common::array::stream_chunk::StreamChunkTestExt;
175 use risingwave_common::util::epoch::test_epoch;
176 use tokio::time::sleep;
177
178 use super::*;
179
180 fn barrier_align_for_test(
181 left: BoxedMessageStream,
182 right: BoxedMessageStream,
183 ) -> impl Stream<Item = Result<AlignedMessage, StreamExecutorError>> {
184 barrier_align(
185 left,
186 right,
187 0.into(),
188 0.into(),
189 Arc::new(StreamingMetrics::unused()),
190 "dummy_executor",
191 )
192 }
193
194 #[tokio::test]
195 async fn test_barrier_align() {
196 let left = try_stream! {
197 yield Message::Chunk(StreamChunk::from_pretty("I\n + 1"));
198 yield Message::Barrier(Barrier::new_test_barrier(test_epoch(1)));
199 yield Message::Chunk(StreamChunk::from_pretty("I\n + 2"));
200 yield Message::Barrier(Barrier::new_test_barrier(test_epoch(2)));
201 }
202 .boxed();
203 let right = try_stream! {
204 sleep(Duration::from_millis(1)).await;
205 yield Message::Chunk(StreamChunk::from_pretty("I\n + 1"));
206 yield Message::Barrier(Barrier::new_test_barrier(test_epoch(1)));
207 yield Message::Barrier(Barrier::new_test_barrier(test_epoch(2)));
208 yield Message::Chunk(StreamChunk::from_pretty("I\n + 3"));
209 }
210 .boxed();
211 let output: Vec<_> = barrier_align_for_test(left, right)
212 .try_collect()
213 .await
214 .unwrap();
215 assert_eq!(
216 output,
217 vec![
218 AlignedMessage::Left(StreamChunk::from_pretty("I\n + 1")),
219 AlignedMessage::Right(StreamChunk::from_pretty("I\n + 1")),
220 AlignedMessage::Barrier(Barrier::new_test_barrier(test_epoch(1))),
221 AlignedMessage::Left(StreamChunk::from_pretty("I\n + 2")),
222 AlignedMessage::Barrier(Barrier::new_test_barrier(2 * test_epoch(1))),
223 AlignedMessage::Right(StreamChunk::from_pretty("I\n + 3")),
224 ]
225 );
226 }
227
228 #[tokio::test]
229 #[should_panic]
230 async fn left_barrier_right_end_1() {
231 let left = try_stream! {
232 sleep(Duration::from_millis(1)).await;
233 yield Message::Chunk(StreamChunk::from_pretty("I\n + 1"));
234 yield Message::Barrier(Barrier::new_test_barrier(test_epoch(1)));
235 }
236 .boxed();
237 let right = try_stream! {
238 yield Message::Chunk(StreamChunk::from_pretty("I\n + 1"));
239 }
240 .boxed();
241 let _output: Vec<_> = barrier_align_for_test(left, right)
242 .try_collect()
243 .await
244 .unwrap();
245 }
246
247 #[tokio::test]
248 #[should_panic]
249 async fn left_barrier_right_end_2() {
250 let left = try_stream! {
251 yield Message::Chunk(StreamChunk::from_pretty("I\n + 1"));
252 yield Message::Barrier(Barrier::new_test_barrier(test_epoch(1)));
253 }
254 .boxed();
255 let right = try_stream! {
256 sleep(Duration::from_millis(1)).await;
257 yield Message::Chunk(StreamChunk::from_pretty("I\n + 1"));
258 }
259 .boxed();
260 let _output: Vec<_> = barrier_align_for_test(left, right)
261 .try_collect()
262 .await
263 .unwrap();
264 }
265}