5 ms·
Interesting. I have the opposite intuition. With a push-based model you end up needing some sort of back-pressure mechanism which gets complicated in a hurry. W
by thinkharderdev 3y ago
Interesting. I have the opposite intuition. With a push-based model you end up needing some sort of back-pressure mechanism which gets complicated in a hurry. With a pull-based model you get back pressure "for free"
- necubi 3y agoYeah, you definitely need backpressure, and doing that well is had. There's a great blog post about how Flink manages this [0] and it involves a complex credit-based flow control system. But in practice the performance advantages of allowing operators to do work as they have capacity (rather than being potentially starved by the capacity of their downstreams) outweighs this. [0] https://flink.apache.org/2019/06/05/a-deep-dive-into-flinks-network-stack/ https://flink.apache.org/2019/06/05/a-deep-dive-into-flinks-...
- scott_s 3y agoI worked on a push-based streaming system which could arbitrarily distribute each operate to another compute node. Backpressure occurred naturally: if the thing you’re pushing into is busy (internal queue is full; TCP socket is busy) you wait. We had no explicit backpressure mechanism, and yet we had it in our system.
- necubi 3y agoThis is an interesting point. In Arroyo [0] (the streaming engine I work on) we also use a simple TCP-based backpressure mechanism: we have NxM TCP connections between each operator subtask, behind small in-memory queues. This essentially is relying on the kernel's flow control and backpressure mechanisms instead of building a custom network stack like Flink does. This has worked well in practice, but I think does leave some performance on the table and likely incurrs higher resource utilization. Flink for example uses a system of "floating" buffers that can be moved between streams as needed to minimize the total number of buffers needed. However it may be that Linux's network stack has gotten good enough at this point that it makes sense to just rely on it directly. I wish there was more research out there on how these sorts of network stacks should be built on modern hardware and OSes. (semi-relatedly, was checking out your CV; a bunch of great looking papers there that I'll definitely be checking out!) [0] https://github.com/ArroyoSystems/arroyo https://github.com/ArroyoSystems/arroyo
- scott_s 3y agoInteresting - I looked into your code a bit. I found your window aggregation library [1]. You may be interested in looking into the Rust implementation of some of the research work I've been a part of [2]. In Flink, I believe the reason they need to implement their own backpressure system is that they multiplex TCP connections. That is, they have multiple logical streams flowing through a single TCP connection. If that's the case, you need to do some work to 1) detect which logical stream is the one that's blocking, and 2) don't block because other logical streams may be able to use the active TCP connection. Thinking it through, I think what Flink's approach buys is not necessarily better performance, but better just a manageable number of connections. That is, imagine you have a process P1 with operators A, B and C. And then P2 has D, E, F. Now imagine that this is a shuffle, where A, B and C are fully connected to D, E and F. In my old system, you would have 9 TCP connections. In Flink, you will have 1. [1] https://github.com/ArroyoSystems/arroyo/blob/master/arroyo-worker/src/operators/aggregating_window.rs https://github.com/ArroyoSystems/arroyo/blob/master/arroyo-w... [2] https://github.com/IBM/sliding-window-aggregators/tree/master/rust https://github.com/IBM/sliding-window-aggregators/tree/maste...
- necubi 3y agoThanks for that link! We've been heavily inspired by the work out of TU Berlin on Scotty (https://tu-berlin-dima.github.io/scotty-window-processor/ https://tu-berlin-dima.github.io/scotty-window-processor/) but exciting to see more work on this problem. Our efficient window handling is a big area of improvement over what Flink is able to do right now. For Yep, Flink has to do this themselves because they multiplex. My sense is that when Flink was originally designed in ~2012 Linux didn't support large numbers of TCP connections well (and memory was more of an issue), so they built a network stack that multiplexes many logical streams on top of a single TCP connection. My (not terribly informed) opinion is that TCP behaviors and semantics are not necessarily optimal for these systems and you get benefits from allowing your systems' scheduler to control prioritization rather than handing it off the kernel. There are also a new class of datacenter-oriented protocols like Homa (https://homa-transport.atlassian.net/wiki/spaces/HOMA/overview https://homa-transport.atlassian.net/wiki/spaces/HOMA/overvi...) that I think are pretty interesting to solve this class of problems.