A single line of code made a 24-core server slower than a laptop (2021)
pkolaczk.github.io
pkolaczk.github.io
Asynchronous code, coroutines, async/await, parallelising problems is my deep interest and I blog about it everyday.
I think the easiest way to parallelise is to shard your data per thread and treat your multithreaded (or multimachine) architecture as a tree - not a graph - where dataflow doesn't need to pass between tree branches. This is similar to the Rust's "no interior mutabiliy" and Rust data structures pattern.
My machine can lock and unlock 61570760 times a second. But it can count to 2 billion in 1 second. So locks are expensive.
I recently worked at parallelising the A* graph search algorithm that I'm using for code generation/program synthesis.
For 16 processes it takes 35 seconds to synthesise a program but with 3 processes it takes 21 seconds. I think my approach to parallelising A* needs a redesign.
We hit Amdahl's law when it comes to parallelising. I need to split up my problem into spaces that don't require synchronization/serialisation.
EDIT: I've mentioned this whitepaper before ("Scalability! But at what COST?") but this whitepaper would be useful reading of anybody working on multithreaded or distributed systems. In summary: single threaded programs can easily be faster and more performant (wall clock time) than multithreaded/multimachine distributed machines, but they don't scale.
https://www.usenix.org/system/files/conference/hotos15/hotos...
Theoretical advancements matter too but usually only in so much as they can translate to hardware. Although, there are some special case where even a slight theoretical gain matters even more than how it translates to hardware, but they're limited.
Anyways, to that end, everything becomes about dividing work in a way that parallelizes nicely, exploits cache well, reduces the need to share information between threads, etc... and this ultimately comes down to data structures. However, these aren't your normal, fundamental data structures. Instead each problem sort of has some kind of exotic, hypothetically ideal data structure, that is fine tuned to exploit the machine's resources to the max for just that problem. By the time you're done they rarely resemble anything intelligible, let alone wha the whiteboard version of the algorithm was.
In that vein, there are a few general trends that appear over and over again. Trees vs graphs is definitely one of of those trends, although that's more of a general theme and not a literal rule.
For reference, we spent the whole class doing SSSP on a graph representing every road in the United States. The naive A*/Dijkstra implementation took like 45 min, and naive bellman-ford never finished (probably days/weeks). By the end of the class I was so proud to have parallelized Bellman-Ford to the point where it took like 30s or something. And my TA's record breaking implementation was <1s. All of this was on a normal university linux desktop.
It's kind of like programming = logic + data. However, there's a continuum between the two, and on the far, extreme end of that spectrum, the structured access of the data can be the program itself.
Sometimes I wonder whether we should not focus more on these micro adjustments. We usually have this “it’s fast enough” attitude, but probably everything we use today (including web services) could be instant if focus was given to optimization (although yes, I understand the drawbacks of focusing solely on that).
Other much simpler optimizations are usually possible in "corporate code". Lots of stupid things being done. To use your python script example, proba ly you could've gotten 2950x improvement rewriting the bad parts better but still in python :)
Do you remember what the python script did/was for? How much time was spent building the original script? What was the guy that built it paid?
How long did it take to build the C++ version and how much was that guy paid?
If these were all open source/unpaid, what would it cost for these different types of people and could the company conceivably pay those salaries/keep the guy happy and busy enough to stay around?
Not that I think that should be priority #1 but it’s staggering how much more efficient things can be than the “naive” approach.
I don't think an example of one and the process taken to define it is that helpful, since the skill is about being able to do this for an arbitrary problem. It also probably requires a lot of trial and error.
When I first got a dual CPU (before "cores" were a thing) computer, I decided I'd try dipping my feet in some "proper" parallel coding.
I started with something simple, parallelizing a quicksort routine. This seemed quite trivial: instead of recursing, add the spans to be sorted to a list. I then spawned a thread per CPU, which fetched a span from the list, did a single quicksort pass on it and added up to two new spans to the list. Rinse repeat until list was empty.
Since each span was non-overlapping, the threads only had to synchronize while accessing the list of tasks.
When benchmarking it became clear that while there was a good performance win for large arrays, for short arrays the multithreaded code was much slower than the non-multithreaded version. At my hardware the threshold was around 50k items for integer elements and 20k or so for string elements, IIRC.
I added a threshold detection, where the thread would do a regular recursive quicksort on the span if the length was below the threshold, and this yielded significantly better results.
And with that the harsh reality of multithreading hit me: no free lunch. It was clear the threshold varies not just with element type (slow/complex comparators would reduce the threshold, and vice versa) but with the details of the hardware. So a hardcoded threshold was out of the picture, and it would have to be dynamically determined at runtime.
Was a great learning experience though.
Thus the majority of spans are small and quickly processed, and hence the naive multithreaded version leads to lock contention as the threads fight to access the shared list of spans.
My neighbour scanning is dynamic and my neighbours and all neighbours from a node is independent from that point forwards, it shall not visit the exact same node. In essence my problem is kind of a tree.
I am trying to infer data flow between two states including hidden states such as functions calls. My dream is that I can provide a start state and end state and the computer writes itself based on type information and data flow analysis of values.
Here's my input data - which is what memory is set to and what registers are set to.
start_state = {
"memory": [0, 0, 0, 0],
"rax": 0,
"rbx": 1,
"rcx": 2,
"rdx": 3,
"rsp": -1,
"rdi": -1,
"rbp": -1
}
end_state = {
"memory": [3, 1, 2, -1],
"rax": 3,
"rbx": 2,
"rcx": 1,
"rdx": 0,
"rsp": 6,
"rdi": -1,
"rbp": -1
}
# these functions take in a value and return another value
minus_1_to_four = Function("minus1", -1, 4)
four_to_five = Function("fourtofive", 4, 5)
five_to_six = Function("fivetosix", 5, 6)
This synthesises the following program in 16 seconds (I improved the heuristic function). Function values input and output can be in any register, but in my example they are all in the same register. With 3 processes it synthesises in 0.6 seconds. [start, mov %rax, (%rdx), mov %rbx, (%rbx), mov %rcx, (%rcx), mov %rdx, (%rsp), mov %rax, %rbp, mov %rdx, %rax, mov %rbp, %rdx, mov %rcx, %rbp, mov %rbx, %rcx, mov %rbp, %rbx, call minus1(rsp=-1) -> rsp=4, mov $-1, %rbp, call fourtofive(rsp=4) -> rsp=5, call fivetosix(rsp=5) -> rsp=6]A* is naturally recursive [1] and so can be parallelized as you go further.
[1] unless I'm mistaking some other algorithm for A*, each new recursive call of all the neighbours can be started on a different thread or node.
If you can restrict the structure of your graphs (e.g. planar) then some very efficient methods exist.
So a tree is just a subtype of a graph?
But yeah. Some graphs are trees. And you can construct trees within graphs for efficiently navigating connected graphs, which is done in various important and famous algorithms.
This is my code generator/program synthesiser, it synthesises assembly instructions from two given states of memory and registers to take it from one to the other, including hidden states such as functions.
This speeds it up from 16 seconds to 320 milliseconds. I am thinking how to make the algorithm scale, by creating better and alternative program candidates in each thread.
This is really for a program I'm writing in Nim, so perhaps it depends on how channels are implemented?
Riscv has an interesting compromise, which is to delineate a subset of ll/sc loops which is guaranteed to eventually make global progress. I do agree that it is better to include real wait-free primitives like cas and faa; but I wish that such guarantees of global progress would be provided to HTM.
Which is true, but eliminating sharing works even better if possible as proved by the article.
The fine article shows that a single lock xadd can destroy perfs on some x86 systems and explain that it is due to cache line bouncing. You would get the same effect with loads: if the loaded data is RO or mostly RO it will of course scale fine. It won't scale as soon as it starts bouncing too much.
Also, FWIW, intel largely does implement operations closer to that way on single socket parts, if you want to see it for real look at on-device atomics on a GPU. Ironically an average laptop chip handles atomics much faster than most servers as a result.
Interesting--can you link a reference for this?
The short version is that if atomics are implemented as part of the memory network, common cache, or memory controller, then atomics of the form “fetch-and-X” can be implemented in roughly equivalent complexity to a load of the current value (plus an instruction for the op, give it take) with the cost only scaling past that as op queues or other implementation-specific limits fire. It’s the infinite consensus ops that just can’t scale no matter what you do. The coherence and memory model matter a lot too of course, which is part of why x86 tends to be slow for atomics, while arm and ppc (with fetch-and-X extensions) or GPUs tend to do much better.
Thank you, Sherlock.
But for the rest of us: when you need shared state, lockfree atomic spinlocks are roughly 1000 times more performant than mutexes. (Not a scientific estimate, numbers taken from real-word experience.)
The slow part of locking is invalidation of cache lines, and this has to happen with spinlocks anyway. Modern mutex implementations also first try to acquire the lock optimistically, so in the uncontended case they are as fast as userspace spinlocks (modulo inlining).
And if you have a contended lock, then userspace spinlocks are a PITA. You need to take care of fairness, ideally deal with the scheduler (yield to a thread that is not spinning on the same spinlock), and so on.
You can do all of that properly, but even then, you're looking at maaaaybe 10-20% performance increase in real-world applications.
Pure spinlocks can win only in contrived cases, like only having exactly two threads contending for the lock, with short locked sections.
Don't you anyway need to drop the scheduler a hint so that the thread holding the spinlock doesn't get scheduled off the CPU, making the contenders wait longer than they ideally should? (Or is this what you meant by your "fairness" reference?)
In my limited understanding, this was the no. 1 reason why userspace spinlocks were discouraged -- because pretty much no scheduler accepted a hint from userspace to not kick a thread off the CPU -- modulo jumping through hoops with priority, et cetera.
If I'm missing something (and I likely am), I would be glad to be educated.
How would you do it? You can change the thread's priority to realtime to prevent the scheduler from pre-empting it while holding the lock, but this requires a kernel roundtrip and several scheduler locks anyway.
You can have a worker thread pool, with individual threads hard-pinned to specific CPUs. Then you can dispatch your work into these threads. This in practice will guarantee that they are not pre-emptied except for occasional kernel housekeeping needs.
But this will make it impossible to use the kernel-level mutexes because they can block your worker threads. So you'll have to reimplement waiting mutexes in userspace, along with a scheduler to intelligently switch to a work item that is not blocked on waiting for something else to complete.
Long story short, you're eventually going to reimplement the kernel in userspace. This can be done, and you can get some performance improvements out of it because you can avoid relatively slow kernel-userspace transitions. DPDK is a good example of this, but at that point you're not just using spinlocks, you're writing software for essentially a custom operating system with its own IO, locking, memory management, etc.
(Edit:) Arguably, I could've been less obtuse in what I wrote.
That is absolutely not universal. See e.g. but there are of course many places discussing this: https://news.ycombinator.com/item?id=21970050
Spinlocks are terrible, and written by people who are trying to do quick hacks because they work in terrible environments and are taught to do bad things.
There is a reason that most GPU drivers are just lists of hacks to get games working correctly.
Uncontested (single thread):
incrementing using atomics took 0.002011 s (0.2513 ns / increment)
incrementing using mutex took 0.005515 s (0.6894 ns / increment)
Contested (8 threads trying to increment a single protected integer): incrementing using atomics took 0.1069 s (13.36 ns / increment)
incrementing using mutex took 1.970 s (246.3 ns / increment)
So mutexes are roughly the same speed in the uncontested case, and about 20x slower in this heavily contested case.
This is on Windows.Unfortunately, the userspace part of pthreads is not tuned for performance and does a lot of ridiculous things if you care about parallelism.
(Mostly I was complaining about the poor quality of userspace system libraries.)
Every feature of programming languages started in this fashion.
We don't have good optimising compiler for that
The memory allocator plays a significant role, since allocation strategy needs to be per-thread-/per-CPU-cache-aware. Choosing and then tuning a different malloc (e.g. tcmalloc, jemalloc) to the one in your platform's default library is a non-trivial matter but may have enormous impact both on overall performance and memory demand.
In addition, when you design computation this way it is relatively easy to hadoopify it later, since it's basically map-reduce writ small.
i feel like the post i was responding to was talking about handling/pinning a thread to a specific CPU core?
IMO It is not enough to know the logical constructs used for synchronization in parallel programs, you have to know the hardware too.
A little bit of everything from high-level parallel algorithms/data structures through memory consistency models, compiler optimizations to processor micro-architecture (cache coherency protocols, atomic instructions, NoC overhead etc.) is needed. Basically we need to be aware of the overheads for contention at every level in the system.
While fooling around I had the idea to exploit that not all atomic operations are equal. So I added an additional "contention" flag. When a thread wanted to aggregate it would do an atomic read of the flag, if it was set it would bail[1] and continue to accumulate to the local buffer. Once done aggregating the flag would be reset after unlocking.
Effectively this was adding a single-iteration spinlock before the "heavy" lock, but even using CriticalSections for the lock on Windows (which does spin before acquiring a mutex lock) it resulted in clear improvements, especially when running on machines with more than 8 cores.
[1]: It would bail unless the local buffer grew too large, so slight memory vs perf tradeoff there.
So what does this mean in practice? In my view, the way to think about it is that atomic writes have non-local side effects. But since atomics are necessary for synchronization, and involves both reads and writes, we should compartmentalize and minimize synchronization as much as possible, to avoid these gnarly issues creeping up and tanking real world performance.
Arc<T> (and it’s relatives in other languages) constitute textbook violations of this rule. In Rust they are everywhere in non-trivial code, including in the async runtimes themselves. Of course, they also violate (or evade if you’re generous) ownership principles of idiomatic Rust, (or “hello world-Rust”, if you will). I think we need to take a hard look as an industry at ref counting as a silver bullet escape hatch to shared data.
This isn't expensive if cache lines are uncontended, though.
> I guess you could say they contend with unrelated data but that would stretch the definition a bit.
I think you might be talking about "false sharing." This is real contention on the cache line due to co-location of apparently unrelated variables.
> Arc<T> (and it’s relatives in other languages) constitute textbook violations of this rule.
Definitely!
> In Rust they are everywhere in non-trivial code
Ehh.. only the hot ones matter. Most are not actually contended much, and the article's solution (unshared clone) is a very reasonable approach to scale these without an API change.
You’re right. And cache lines are quite small, so this is probably less common. Yet, it’s another potential source of perf regressions in concurrent code, as if it wasn’t incredibly complex already.
> Ehh.. only the hot ones matter.
Well.. first atomics have even more non-local effects, such as barriers on instruction reordering. So Arcs that are cloned willy nilly can still be significant, with no contention.
But let’s ignore that and focus on the contended case: when you hear “uncontended X are basically free” it (subjectively, imo) downplays the issue, like contention is some special case that you can compartmentalize and only worry about when you consciously decide to write contended code. The blog post demonstrates exactly how this is so easy for contention to creep in, that you have to be superhuman levels of vigilant and paranoid to spot these issues upfront. Extremely easy to miss in eg code review.
I think both compile- and runtime tooling could help at least partly here. I’d also give rust some credit for having explicit clone instead of hiding it.
The latter is a bit clunky but the core more or less implements them in the same way. Acquire a line exclusive, load value, increment it, write it back. And you can hold the line exclusive such that the conditional store failure cause is mostly a formality, and can't actually become an infinite loop.
No general purpose atomics are done by shipping the operation to the cache or to memory controllers, it just doesn't work[*]. So even if they look slightly different in the core, they all end up looking exactly the same at the caches and coherency protocols, and that is where atomics are slow. Well any sharing of cache lines updates really.
[*] EDIT: That is to say it doesn't work for performance, for many reasons. Some CPUs do have "remote atomics" something like that which does exactly this, but they are not intended to be broadly used.
Only if your software is badly implemented. If you follow the requirements specified by the architecture, forward progress is guaranteed. Of course there is no guarantee how long it will take, but the things that make it slow are essentially the same things that make atomics slow.
Of course not sharing at all is of course best, but often you have no choice in that.
This is significant because if you profile your code and find that a mutex is expensive, you should change your algorithm to avoid contention rather than blindly trying to change the code to use atomics instead of mutexes.
I found a contention bug inside of Wine a few weeks ago. Something that is supposed to be "lockless" really had three nested spinlocks. With many threads contending for a lock, performance would drop to about 1% of normal.[1]
Can I get the HN comment length explanation of what this means?
I thought the title sounded familiar, and the culprit is more or less the same (false or in this case unnecessary sharing). But I didn’t think it was quite that long ago, so maybe it’s two articles about the same classic blunder.
(Before reading the paper I was expecting that the additional bytes were put for the split counter, plus thread id - but it actually packs them using lower bits for reference counting).
I wonder what abseil/folly/tbb do - need to check (we are heavy std::shared_ptr users, but I don't think 14 bits as described in the paper above would be enough for our use case)
The most performant solutions are RCU (https://github.com/facebook/folly/blob/main/folly/synchroniz...) and hazard pointers (https://github.com/facebook/folly/blob/main/folly/synchroniz...), but they're not quite as easy to use as a shared_ptr [1].
Then there is simil-shared_ptr implemented with thread-local counters (https://github.com/facebook/folly/blob/main/folly/experiment...).
If you absolutely need a std::shared_ptr (which can be the case if you're working with pre-existing interfaces) there is CoreCachedSharedPtr (https://github.com/facebook/folly/blob/main/folly/concurrenc...), which uses an aliasing trick to transparently maintain per-core reference counts, and scales linearly, but it works only when acquiring the shared_ptr, any subsequent copies of that would still cause contention if passed around in threads.
[1] Google has a proposal to make a smart pointer based on RCU/hazptr, but I'm not a fan of it because generally RCU/hazptr guards need to be released in the same thread that acquired them, and hiding them in a freely movable object looks like a recipe for disaster to me, especially if paired with coroutines https://www.open-std.org/jtc1/sc22/wg21/docs/papers/2020/p05...
folly::SharedMutex, the main RW lock in folly, tries to shard them by core when it detects contention (in fact it is the OG core-sharded primitive in folly) and when that works it is virtually linearly scalable, but the detection is a heuristic (which also has to minimize memory and writer cost) so there are still access patterns that can be pathological.
That's very interesting. I dealt with the scenario where I had to scale the hash-map to support 1k-5k concurrent readers and a few sporadic writers. I ended up sharding the hash-map with each one using the RW lock to guard the access to the underlying data. This essentially declined the contention within the RW lock itself, or at least I wasn't able to measure it.
Instead, if a mutex showed up as hot in a profile, I'd look at things like RCU/hazard pointers for read-biased data structures in most situations, or trying to shard or otherwise split data between cores such that there isn't much contention on the boring, vanilla mutex.
They don't have to, see my sibling comment about folly::SharedMutex.
> (Ok, there are sharded reader count variants that reduce the cost of the reader count, but they're still only useful for relatively long reader lock sections.)
RW locks are a code/design smell, even with a cheap reader count.
What you're saying is wrong: if the reader section is long, you have no problem amortizing the cache invalidation. Sharded counts are useful when the reader section is small, and the cache miss becomes dominant.
Also I don't get this "RW locks are code smell" dogma, not having RW mutexes forces you to design the shared state in a way that readers can assume immutability of the portion of the state they acquire, which usually means heavily pointer-based data structures with terrible cache locality for readers. That is, in order to solve a non-problem, you sacrifice the thing that really matters, that is read performance.
I've heard this thing from Googlers, who didn't have a default RW mutex for a while, then figured out that they could add support for shared sections for free and suddenly RW mutexes are great.
[0] https://sites.cs.ucsb.edu/~ckrintz/racelab/gc/papers/levanon...
However, HybridRc would still be as contended in the scenario in the blog post, wouldn't it (before the fix that solved it)? Just checking my understanding.
> In Rust, it is very easy to generate flamegraphs with `cargo flamegraph`.
... Also in pretty much every other language that can generate perf stacktraces, because this is just a wrapper around Brendan Gregg's FlameGraph visualizer: https://github.com/brendangregg/FlameGraph
During execution you had two kinds of memory locations, some in CPU caches and some in RAM. By running all the threads on one socket, everything accessed from the cache was just a fast cache access. Everything accessed from the memory was a slower memory load. Frequently loaded/stored locations will tend to go to the cache.
In the NUMA setup, you would have a larger cache (more than one socket) which would mean that more locations were likely to be in the cache. However, if a core on a socket tries to access a location which is on another socket's cache, it will use the interconnect between them to access it.
If you have an unfortunate memory layout, this can make it so that you end up having a large percentage of the accesses using the interconnect (slower than cache access) and values get swapped between the caches constantly, which forces subsequent accesses to also use the interconnect.
Another way to avoid this except using just one socket is for the designer of a program to consider NUMA nodes as separate processing units and design around that. Both should be processing separate data and they should only share small amounts of data for synchronization/communication. Then the caches will be much less affected.
With the way that processors are going, with this focus on increased core counts etc, NUMA is increasingly being important to understand and account for, as processors are getting more "NUMA-ish" (to borrow a co-workers apt description). Especially Neoverse/Arm CPUs, etc.
Lockless/lockfree refers to the fact that there are no deadlocks.
Lock-freedom is a much stronger statement. It guarantees progress of at least on thread in finite time, regardless of failures.
Source: Dan Alistarh, Keren Censor-Hillel, and Nir Shavit. "Are lock-free concurrent algorithms practically wait-free?" December 2013
The implied expectation underlying "beginner-friendly" seems naively misguided; it's an advanced undergrad computer/software engineering topic in the most permissive sense, and the blog's prose appears to have been tailored with that minimum target audience in mind.
Example: https://ieeexplore.ieee.org/document/7804711
Edit: looks like there's only 12 cores per CPU so that's 24 physical cores. 48 HT cores. So the drop must be cache trashing?
For throughput tasks it’s often the case that you go with less parallelism to reduce Amdahl’s law a few percent, and instead investing in keeping the pipeline saturated, so that the variance in concurrent tasks is lower. Work stealing being one of the more notable tricks.
apparently 2 process per core is more efficient than 4.
It’s very complicated and I don’t blame anyone in particular. It’s an ok solution, but it’s definitely not in the zero-cost category. I’d rather accept that this is the state of things and working towards systemic solutions in allowing borrowing in more situations, but that requires a compile-time verifiable hierarchical thread- and task model.
Fun fact: Arc was partially the reason Rust disallowed borrowing across threads. Arc had already become popular and combined with thread borrowing someone demonstrated use-after-free. Borrowing across threads seemed less important at the time, so Arc was kept. This led to the famous “leaks are safe” rule, which was sold as a mere clarification of an inherent truth.
It’s possible that all that was inevitable and correct, but I was never convinced by the arguments made in those old threads or the subsequent writings about it. To me, it looked like details were glossed over in order to get to a swift resolution. I’m quite content waiting for someone else to figure out whether Rust could have gone down a different path. Worst case, I’ll come down peacefully from this tiny hill. Best case, there’s an alternate timeline where Rust could live out its full potential in concurrent environments.
(Started last month after an update. They are still trying to figure it out. I instead not use ctrl-c to copy files as that doesnt need the little menu to load. I kick myself every time i forget and accidentally right-click.)
Try using ShellExView to disable all 3rd party shell extensions. There's a good chance that something like Acrobat is running some code that doesn't play nice with your windows version (or hardware etc).