Skip to main content

Flatten

Trait Flatten 

Source
pub trait Flatten<K, VI, T, VO, I>: Sealed {
    // Required method
    fn flatten(self, name: &str) -> StreamBuilder<K, VO, T>;
}
Expand description

Flatten a stream of iterables by emitting each element of every iterable as a distinct message.

Required Methods§

Source

fn flatten(self, name: &str) -> StreamBuilder<K, VO, T>

Flatten a datastream. Given a stream of some iterables, this function consumes each iterable and emits each of its elements downstream.

§Key and Time

If the message containing the iterator has a key or timestamp, they are cloned and attached to every emitted message.

§Example

Only retain numbers <= 42

use malstrom::operators::*;
use malstrom::runtime::SingleThreadRuntime;
use malstrom::snapshot::NoPersistence;
use malstrom::sources::{SingleIteratorSource, StatelessSource};
use malstrom::worker::StreamProvider;
use malstrom::sinks::{VecSink, StatelessSink};

let sink = VecSink::new();
let sink_clone = sink.clone();

SingleThreadRuntime::builder()
    .persistence(NoPersistence)
    .build(move |provider: &mut dyn StreamProvider| {
        provider.new_stream()
        .source("numbers", StatelessSource::new(
            SingleIteratorSource::new([vec![1, 2, 3], vec![4, 5], vec![6]])
        ))
        .flatten("flatten")
        .sink("sink", StatelessSink::new(sink_clone));
    })
    .execute()
    .unwrap();

let expected: Vec<i32> = vec![1, 2, 3, 4, 5, 6];
let out: Vec<i32> = sink.into_iter().map(|x| x.value).collect();
assert_eq!(out, expected);

Dyn Compatibility§

This trait is dyn compatible.

In older versions of Rust, dyn compatibility was called "object safety".

Implementors§

Source§

impl<K, VI, T, VO, I> Flatten<K, VI, T, VO, I> for StreamBuilder<K, VI, T>
where K: MaybeKey, I: Iterator<Item = VO>, VI: IntoIterator<Item = VO, IntoIter = I> + Data, VO: Data, T: Timestamp,