9 ms·
Streaming here has a different meaning than perhaps what you're used to. It's not referring to online processing where you maintain aggregates/state while an en
by orlp 13d ago
Streaming here has a different meaning than perhaps what you're used to. It's not referring to online processing where you maintain aggregates/state while an endless stream of data comes in.
The name was chosen early on to contrast with the old execution model, which was essentially all-data-in-memory, column-at-a-time. That engine still exists, we use it as a fallback mechanism for things that aren't supported yet in the new engine (or if you explicitly ask for `engine="in-memory"`).
The new execution model first constructs a computational graph of nodes which communicate in streams of in-cache batches (morsels) of data, meaning the full dataset will never be held in memory if not necessary. This was called the streaming engine for that reason in an early prototype and the name stuck. In hindsight I do admit the naming choice is somewhat confusing.
- arn3n 13d agoCool, thanks for the explanation!
- sanderjd 13d agoWhen you say "in-cache batches", you mean that this cache is on disk? Is that only the case when data is quite large? (Or a more general question: What is the best resource for me to read about how the streaming engine and cache work?)
- orlp 13d agoWell... once my recent work on out-of-core lands the batch could be on disk when we run out of memory budget ;) But no, that's not what I meant. I meant that the batch is meant to be of a size that fits in your CPU cache. This can be a huge throughput improvement as each bit of data stays in cache as it moves from data source to sink. Compare this to column-at-a-time execution: by the time you start the next operation on this column the start of the column will be out of cache again, meaning you operate at RAM speed (or worse, disk speed) rather than cache speed. I gave a (fairly surface-level) talk on the streaming engine a bit over a year ago: https://pola.rs/posts/talk-polars-meetup-1-streaming-engine/ https://pola.rs/posts/talk-polars-meetup-1-streaming-engine/.
- sanderjd 13d agoAaaaah, the cpu cache aspect is what I was missing! This makes tons of sense! I'll check out your talk as well.