Skip to main content

risingwave_stream/executor/
barrier_align.rs

1// Copyright 2022 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
15use 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                // left stream end, passthrough right chunks
80                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                // right stream end, passthrough left chunks
95                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                    // received left barrier, waiting for right barrier
114                    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                    // received right barrier, waiting for left barrier
142                    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}