Performance: Adventures in Thread-per-Core Async with Redpanda and Seastar
infoq.com
infoq.com
An issue is that there are few libraries that are designed or optimized for thread-per-core architectures. Storage management is a reusable concept in these architectures but there is a lack of competent and scalable libraries that do this optimally. For someone coming into the space, you have to learn how to design e.g. high-performance storage allocators that can keep up with a JBOD of NVMe. I have mature libraries like this but they are tightly coupled to the rest of the software I design.
When I first designed these architectures I sharded up resources across cores but there are many cases where this requires giving up a lot of performance. The problem of shedding load and resources across cores is really interesting. An implication is that some structures should be shared globally but you really want the performance to be as close to contention-free as possible or you defeat the purpose of thread-per-core. This is possible and there are design heuristics that approximate it but it isn’t discussed much.
As mentioned in the article, no matter how you build these things your software will require a sophisticated understanding of lifetimes. Sometimes the builtin tools of the language will help with this, other times they can’t and you will have to build your own tools.
I have never used Seastar. Not because I am not familiar with it but because it makes tradeoffs that are not appropriate for the workloads I tend to support. It is a good choice for many workloads. The only thing this indicates is that we are still in the early days of these types of architectures. Even if you don’t care or need to improve your architectural efficiency, these types of architectures significantly reduce the carbon/cost footprint of software systems.
I’ve only dabbled with them on hobby projects, but I absolutely love the idea. I’ve been on the hunt for more content about this, but I have trouble finding much, I suspect a lot of content is a bit “if you know you know” sort of insider knowledge stuff, so if you know of any other good content lurking out there I’d be so keen to know.
> this is the article for you. I’ve designed thread-per-core architectures for 15 years now
How’d you get into this/what do you do for work? As I said, I’ve dabbled with it on personal projects, but at least on everything I do for work it’s not always a good fit, or it is, but it’s too weird for my teammates.
It is, the kind of knowledge that takes decades to accumulate you won't see on people's blog posts. This is why ultra low-latency and high-performance systems engineers can ask more for their services.
> How’d you get into this/what do you do for work
I'm not the OP but for HFT shops thread-per-core design is the only right way. Prop trading, big market makers etc, we all go to great lengths to isolate the workloads as much as possible.
I really want the next generation NodeJS equivalent to take advantage of all these learnings and learnings from Erlang. I think nodejs showed how awesome event loops and async can be and simple.
The TFA (The Friendly Article) makes it clear there are sharp edges with coroutines and lambdas in C++:
- something stored in a temporary between a yield barrier will be destroyed by the time the coroutine is resumed
- you have to make sure you store things inside the coroutine promise object
- lambda capturing rule complexity associated with copying and moving semantics
I am just learning but I've written a multithreaded barrier in C and Java that is mutex free. In 6×2 thread pairs it can communicate between threads in 42 nanoseconds.
I have the beginnings of a JIT compiler. But there's so much work to do.
My barrier resembles a "phaser" pattern in a loop and bulk synchronous parallel.
Specifically our device/firmware is very networking/IO and runtime-configurability focused, which makes for some fun challenges. Squeezing the most out of the microcontroller and peripherals has been important, and keeping data/code locality has helped lots.
I wonder if there is a level inside of modern CPU cores where cache is shared by different parallel pipelines or sub-caches for each.
It's sort of fractal. Moving the data around is the bottleneck, so the most efficient way is to arrange the data so that data movement is minimized, and that applies at different zoom levels.
I suspect that there is a lot of research in memory-optimized computing that are possibly different enough from typical approaches to be considered new paradigms.
I wonder if this could intersect with something like Mixture of Experts in LLMs. I don't actually know how that works, but maybe it could allow for grouping the related expert weight data closer together to speed things up.
Also, think about the complexity of real neurons when compared to something like those in an MLP. My impression is that real neurons pack in (locally) much more information and compute per unit.
I wonder if you could use some type of model to predict what data you will need for text generation (aside from the KV cache of course), and then have a fast cache for each task.
Is there a large open LLM or LMM trained using something like https://github.com/lucidrains/st-moe-pytorch ?
Actors are asynchronous objects with a message queue, designed to run on a thread pool. You can pin them to a core, or allow migration via work stealing (to spread the workload). Did you know that in modern many core systems, a hot CPU core will halt to cool down? Work stealing will offload non-pinned actors, balancing the load.
There is one stumbling issue with Actors which is “sync” work, a transaction which requires one actor to update anothers state before continuing. This can be resolved by “locking” an actor, but this mechanism is consceptually “dirty”, but solves real world problems.
FWIW, work stealing is often not recommended in thread-per-core architectures because it introduces quite a bit of unnecessary thread contention. You are correct that load balancing is a central problem in these architectures but it is usually achieved by shedding data since that can be done with minimal locking and inter-thread coordination. This moves the problem to figuring out what data to shed but this has satisfactory inexpensive solutions in many system designs.
Now to be fair, a lot of actor systems were about correctness and not squeezing the most performance of a platform. Not sharing data means no explicit locking with mutexes (though one can still have deadlocks, as in several actors stuck waiting for each others) and something simpler to analyze: the system is made of communicating state machines. Actors were also used in high performance system like telecom switches (high performance, for their time), but here too correctness was probably the main concern. At the time actors were first used accessing even far memory was not as costly as today, so the cost of synchronization wasn't as bad. But it was already tricky to get right.
Still, actors are perfectly on topic when discussing a shared nothing, message based architecture. They're worth mentioning IMHO, to give some historical perspective. And then one can combine both approaches: a pinned thread per core can dispatch its messages to a set of cooperating actors for example.
I was just trying to answer your question on why agents could be relevant here, nothing more.
There is a core to core visualiser here.
That being said, because hyperthreaded workloads share a pipeline and a cache, there might be benefit for memory-constrained applications to pinning pairs of processes to the 2 logical cores on the same physical core if it's highly queue-like and you can process the data in similar numbers of instructions for codestream.