Why events are a bad idea for high-concurrency servers (2003) [pdf]
people.eecs.berkeley.edu
people.eecs.berkeley.edu
The test setup used for this paper was a "2x2000 MHz Xeon SMP with 1 GB of RAM running Linux 2.4.20". Solid PC server iron for 2003, but basically equivalent to a $5/month server from Digital Ocean today.
If you're looking to squeeze 100,000 concurrent tasks from that $5 server, this paper is relevant to you.
It's really quite simple: threads encourage the use of implicit program state via function calls, with attendant state expansion (stack allocation), whereas evented I/O encourages explicit program state, which means the programmer can make it as small as possible.
Smaller server program state == more concurrent clients for any given amount of memory. Evented I/O wins on this score.
But it gets better too! Smaller program state == less memory, L1/2/3 cache, and TLB pressure, which means the server can take care of each client in less time than an equivalent thread-per-client server.
So evented I/O also wins in terms of performance.
Can you write high-performance thread-per-client code? Probably, but mostly by allocating small stacks and making program state explicit just as in evented I/O, so then you might as well have done that. Indeed, async/await is a mechanism for getting thread-per-client-like sequential programming with less overhead: "context switching" becomes as cheap as a function call, while thread-per-client's context switches can never be that cheap.
The only real questions are:
- async/await, or hand-coded CPS callback hell?
- for non-essential services, do you start with
thread-per-client because it's simpler?
The answer to the first question is utterly dependent on the language ecosystem you choose. The answer to the second should be context-dependent: if you can use async/await, then always use async/await, and if not, it depends on how good you are at hand-coded CPS callback hell, and how well you can predict the future demand for the service in question.On the other hand, simplicity in a code base usually matter. Code written with an evented API, littered of callbacks, is usually harder to read and maintain than that written in a sequential way with a blocking I/O API.
You can recreate a sync API on top of an evented architecture using async/await, but then you have the same performance characteristics of a blocking API, but with all the evented complexity lurking underneath and leaking here and there. Seems to me a very convoluted way to arrive to the point from where we started.
A GB of ECC server RAM costs more. An extra GB of RAM in the cloud can even cost you $10/mo if you have to switch to a beefier instance type.
I don’t have the answer, but you would want to measure dollars to buy it, and nanoseconds to refill it.
One lesson I've learned is: a) make a library from the get-go, b) make it async/evented from the get-go. This will save you a lot of trouble down the line.
Let's say you configure it for 2000 max connections (really not much) so that's 2000 threads, so 20 GB of memory right away because the thread stack is 10 MB on Linux. It's a lot of memory and it's obliterating all caches.
You can reduce the thread stack to 1 MB (might crash if you go lower) but any caching is still trashed to death.
Next challenge. How do you think concurrency work on the OS with 2000 threads? Short answer is not great.
The software making heavy use of shared memory segments, semaphores, atomics and other synchronization features. That code is genuinely complex, worse than callbacks. Then you're having issues because these primitives are not actually efficient when contended by thousands of threads, they might even be buggy.
* On Unix, only socket IO is really evented; this can be solved by using a thread pool (the Flash paper describing this problem and solution for httpds dates to 1999[1]) but is inelegant. This is the approach Golang takes behind the scenes.
* How to scale event loop work across cores; this can be solved, but adds complexity. In contrast, threads are quite simple to scale.
* How to share accept() workload across cores; this can be solved, too, but not portably. I.e. your 100 core server may very well be accept-limited if you have a single-core accept loop, or highly contend on a single socket object.
* Threadpools are definitely inappropriate if you have a ton of relatively idle clients (long-poll server, or even typical webserver), but are less wasteful if you have relatively few, always-busy clients.
I don't disagree that the synchronous threadworker design isn't the best choice for high performance services that can afford a lot of engineering time to design and find all the bugs. But thread-per-client is often a completely acceptable place to start.
[1]: https://www.usenix.org/events/usenix99/full_papers/pai/pai.p...
Threads haven't exactly stood still in that time either, especially if you include green threads, fibers, coroutines, etc.
> If you're looking to squeeze 100,000 concurrent tasks from that $5 server, this paper is relevant to you.
It's relevant regardless, as part of a long-running back and forth between threads and events. For example, Eric Brewer was one of the co-authors of this paper, but also for the SEDA paper which was seminal in promoting event-based programming. I highly doubt that we've seen the last round of this, as technology on all sides continues to evolve, and context is a good thing to know. Those who do not learn the lessons of history...
If it's a web server then obviously you want more cores, and an event driven server would make a lot of sense.
Basically if you need concurrency then you want events, if you're compute bound than you don't want that overhead.
EDIT: Instead of i7 just imagine any high end (high frequency and IPC) consumer chip
2014: https://news.ycombinator.com/item?id=7684163
2011: https://news.ycombinator.com/item?id=2907415
Smaller but interesting threads from 2010-12:
https://news.ycombinator.com/item?id=3482002
https://news.ycombinator.com/item?id=3101451
https://news.ycombinator.com/item?id=2910849
https://news.ycombinator.com/item?id=1547353
I worked on this subject during my Phd, and the result is a paper that will be published in sigmetrics2020.
We developed an M:N user-level threading library and exhaustively tested it against event-based alternatives and pthread based solutions.
We used both memcached and webservers to test it on 32 core and 64 core servers.
Even connection/pthread looks promising in terms of performance.
You can find the paper and source files here: https://cs.uwaterloo.ca/~mkarsten/papers/sigmetrics2020.html
[1] https://github.com/scylladb/seastar/wiki/Memcached-Benchmark
Regarding memcached, the first reference you posted is from 2015, yes in 2015 memcached was in a very bad shape in terms of synchronisation and things have significantly changed from that version with locks per hash bucket rather than a global lock, avoiding try_lock and .... So those results are too old to rely on. Seastar moves network stack to user-level and if I remember correctly the scheduler consisted of multiple ever looping threads even if there was no work to do. Considering in our experiments 40-60% of the memcached overhead were coming from network I/O, there is no surprise about their results. I would call this a hack for sure.
I have not read the Anna paper to be able to comment, but it seems that they are creating a key-value store using the actor model. I briefly schemed through the paper, and I did not find anything that points to "killing any possibility for shared memory multi-threading as a concurrency model"; this is a very bold claim and if they do claim that, they should have really strong results.
But my guess is anna's whole actor model was implemented on top of threads and shared memory multi-threading? Which will be in contrast of that bold claim. I worked with actor models and implemented one as well, it is a perfect fit for many use cases but threads at least for now are the bread and butter of multicore programming.
Having said that, all these models have their respective place in the software world. What we are trying to show in our paper , through thorough experiments, is that the misconception that event-driven has higher performance than thread programming is not fundamentally true. Therefore, falling into asynchronous programming and create hard to maintain applications [1] only due to better performance has no merit.
[1] https://cacm.acm.org/magazines/2017/4/215032-attack-of-the-k...
> NGINX scales very well to support hundreds of thousands of connections per worker process. Each new connection creates another file descriptor and consumes a small amount of additional memory in the worker process. There is very little additional overhead per connection. NGINX processes can remain pinned to CPUs. Context switches are relatively infrequent and occur when there is no work to be done.
> In the blocking, connection‑per‑process approach, each connection requires a large amount of additional resources and overhead, and context switches (swapping from one process to another) are very frequent.
This is their explanation on events vs threading approach. Still, a lot of web servers today use a thread-per-connection which is acceptable since a database(e.g. postgres) performance degrades slowly as more active connections are introduced.[1]
I guess I’m somewhat used to the asynchronous I/O model that things like Win16 and JavaScript use. You have to think a bit more about the unpredictable order that things happen in, but not about preemption race conditions and data corruption.
The only threading model I have seen that I care for is Erlang’s. The C ... Java model is a data corruption death trap.
At that, Erlang kind of straddles the gap between independent execution and communicating sequential processes.
Meh.
Total gold, nothing in this world is clean or definitive.
https://news.ycombinator.com/item?id=22174201
People only now start to understand that the problem with multi-core is memory speed and the only way to solve that is by using a virtual machine with a concurrent capable memory model = Java.
In-lined memory does not play well with multiple threads. Cache-misses is the way to parallelize code!
* J. K. Ousterhout. "Why Threads Are A Bad Idea (for mostpurposes)". (http://web.stanford.edu/~ouster/cgi-bin/papers/threads.pdf)
* V. S. Pai, P. Druschel, and W. Zwaenepoel. "Flash: An Efficientand Portable Web Server". (https://www.usenix.org/legacy/events/usenix99/full_papers/pa...)
* M. Welsh, D. E. Culler, and E. A. Brewer. "SEDA: An architecturefor well-conditioned, scalable Internet services". (http://www.sosp.org/2001/papers/welsh.pdf)
Those are some of the more influential papers in the development of event-driven designs.
Thread-per-connection is just plain dumb. Threads aren't free. Context-switching is a thing. Modern advice is to try not have more threads than cores. Thread-local storage is a very useful, having 1000's of thread may make TLS unfeasible.
And mention 2003 :)