Scipio: A Thread-per-Core Crate for Rust and Linux
datadoghq.com
datadoghq.com
While it increases the data locality, I have seen a few software following this sharding model (notably scylla) that work really bad once the load is not evenly distributed across all shards
When that happens it can be a huge waste of resources and can give lower performance (depending on the type of load)
Imho unless you are absolutely sure about the type of load, leave the sharding to dividing data between servers, or have some mechanism that can shift to sharing the load between threads if the system imbalance is too great
Serial access to hot keys is a hard thing to design around, and you're right that sharding doesn't solve that problem. Worse, it exposes other keys that just happen to share the shard (or the core) to poor performance.
There are a couple of well-understood solutions to this problem. The obvious one is to dynamically re-balance shards based on heat, either moving some keys or split/merge. This is the same tactic that many distributed databases use, and while it's complex to do, it's easier to do locally than distributed. Another option is stochastic re-balancing, like the Stochastic Fairness Queuing (https://ieeexplore.ieee.org/document/91316) model used in networking. Here, shards are randomly re-shuffled occasionally. Doesn't fix the noisy neighbor problem, but does mean that the noisyness moves around. That might seem silly, but it's pretty much what's going to happen under the covers of the non-sharded version of the code when the scheduler gets involved, only easier to reason about.
> When that happens it can be a huge waste of resources and can give lower performance (depending on the type of load)
Lower apparent performance for neighbors of the heavy hitter, sure. Under which other circumstances does it reduce performance?
> Imho unless you are absolutely sure about the type of load, leave the sharding to dividing data between servers, or have some mechanism that can shift to sharing the load between threads if the system imbalance is too great
I'm a bit puzzled by this. Distributed systems have exactly the same problem, and solving that problem is much harder there because the cost of contention is higher, data movement is more expensive, and you have to deal with a lot more failure cases.
The statistics may be better for distributed systems because a hot tenant has to be a lot hotter to make hot box than a hot core. But that's a very specific kind of bet, and if you end up with a tenant that does cause a hot box you have an even harder problem to solve.
Dynamo (https://www.allthingsdistributed.com/files/amazon-dynamo-sos...) solves this by moving the key space around, as do many similar kinds of systems. It's not easy, though, and filled with caveats. If you're scared of sharding, distributed sharding should be scarier than on-box sharding.
https://docs.scylladb.com/operating-scylla/nodetool-commands...
It is possible to design data sharding architectures where balancing of both data and load across cores is continuous and smooth, with surprisingly minimal overhead and coordination cost. In fact, there are multiple ways of doing it, depending on your workload characteristics. At this point, the design idioms for this style of architecture are refined and robust, so there is no computer science reason the problems you raise need to exist in a real implementation. The reason it seems "difficult" in practice is because so many designs insist on loosely coupling and weakly scheduling storage, execution, network, etc. If your architecture concept is slapping a thin layer on top of RocksDB, it won't be feasible. Every part of the stack needs to understand and be designed to the model. The end result is actually quite elegant in my view, and with unmatched throughput.
People design distributed systems with simple static sharding schemes because they are obvious, easy, and it lets you cut a lot of corners on the rest of your system design. It is not the only way to design a distributed system, continuous adaptive resharding and load shedding is a demonstrably viable option, and it is much easier to implement within a single server than on an actual cluster of networked computers.
Would you be willing to point us toward at least the zipcode of some citations which could outline these refined and robust design idioms?
Methods and apparatus for optimizing resource utilization in distributed storage systems
https://patents.google.com/patent/US9990147B2/en?inventor=Ja...
Load rebalancing for shared resource
https://patents.google.com/patent/US8539197B1/en?q=monitorin...
Amazon Aurora: Design Considerations for High Throughput Cloud-Native Relational Databases
"Since our system has a high tolerance to failures, we can leverage this for maintenance operations that cause segment unavailability. For example, heat management is straightforward. We can mark one of the segments on a hot disk or node as bad, and the quorum will be quickly repaired by migration to some other colder node in the fleet."
http://news.cs.nyu.edu/~jinyang/ds-reading/aurora-sigmod17.p...
The parent developed a high perf distributed GIS database. I assume you knew this.
If the number of connections is not small, Scylla will crank up any other implementation with traditional locking.
Async/await will make complexity explode because of the colored function problem [1].
The solution to expensive context switches is cheap context switches, plain and simple. User-mode lightweight threads like go's, or upcoming Java's with Loom [2] have proven that this is possible.
Yes, it does mean that it can only happen in a language that controls its stack (so that you can slice it off and pop a continuation on it). I sincerely believe this is Rust's ballpark; hell they even started the project with that idea in mind.
[1] https://journal.stuffwithstuff.com/2015/02/01/what-color-is-...
> The solution to expensive context switches is cheap context switches, plain and simple.
Except the performance difference between a kernel-mode context switch and a user-mode one is only going to narrow in the future. The overhead that cannot be eliminated from context switches is their effect on the cache, since you start to run into the laws of physics at that point...
The real solution to expensive context switches is to just do fewer of them... No context switch is always faster than a "fast" context switch.
> I sincerely believe this is Rust's ballpark
I think it's plausible that Rust could get a library-level solution for fibers that does not rely on unstable details of the compiler. Rust will never again have that baked in to the language, as it would make the language completely unsuitable for many low-level tasks.
Fibers, especially the way they are implemented in Go, come with a lot of their own complexity.
Just look at this issue: https://marcan.st/2017/12/debugging-an-evil-go-runtime-bug/
This is just one of the segfaults caused by Go's complex stack control. I don't want to rely on a runtime that contains these sorts of bugs, and the best way to avoid that is to avoid having a runtime in the first place.
Blocking an OS thread as a mean to be compatible is not exactly what we're trying to do here.
> Except the performance difference between a kernel-mode context switch and a user-mode one is only going to narrow in the future
OS overhead can be minimized, but program stacks are a function of the language's. And if you're not right sizing preemption points in your stack, you'll be switching large parts of it. This means you _must_ have stackful coroutines if you want to keep switching threads.
> The real solution to expensive context switches is to just do fewer of them... No context switch is always faster than a "fast" context switch.
Sure, but writing the perfect assembly and using gotos has always been the fastest. Abstraction has a cost, and some runtimes/languages are currently proving that they can reduce this cost to zero in the current conditions of IO being much costlier than a few 100s of nanos. We're just happening to be at a time where the compiler is starting to be smarter than the user. But I guess the benchmarks will settle all this.
> Just look at this issue: https://marcan.st/2017/12/debugging-an-evil-go-runtime-bug/
So that's a bug in the go compiler. They can either fix it or pay a "a small speed penalty (nanoseconds)" as a workaround, which the author qualifies as "acceptable".
Yes, that's not the absolute performance possible. But why care about that? At some point it all comes down to TCO (except for latencies in HFT); and TCO tells you that it's ok. Development complexity and maintainability matters. Especially when you can max out your IO usage for the decades to come.
It's what Go does whenever you call into a C function or make a system-call. As long as those blocking functions are not the bottleneck then it works fine.
My problem with the "colored function" analogy is that it implies that the problem is somehow due to the surface syntax, when in reality the problem still exists in all languages that support procedural IO: some of those languages like to just pretend that the problem doesn't exist.
The only language I'm aware of which truly solves that is Haskell, since all IO happens via a monad.
> Yes, that's not the absolute performance possible. But why care about that?
This point was not about performance. It's about the pitfalls of writing all of your code on top of a complex and buggy runtime.
Programming is a lot simpler, and development is a lot faster, when I don't have to worry about that.
There are also several things that are a lot more complicated when you bake a complex runtime into the language like Go does. Thread-local storage is completely broken for one. If you do any kind of GUI programming, you may need to use `runtime.LockOSThread` as most GUIs expect function calls from a single thread. etc. etc.
Source: https://devblogs.microsoft.com/dotnet/configureawait-faq/
Anyway, in both models you can have deadlocks. And even if there is no deadlock, blocking the eventloop is still an antipattern, since it prevents other tasks which might be able to make progress from running.
The calling thread immediately yields until the blocking work task completes and wakes it up again.
First, in Rust this isn't really a problem. You can always turn async calls into blocking ones in Rust by calling block_on [0]. In some languages block_on doesn't exist, like in in-browser js, because here, code is supposed to be async. But in Rust there is no requirement, so there's no colored function problem here.
Second, I don't think it's a big problem in the first place. In one of my projects, I'm using an async library and have isolated the async-ness by creating a dedicated thread that communicates with the library. The thread provides a queue of messages that the remaining code of my project can handle.
[0]: https://docs.rs/tokio/0.3.3/tokio/runtime/struct.Runtime.htm...
A viral one, and number-of-possible-states-exploding one at that.
> First, in Rust this isn't really a problem. You can always turn async calls into blocking ones in Rust by calling block_on
That's not exactly what we're trying to do here.
I truly don't care about this issue...
Of course you can add async to the whole callstack, but it could be third party code and it might require code duplication if async adds a penalty to compared to non-async code.
Ideally the fixed stack size/conversion to state machine would be an optimization that the compiler would apply if it can prove that the coroutine ever yields form top level (or from a well known and fixed stack depth) and resort to dynamic stacks otherwise. I have been thinking a lot about this, and I think the key is reifying the incoming continuation and, as long as it doesn't escape the called coroutine , the optimization can be guaranteed. I believe that rust lifetime machinery might help, but it is something I'm not familiar with.
edit: it also requires heap allocating activation frames in the most general case, which is slow.
async func foo(x int) task[int] { return await(2 * x); }
// No "async" only because cpu_bound_future_factory() returns a task
func bar(y int) task[int] { return cpu_bound_future_factory(async func() task[int] { return await(y - 5); }); }
async func baz(a,b int) task[int] { return await(await foo(a) + await bar(b)); }
is not fun, especially when you have to await the result of the arithmetic operators, and write all those now-redundant async and await, so let's also drop "async" and make "await" implicit--although we'll need something for non-awaiting. Let's call it "nowait". Now we can write code like this: func foo(x int) task[int] { return 2 * x; }
func bar(y int) task[int] { return nowait cpu_bound_future_factory(func() task[int] { return y - 5; }); }
func baz(a, b int) task[int] { return foo(a) + bar(b); }
A-a-and we're back to multithreading, basically. So yeah, looks like cheap context switches is the solution.My idea is to make the continuation explicit and CPS transform all and only the functions that have a continuation as parameter (and any generic function with a type that is a continuation or contains a continuation). Fall back to dynamic stacks if the continuation escapes.
It is always continuations all the way down.
So, leaving aside the problem of interacting with the OS for I/O, you can do interleaved CPU-bound in a single thread, by instrumenting your code with basically an instruction counter. When it grows a bit too large, stop, reset it, and switch to another task. Erlang does this, although being a functional language, it counts only function calls (no other way to make a loop than to (tail-)call itself). No CPS needed.
If you compile both for the same goal the system will have suboptimal memory usage or performance, and in general you don't know whether a function will be interrupted in advance if you don't have colored functions (also you can't move the stack if you have references unless you have a GC, which is terrible).
That said, I am not claiming thread-per-core is _the solution_ either, just saying that if you can partition data at application-level, you can make things run plenty fast. Of course, you’re also exposing yourself to other issues, like “hot shards” where some CPUs get disproportionately more work than others. However, as we scale to more and more cores, it seems inevitable that we must partition our systems at some level.
In fact a well designed generic executor can be completely oblivious to whether it is scheduling plain closures, async functions or fibers. See for example boost.asio.
I would go further, and say assuming your program is doing "interesting work" (not just calling a database / reading a file and returning the results), use a threadpool. In almost all applications threads are fine, certainly when you are comparing threaded Rust to people using Python/Ruby.
The problem with custom stacks is you essentially need to have a custom runtime and it makes calling into C functions a lot more difficult (cgo isn’t a cakewalk)
This runs counter to Rust’s C++ inherited belief that you don’t pay for what you don’t use and would have have made Rust less feasible for all kinds of other projects.
async/await might not have the best developer ergonomics, but it does have the best implementation ergonomics from a language point of view
Python's GIL comes to mind :p
Goroutine (loom) style concurrency only makes things worse both ergonomically and performance-wise, pushing programmers towards slow buggy lock-riddled code as the only way to use such models.
Dont use async-await then. Its just syntactic sugar (minus a few details that make it so ergonomic that you can do things that otherwise would be more work than they are worth). Just use plain futures. Be joyful as there will be no "colored" functions in your editor - only 10x more code (and less readable!).
I think there is a place for arguing about async await syntax and whether fibers are the better solution (try implementing them without a large runtime! Rust got rid of theirs.). Clearly, the benefits provided by continuations are worthwhile in avoiding context switches - to argue otherwise is wild. Context switches are very expensive and, in performance critical code, a huge red flag.
To actually answer your question a bit: Threads allow all the benefits of processes (except the two I just mentioned), with the added benefit of being lower overhead and allow sharing information more easily (but you can still use full-blown IPC to communicate between them if you wish).
Threads also introduce security bugs and possible instability (one thread crash brings the whole thing down).
> all the benefits ... (except the two I just mentioned)
One of those two exceptions was indeed memory protection.
Multiple processes works better if you have little or no sharing or communication between the threads.
Indeed without that, you might as well use multiple processes
All of these languages/libraries use a dynamic scheduler for load-balancing:
* Rayon (Rust) [https://github.com/rayon-rs/rayon]
* Goroutines (Go) [https://golangbyexample.com/goroutines-golang/]
* OpenMP [https://www.openmp.org/]
* Task Parallel Library (.NET) [https://docs.microsoft.com/en-us/dotnet/standard/parallel-pr...]
* Thread Building Blocks (C++) [https://software.intel.com/content/www/us/en/develop/tools/t...]
* Cilk (C/C++) [http://cilk.mit.edu/]
* Java Fork-Join and Parallel Streams [https://docs.oracle.com/javase/tutorial/collections/streams/...]
* ParlayLib (C++) [https://github.com/cmuparlay/parlaylib]
But if thread-per-core is fundamentally tied to the idea of sharding, then I think I see what you're saying.
Any sufficiently complicated concurrent program in another language contains an ad hoc informally-specified bug-ridden slow implementation of half of Erlang.
http://rvirding.blogspot.com/2008/01/virdings-first-rule-of-...
If you are interested, look into the new Windows I/O scheduler that was implemented for GHC 9.0 (in Haskell, obv). There is also some preliminary work to integrate io_uring into the Haskell runtime if it' available, but the Haskell runtime philosophy is rather different from the Rust approach.
Eg https://docs.microsoft.com/en-us/windows/win32/fileio/i-o-co...
>The concurrency value of a completion port is specified when it is created with CreateIoCompletionPort via the NumberOfConcurrentThreads parameter. This value limits the number of runnable threads associated with the completion port.
>[...]
>The best overall maximum value to pick for the concurrency value is the number of CPUs on the computer.