risingwave_stream::executor::source

Function barrier_to_message_stream

source
pub fn barrier_to_message_stream(
    rx: UnboundedReceiver<Barrier>,
) -> impl Stream<Item = Result<Message, StreamExecutorError>>
Expand description

Receive barriers from barrier manager with the channel, error on channel close.