It's "locking" if it's blocking
yosefk.com
yosefk.com
In many modern processor architectures (x86, for example), a cache-coherence protocol is used to ensure that cache lines provide data according to the memory discipline. On x86, that discipline is Total Store Ordering. (See http://en.wikipedia.org/wiki/Memory_ordering).
This means that if two processors are contending for the same value (eg. trying to increment it, set it to true, or even read it), they will force the CPU to invalidate the cache line containing the value on all the other cores, leading to massive scalability bottlenecks. If multiple cores are contending for the same location in memory, whether it's lock free or not, performance will suffer.
More deviously is the case of false sharing, where 2 different values just happen to fall on the same cache line. Even though they don't conflict, the processor must still invalidate the line on every core. Modern compilers do their best to prevent this, but sometimes they need a little help.
The takeaway is this: Don't try to implement your own locks using CAS -- even something as simple as a lock is very hard to get right (performance-wise) when scaling to dozens and hundreds of cores / threads. People have solved this problem (people.csail.mit.edu/mareko/spaa09-scalablerwlocks.pdf). Writing fast concurrent code (especially lock-free code) is a minefield of weird architecture gotchas. Watch your step.
It's possible to implement algorithms independently of data order on the memory bus, by dependence on execution order, and can therefore be barrier free, but cache lines are synced. Complete lines can be dedicated to single loads, a case where performance increases by sacrificing cache effective size.
The different Sparc memory write systems were always useful for explaining the performance effects of cache coherency and memory barriers.
Either of message-passing concurrency and data immutability would trivialize the problems discussed here.
EDIT: I also fail to see how this is academic atomic instructions are needed to implement any basic concurrency construct.
But it's a feature of small systems that you THINK you have a flat address space and not a message passing system. CPU cores send messages to L1 caches to fetch memory, and the messages go outward to L2 and RAM, all in hardware. With a modern OS kernel this message will often be intercepted and fulfilled with disk IO.
So if you pretend you have a flat address space you won't be miserable or anything, because the entire system (hardware + software) is designed to make it look that way.
Once you get interested in performance details the flat memory space abstraction goes out the window. L1, L2, RAM, and disk latencies are really message passing latencies for messages passed between physically separate devices. Designing your algorithm to minimize the number of messages passed between devices (a.k.a. the number of "cache misses") will improve its performance. Optimizations meant to minimize message passing can have unintended consequences. For example, Linux once allocated memory on the region of RAM closest to the core running the process which requested the memory -- I'm sure it made some benchmarks run faster, but it turned into an absolute disaster for MySQL, because it meant that MySQL would only use one region of RAM. (Google "MySQL NUMA" for the full story. You can probably imagine most of it if you just think about the consequences of mixing on-die memory controllers with multi-socket systems.)
The degree to which the flat memory abstraction applies is related to the size of the system. An single-core, embedded microcontroller with 256 bytes of RAM really does have a flat address space and you can often count processor cycles just by looking at the assembly, the abstraction is basically perfect. A modern desktop with a quad-core processor is going to act less like an ideal shared memory machine because the message passing overhead of using contested cache lines can affect performance. And a world-class super computer might be split into N units of M cards of K chips with L cores on each chip -- the shared memory abstraction will only extend so far before it's all explicit message passing in software.
So the real question is not "how do you plan to write a queue without mutable shared state?" The real question is, "How do you implement mutable shared state using messages?"
Footnote: I think it's very telling... the shared memory abstraction is so convenient, the work behind the scenes is so good, we almost want to argue and say "Look, it's shared state, and everybody knows that it's shared state." Buddy, it's one fine illusion.
I have struggled to find the right ways to think about the phenomena underlying the shared memory abstractions. The limited visibility makes this harder than understanding other aspects of software. You can run instrumented simulations like cachegrind; you can query arcane CPU performance counters (Intel's Nehalem optimization guide was eye-opening regarding how much goes on inside vs. what you might learn in an undergrad CPU design class) -- but at the end of the day you have to be guided by experimentation and measurement. And so many times those leave us with no good theories to explain the observed behavior.
(War story, feel free to skip: One time we were trying to speed up a datafile parser by any means possible -- which was already split into a producer-consumer thread pair, one thread running the unfortunately complex grammar and producing a stream of mutation commands, the other thread consuming that and building up in-memory state. The engineer working on this found that adding NOPs could speed this up, and he measured and charted a range of # of NOPs and chose the best one. Our best guess was "something to do with memory layout?" The outcome of the story was that we tore out a bunch of abstraction layers and ended up with a simpler, single-threaded parser that didn't need such heroic and bizarre efforts, but it also left us feeling a bit of vertigo with regard to memory hierarchy behavior.)
Your pointing out that the shared memory abstraction is backed by message-passing between hardware components (which each represent concurrent processes) is really interesting - thanks!
We are writing such code because it is an established way of doing things that also happens to be more convenient to a lot of programmers in the field.
The big disadvantage of locks is that performance decreases with contention, and performance of a shared-memory system in general degrades as the number of nodes increases. So every supercomputer in recent history uses a hierarchical approach: a network of multi-core units. Shared memory and locks for sharing data with cores in the same unit, message passing for sharing data with other units.
Just imagine trying to use system-wide locks on the IBM Sequoia. It has something like a million cores.
Also, advances have been made in manycore shared memory systems. The Cray XE6 (the hardware behind, e.g. HECToR [1]) has a hardware accelerated global address space with remote direct memory access that allows PGAS [2] to outperform MPI [3].
By the way, system wide locks are a red herring. At these scales, you avoid global data as much as you can, regardless of what your programming model is.
[1] http://www.hector.ac.uk/ [2] http://en.wikipedia.org/wiki/Partitioned_global_address_spac... [3] http://upc.lbl.gov/publications/pmbs11.pdf
The comment about global locks was intended to be silly, because comparing locks to message passing without talking about what you're doing with them is also silly.
The linked paper compares MPI message passing to an alternative hardware-accelerated message passing, which is interesting, but the choice of micro-benchmarks is not very exciting. To be clear, while the GP was really comparing the actor model (private memory + message passing) against the shared memory + locks model, I was only responding to the parent comment, and when I think "message passing" I don't automatically think "private memory".
You can have, say, a 100,000 x 100,000 matrix represented as an array over thousands of processors, where each processor can read and write each array element individually.
Because not all concurrency problems can be solved with message-passing and immutable data?
Not to mention that you also need to implement said mechanisms in the first place, and that needs to take into account what the post describes.
Oh, and it's not "academic" at all. It's a rather casual post. Nothing academic about it.
Example 1: Algorithms that operate on large matrices. Immutability means (1) that you can't have destructive updates of small parts, but need to copy the matrix in its entirety, and (2) that you can't have in-place operations, increasing your memory usage (and thus decreasing L2 cache performance).
Example 2: Caching expensive operations. You want to share these among as many threads as possible to maximize hit rates, but you also need to avoid serialization bottlenecks. This is incompatible with both pure message passing and, say, the standard implementations for immutable maps.