Skip to main content

Map

Trait Map 

Source
pub trait Map<K, V, T, VO>: Sealed {
    // Required method
    fn map(
        self,
        name: &str,
        mapper: impl FnMut(V) -> VO + 'static,
    ) -> StreamBuilder<K, VO, T>;
}
Expand description

Apply a function to every message in a stream

Required Methods§

Source

fn map( self, name: &str, mapper: impl FnMut(V) -> VO + 'static, ) -> StreamBuilder<K, VO, T>

Map transforms every value in a datastream into a different value by applying a given function or closure.

§Example
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(0..100)))
        .map("map", |x| x * 2)
        .sink("sink", StatelessSink::new(sink_clone));
    })
    .execute()
    .unwrap();

let expected: Vec<i32> = (0..100).map(|x| x * 2).collect();
let out: Vec<i32> = sink.into_iter().map(|x| x.value).collect();
assert_eq!(out, expected);

Dyn Compatibility§

This trait is not dyn compatible.

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

Implementors§

Source§

impl<K, V, T, VO> Map<K, V, T, VO> for StreamBuilder<K, V, T>
where K: MaybeKey, V: Data, VO: Data, T: Timestamp,