Why you can have millions of goroutines but only thousands of Java threads
rcoh.me
rcoh.me
The magic, as the author points out, is dynamic stack allocation. If you disassociate a thread with a required stack size, then you are back to using just the amount of memory that your object heap and reference stores are using. GC gets to be trickier to of course.
This should nominally be a non-issue on 64 bit address machines as you only need map the pages that actually have data in them (minimum one page) but that still puts a huge load on the page tables. This is especially true with huge pages which are used to avoid tlb thrashing and to save on the number of levels you have to go through to get back the physical address of the page you are trying to access.
Its the kind of optimization problem that systems folks really like to dig into and optimize around.
Where else can you allocate stack frames except for in userland?
So yeah, they were pretty new and they were blowing up things left and right. Sun at the time had a leg up on most (if not all) UNIX vendors with a threading model that worked well. But it made porting the JVM harder and green threads lived for a long time as a fallback.
of course it's not something a developer might want to do, but might be done as part of an user-land library. no idea why one would want that today outside of highly specialized code, but it might be that it could be more portable to other kind of non-preemtible system (i.e. os with fixed time sharing)
I would imagine register based VMs have a leg up in this regard?
That doesn't make any sense - register VMs have exactly the same requirement for stack frames as stack based VMs do.
There have been other projects like Quasar from ParallelUniverse [2] that retrofits them by processing the compiled bytecode to enable saving/restoring the IP & stack.
[1] http://cr.openjdk.java.net/~rpressler/loom/Loom-Proposal.htm...
In implementations of languages like Scheme or Haskell where no stack may even be allocated, or where stacks might be allocated in chunks, the cost of a co-routine can be very small -- as small as the cost of a closure. If you take that to the limit and allocate every call frame on the heap, and if you make heavy use of closures/continuations, then you end up with more GC pressure because instead of freeing every frame on return you have to free many of them via the GC.
In terms of ease of programming to deal with async events, the spectrum runs from threads on the one hand to callback hell on the other. Callback hell can be ameliorated by allowing lambdas and closures, but you still end up with indentation and/or paren/brace hell. Co-routines are somewhere in the middle of the spectrum, but closer to threads than to callback hell. Await is something closer to callback hell with nice syntax to make things easier on the programmer.
Ultimately it's all about how to represent state.
In threaded programming program state is largely implicit in the call stack (including all the local variables). In callback hell program state is largely explicit in the form of the context argument that gets passed to the callback functions (or which they close over), with a little bit of state implicit in the extant event registrations.
Threaded programming is easier because the programmer doesn't have to think about how to compactly represent program state, but the cost is higher resource consumption, especially memory. More memory consumption == more cache pressure == more cache misses == slower performance.
Callback hell is much harder on the programmer because it forces the programmer to be much more explicit about program state, but this also allows the programmer to better compress program state, thus using fewer resources, thus allowing the program to deal with many more clients at once and also be faster than threaded programming.
Everything in computer science in the async I/O space in the last 30 years has been about striking the right balance between minimizing program state on the one hand and minimizing programmer pain on the other. Continuations, delineated continuations, await, and so on -- all are about finding that sweet spot.
The point is that thread-per-client is just never going to scale better than C10K no matter what, and that the work we see in this space is all about finding the right balance between the ease-of-programming of thread-per-client and the efficiency of callback hell (typical C10K).
For any size computer (larger than tiny), C10K designs will wring out orders of magnitude better scaling than typical thread-per-client apps no matter how much you optimize the hardware for the latter.
Thread-per-client can only possibly compare in the same ballpark as C10K when the actual stack space used is tiny and comparable to the amount of explicit state kept in the C10K alternative, and even then, the additional kernel resources consumed by those per-client threads will dwarf those needed for the C10K case (thread per-CPU), thus adding cache pressure. In practice, thread-per-client is never that efficient because the whole point of it is that it makes it trivial to employ layered libraries of synchronous APIs.
Now, I'm not saying that the alternatives to thread-per-client are easy, but we should do a lot better at building new libraries so that async I/O can get easier.
Proposal: http://cr.openjdk.java.net/~rpressler/loom/Loom-Proposal.htm...
It's in the early phase, but there is a prototype currently.
The commits in the mailing list mention a recent zero copy continuation thawing if I'm not mistaken.
http://mail.openjdk.java.net/pipermail/loom-dev/2018-October...
EDIT: I'm mistaken, it is about lazy stack walking (but I remember zero copy continuation thawing might be attempted)
We think we can achieve the level of performance we want with stack copying.
Anyways, allocating frames or stacks on the heap doesn't seem like it should require teaching the GC all that much. They would essentially be objects of such a class that the GC can traverse them... just like a normal stack (from current frame through return closures), which it already has to know how to do. Their location shouldn't be so important.
That's wrong. There are a variety of different Apache mpm's: https://httpd.apache.org/docs/2.4/mpm.html
In the footnotes, it also points out that Erlang uses a similar system, which is true, and worth a look as well.
Also in Java fashion, the support isn't perfect, and requires some fiddling with JVM parameters (and annotations, obivously). It's still really solid from my experience however, and a shame that it isn't more widely known.
A partner project Comsat provides fairly comprehensive library support for databases and such. I've used it before with Vert.X and the performance was insane. Possibly faster than Go.
Another aside, there's an active research project to add greaan threads to Javas core as well. Project Loom
Since it's Java, I expect it to be slightly faster than Go while using double the amount of memory.
Could you write more about it? Maybe a blog post?
* Thread Local storage. A million Goroutines means nothing is they are all fighting over shared storage. Consider trying to implement Java's LongAdder class, which does intelligent sharding based on contention.
* ConcurrentHashMap. sync.Map is horrendously contentious for writes. It's not even possible to write one yourself because Go hides the built in hash function it generates for each type.
* Goroutine joining. This is between a rock and a hard place, because if you don't wait till Goroutines are done, you risk leaks, and if you do wait, you need to put WaitGroups all over the place.
All those millions of goroutines don't help if you don't have the tools to coordinate them.
On a 64-bit machine, with 1MB stack size (this is a claim on the virtual address space, not the physical memory), you can have millions of threads too (that you may hit a lower limit due to other OS level knobs is another matter).
[1] for classic contiguous stacks of course. That doesn't apply for separately allocated frames or segmented stacks.
If you write your app with more explicit state rather than implicit state (bound up in the stack), such as by writing a C10K style callback-hell, thread-per-CPU application, you're going to get much much better performance. The reason is that you'll be using less memory per-client, which means fewer cache fills to service any one client, which means less cache pressure, all of which means more performance.
The key point is that thread-per-client applications are just very inefficient in terms of memory use. The reason is that application state gets bound up in large stacks. Whereas if the programmer takes the time to make application state explicit (e.g., as in the context arguments to callbacks, or as in context structures closed over by callback closures) then the program's per-client footprint can shrink greatly.
Writing memory-efficient code is difficult. But it is necessary in order to get decent throughput and scale.
The thread switching costs itself could be prohibitive to performance. IIRC, goroutine switching cost is roughly 1/10th of linux thread switching costs.
[0] http://man7.org/linux/man-pages/man3/pthread_attr_setstacksi...
The problem is the scheduler doesn't know when a goroutine can do useful work. It thinks that if a new value arrives to a channel there is a good chance the goroutine will do something useful but that's not the case in general. My favorite example is Game of Life (CA) implemented as a network of goroutines communicating their state changes through channels: you have to collect 8 updates from predecessors to update your state, meaning you'll be scheduled to run 8 times before you finally make a real progress. Not good for scalability.
It's true that it's hard for a scheduler to "know" that useful work is going to be done, but that's just a fact of life, just like the CPU doesn't "know" whether or not you're going to use the next page of RAM or go flying off in an effectively random direction.
When would millions of goroutines be useful if you can only run ~4 threads in parallel?
Imagine that you want to model a quarry. You have a pile of rocks, a rock crusher, a pile of gravel and a whole bunch of trucks. A truck takes a rock from the rock pile. it gives it to the crusher. The crusher works on the rock for a while and outputs some gravel. A truck takes the gravel to the gravel pile. You want to simulate different ways of scheduling the trucks so as to be optimal.
You could write your simulation so that you update every object in your system at every time stamp, but that's really complicated and it also makes it really difficult to collect statistics about processes. Another way is to create a kind of "process" for each action in the system and a "resource" for each thing that could be used in the system. You then write the code for each process as a continuation and you request resources, blocking when they aren't available. Your overall simulation loop is very simple: it runs a single continuation until it requests a resource, and then it blocks. It then runs the next continuation until it requests a resource. Then it blocks. You just keep doing this, unblocking continuations when the resources become available (and all other continuations' simulation clocks are passed its current clock).
The nice thing about this is that you need no concurrency at all -- or you can have a concurrent thread of execution for each continuation. The latter will obviously run faster, but that's not necessarily the point. The point is that the code was dramatically easier to write.
For some kinds of simulation, a million continuations is not unreasonable at all. In fact, for a lot of applications it would be great if you could have a lot more.
Concurrency.
Have you heard of the 10k problem?
Some people think the best way to handle tends of thousands of concurrent clients is to have each run in a lightweight thread (or goroutine, or whatever term for the same basic abstract concept.) This is as opposed to using an event-based system like Node.js does. Maybe they're right or maybe they're wrong, but that's where it's useful.
Lets say an IO call to high latency connection takes 4 seconds. With no concurrency in between that IO call the program just sits and waits for a response. Under concurrency the program can switch to another thread while it's waiting for the IO call to complete.
Typically multithreading adds way more complexity and a whole new class of bugs to the code.
> So I think I've identified the culprit and I believe the culprit is XFree86.
> it took seven years after this change went in for someone to figure out all the different parts that were damaged and fix them later, right? So threads were broken for a really long time.
https://www.deconstructconf.com/2017/joe-damato-all-programm...
And people were saying it on platforms other than Linux.
Another use-case for threads that I see is when you have work-queues that have different priorities, so you put each queue into its own thread with a matching priority. For example, 1 thread for responding to UI events, and 1 lower priority thread for all background work.
Anyways, none of these use-cases require many threads. What use-case am I missing that makes go-routines useful?
I guess doing something in a thread is a fail-safe way to make sure that 'some time' is spent processing it, so for example for the UI tasks, it would be easy to just dump all handling of UI events in virtual threads.
Processes or co-routines are often used more as concurrent objects. That is, they're used to separate data and concerns into separate, isolated memory areas. However, they're also sequential programs in their own right, so the system becomes much easier to reason about. If you have a few OS thread per core, you are responsible for distributing work among those threads, whereas if you have many co-routines per "job" you can let the runtime decide and schedule intelligently.
A concrete example is a web server. Implemented in Erlang or Go you would typically spawn a co-routine for each incoming connection. There you would get isolation and a simple implementation (parsing the request, performing the task, returning response and exiting). A crash or bug handling a specific request would not affect the other clients.
Other great use cases are messaging and communication apps (examples like Rabbit Q or Discord come to mind), concurrent databases, etc.
Sure, it can be handled as a big array where every element is a struct consisting of the ingress and egress socket/fd referencse, read and write buffers, and some saved state label, and then crunching all of it in a big while loop, in C, with epoll. Nginx does this, but that big array can become the bottleneck quickly. (Lock contention, cache line conflicts resulting in poor memory bandwidth utilization, etc.)
But if you have low-overhead async-await, then you don't have to solve the big array problem and you don't have to write that huge state machine thing either, things become easy[-er] to reason about again.
You can get high concurrency in Java by avoiding the 1 thread per request/connection model and using NIO with easy-to-use libraries like Netty.
The thread-per-client sin derives in large part from having synchronous APIs first.
Every API, every project, needs an async-first approach.
Even then, many many all-synchronous-all-the-time thread-per-client apps will get written, but at least it will be possible to write thread-per-CPU, async-I/O apps where those apps need to scale.
It's much harder to rewrite software to use async I/O than it is to write it to use async I/O in the first place. It's also much harder to write software using async I/O than sync I/O. Therefore it's all about economics and forecasting. If you can forecast needing to scale soon enough, then start with async I/O, otherwise increase revenue and profits first then rewrite. For many companies developer time is a larger item on their budgets than the cost of extra servers and power and all that, which is why we have so much thread-per-client software -- it's a natural outcome.
"each OS thread has its own fixed-size stack."
100% wrong.
From a virtual memory standpoint, every thread would appear to have a fixed size which can be set when you create it. [0]
From a physical memory standpoint, every thread would appear to have a dynamic size (but bounded by the virtual memory size of course).
[0] http://man7.org/linux/man-pages/man3/pthread_create.3.html
FTA: "With 4KB per stack, you can put 2.5 million goroutines in a gigabyte of RAM – a huge improvement over Java’s 1MB per thread"
You can implement Go's thread model (more generally known as green threads, coroutines, fibers, lightweight threads) on the JVM. This is precisely what Project Loom intends to do and what Scala effect systems already do. You construct a program in a type called IO that only describes a computation that can be suspended and resumed, unlike Java Runnables. These IO values are cooperatively scheduled to run on a fixed thread pool. Because these IO values just represent continuations, you can have millions of them cooperatively executing at the same time. And there is no context switch penalty for switching the continuation you are executing for another.
Technically os threads are also simply continuations.
But don't OS threads work like this as well, by being integrated with file descriptors among other things? If a thread is blocked on a read, the OS knows this and won't keep scheduling that thread - right?
If so, the article's argument about time wasted on scheduling doesn't make sense.
You're basically asking if generalized preemptive multitasking that has to solve every problem equally well is really that much slower than specialized cooperative multitasking that is tailored to specific language?
Of fucking course user level scheduling is going to be superior. Just think about it. Languages like erlang and ponylang dispatch to the scheduler after finishing the execution of every single function. Do you really think that switching to a new thread on every function call is going to be faster than not doing that? Consider that the way you're supposed to do iteration is via recursion which means you will have one function call on every iteration and therefore invoke the scheduler on almost every iteration. The user level scheduler can instantly switch to the next pending actor/goroutine/whatever as if it was just a regular function call. Meanwhile your regular threads will have to switch to the kernel first, now has to load the scheduling tree from main memory because the cache are filled with user level data, then has to make a complex scheduling decision and finally switch back and reload whatever data you just flushed out of the cache.
So yes the article's argument about time wasted on schedule makes a whole lot of sense because not only is the scheduling of preemptive threads itself more expensive, it also incurs extra costs through context switches and emptying the cache.
No.
> So yes the article's argument about time wasted on schedule makes a whole lot of sense because ... context switches ...
Right. I understand that doing multi-tasking closer to the application code can be more efficient by avoiding involving the kernel and context switching.
What I was wondering about was this specific argument that the article was making:
Suppose for a minute that in your new system, switching between new threads takes only 100 nanoseconds. Even if all you did was context switch, you could only run about a million threads if you wanted to schedule each thread ten times per second. More importantly, you’d be maxing out your CPU to do so. Supporting truly massive concurrency requires another optimization: Only schedule a thread when you know it can do useful work! If you’re running that many threads, only a handful can be be doing useful work anyway. Go facilitates this by integrating channels and the scheduler. If a goroutine is waiting on a empty channel, the scheduler can see that and it won’t run the Goroutine.
Doesn't the OS also perform this optimization, i.e., only scheduling threads that can do useful work? If you have thousands of threads that are all sleeping or blocking on read, it was my understanding that the OS will not schedule them – contrary to what the article is saying above. Am I wrong about this?
Additionally, as others have pointed out, Java (and, for that matter, all native code on Linux) used to use the Go model and switched to 1:1.
http://www.wirfs-brock.com/allen/things/smalltalk-things/eff...
"The paper contributes Biscuit, a kernel written in Go that implements enough of POSIX (virtual memory, mmap, TCP/IP sockets, a logging file system, poll, etc.) to execute significant applications. Biscuit makes liberal use of Go's HLL features (closures, channels, maps, interfaces, garbage collected heap allocation), which subjectively made programming easier. The most challenging puzzle was handling the possibility of running out of kernel heap memory; Biscuit benefited from the analyzability of Go source to address this challenge."
By my count, Go is a fifth system. (CTSS, Multics, Unix, Plan 9.) It's interesting to observe how the second system effect progresses when the same team has been working on the same problem for 60 years.
1. is Go scheduler better than Linux scheduler if you have a thousand concurrent goroutines or threads?
2. is really creating goroutines that much faster than creating a native thread?
3. are gorutines relevant for long lived and computationally intensive tasks or just for short lived I/O task?
4. What is the performance of using channels between goroutines compared to native threads?
tbh, I have read several of respected articles that criticize Golang goroutines in terms of performance and I am not really sure that Golang's only virtue imho which is simple concurrency is performant at all
package main
func main() {
go println("I ran")
for {}
}
If you run with: GOMAXPROCS=1 go run main.go
It will never print the statement. This doesn't come up frequently in practice, but Linux does not suffer from this case.As a sibling correctly notes, `select {}` is what you want to do instead. You need to do something that can block the goroutine in such a way that control returns to the scheduler. Selecting is one of those ways.
1. Go's scheduler is better at scheduling Go. There are well defined points at which a context switch can occur, and the scheduler is in userland so you never really have expensive context switches.
2. Not sure. I would guess so, if only because it saves a couple context switches.
3. Goroutines are absolutely relevant for both. The stdlib webserver is a high-performance, production-grade implementation and it forks a goroutine for each request, and applications frequently fork more goroutines within a request handler.
4. I assume you're comparing Go channels to pipes? In which case the answer is probably always "faster", but how much faster depends on the size of the thing you're transferring since IPC requires serialization and deserialization while in Go the overhead is just some locking and a pointer copy.
5. Go has dynamically-sized stacks. While you can change stack size on Linux (to use very small stacks and thus spin up more threads), I don't think you can change the size of an individual stack. Plus with Go, you don't have to configure stack size at all; it works out of the box on every platform.
For me, that Go's goroutines are performant is just gravy. I like that they're a platform agnostic concurrency abstraction. I don't have to deal with sync vs async ("What color is your function"; which incidentally cost me the last 1.5 workdays tracking down a performance bug in Python) or posix threads vs Windows threads.
GCC also, on some platforms, allows using segmented stacks on any C and C++ programs which means that a thread will only use only as much physical and virtual memory as required. I don't think segmented stacks have seen much use though as they have non zero cost and, for optimal space usage, ideally require recompiling any dependncy, including glibc with them.
Interestingly segmented stacks were implemented to support GCCGO.
This took some of the bloom off the rose for me. For each goroutine I had to predict ahead-of-time whether it would require a kernel thread, and if so send it through a rate limiter. Effectively it was sync-vs-async but hidden.
An expected way to manage this would be to have a limit on the number of threads to run blocking syscalls. Up to you how to do it (threadpool for syscalls, anythrrad can run a syscalls but checks a mutex first, etc)
I don't think there's a danger of deadlocks here -- your blocking syscalls shouldn't be dependant on other goroutines. Eventually the calls will succeed or timeout (one hopes) and the next call will commence.
In my experience, you can usually find a safe number of parallel syscalls to run -- and it's often not very many.
https://www.ardanlabs.com/blog/2018/08/scheduling-in-go-part...
To be more specific, most socket calls should be async + selectish, but file system calls would likely be done as synchronous, because async doesn't generally work for that -- anyway limiting the number of outstanding file I/O requests tends to help throughput.
https://stackoverflow.com/questions/41627931/in-golang-packa...
https://www.techempower.com/benchmarks/
http://marcio.io/2015/07/handling-1-million-requests-per-min...
It can be possible only if there are no more "physical" threads than cores. And anyway a switching between async routines on the same "physical" thread requires a switching of routines contexts (go to the memory and update cashes).
They are definitely not more performant than native threads with a work-stealing approach (producer - consumer pattern) that is tuned to the number of cores available. How could they be more performant? Even the simplest switching adds some cost, from a performance perspective it would never make sense to run thousands of threads on 8 cores, whether they are scheduled by the OS or by the language implementation.
So to answer your 3rd question: short lived I/O tasks.
I keep expecting someone to write one and put it up on NPM (I haven't looked recently, though). Perhaps I should clean up my and stick it up there. It provides a very nice abstraction of continuations. I did it merely as kata to show that continuations do not require concurrency.
https://github.com/olahol/node-csp
Are you asking to compare mutexes and channels? The article goes through multiple reasons why the go scheduler has more information about how something should execute than the Linux kernel. Reread and then maybe rephrase?
Now you want to scale, so you need to abandon the thread-per-client, everything-is-synchronous-looking model and start doing async I/O the hard way (perhaps with await, perhaps as callback hell). That's a huge refactor / re-write exercise.
But you'd be surprised how many thread-per-client apps exist. It's so easy to write them...
Threads are quite useful in many situations and have great performance characteristics. It comes down to people knowing the language they are using and using it correctly.
Means higher throughput with lower jitter, that in our case was what we were looking at that time. For go defense we were using PooledByteBufAllocator for recycle Output Streams but memory usage was not a concern in benchmark as much as GC does not affect throughput. But I also believe that go http stdlib memory usage is better than netty, sorry for not being more helpful with this topic.
We have both java8 and go, high throughput critical microservices with excellent 99pctl, and in general memory is not a concern (as much as we do fine tuning on gc and don't have any memory leak) and generally for the really critical and portable solutions we choose java over go (unit testing and library versatility is a big player in this discussion)