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§
Sourcefn flatten(self, name: &str) -> StreamBuilder<K, VO, T>
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".