Aggregate streaming data in real-time with WebAssembly
infinyon.com
infinyon.com
"Aggregates let you define functions that combine each record in a stream with some long-running state, or 'accumulator'."
In Rusty pseudocode, reduce requires a function with two inputs and an output of the same type:
fn reduce<T>(f: Fn(T, T) -> T)
Whereas fold may use one type for the accumulator and another type for the elements, but requires an initial accumulator value to be given explicitly:
fn fold<A, T>(init: A, f: Fn(A, T) -> A)
The Aggregate SmartStreams discussed in the blog follow this fold pattern, applied to a distributed persistent log as the stream and using WebAssembly modules as the functions.
https://developer.apple.com/documentation/swift/array/229868...
Although, I dont understand whats the value add of wasm (apart from security) if the user still has to write code in Rust -> wasm. Why not just execute in rust alone?
You can compile almost any language to WASM not just Rust. For example, Python, Go, Javascript: https://github.com/appcypher/awesome-wasm-langs.
Taking the examples from the article, equivalent schemas might be:
type: object
properties:
mySum:
type: number
reduce: { strategy: sum }
Or even: type: object
properties:
myDeeplyAggregatedMap:
type: object
reduce: { strategy: merge }
additionalProperties:
reduce: { strategy: sum }
Use of WASM is really interesting. We've been exploring it as a means for powering user-defined reduction strategies, in cases where the built-in strategies are insufficient.[1] https://estuary.dev [2] https://docs.estuary.dev/reference/catalog-reference/schemas...
Security is certainly one of the reasons to use WASM, the ability to run it in a sandbox means that untrusted user code can be uploaded to Fluvio's Streaming Processing Units and do the processing inline, rather than on the client side. This can save big on network bandwidth, especially with a dataset where filtering whittles down a lot on volume.
Other reasons include that WASM is a fast and portable bytecode format and that there is very good tooling and support for compiling Rust to WASM as well as embedding WASM runtimes in Rust, which works well for Fluvio as it's written in Rust.
Here's another post with a bit more detail about some of the design and motivational factors if you're interested: https://www.infinyon.com/blog/2021/06/introducing-fluvio/#fl...
CQRS reactive patterns with Flink or Spark computing a result to be send to the client could benefit from it: you could decide to move some aggregation client-side in the same business language.