Distributed Systems Reading List
dancres.github.io
dancres.github.io
The actual hard stuff is not even these papers, it's the implementations that are way more complex than some algorithm or architectural pattern. Anyone who says "X is better than Y" is fooling themselves because it's only the implementation context that matters.
The only thing you can say for certain is that reducing the amount of components and complexity in the system often results in better outcomes.
No, there are a few other things that you can say for certain:
Watch out for positive-only feedback loops, you absolutely need negative feedback as well - or only. Eg. exponential back-off.
Sometimes, you just need a decentralized solution, rather than a distributed one, and you don't have to have the same answer at every scale (eg. distributed intra-datacenter, decentralized inter-datacenter, or vice-versa).
Loose coupling is your friend.
Sure, add an extra layer of indirection, but you probably need to pay more attention to cache invalidation than you think.
Throughput probably matters more than latency.
Reducing the size/number of writes will probably help more than trying to speed them up.
Multi-tenancy is a PITA for systems in general, and distributed ones are no exception (aside: there is probably a huge business for multi-tenancy-as-a-service, if anyone manages to solve it in a general-purpose way), but a series of per-customer single-tenant deployments may be worse, especially if they are all on different versions of the code. Here be dragons.
Don't overthink it. Start with a naive implementation and go from there (see loose coupling above).
> Watch out for positive-only feedback loops, you absolutely need negative feedback as well - or only. Eg. exponential back-off.
Agreed it may need negative feedback, but I'm not sure about always.
If your service has a latency SLA, exponential back-off might kill your SLA (depending on wordage and where the back-off is). The fix is to soft reject requests (RST rather than dropping packets) when you can't meet the demand. This change may allow you to meet your SLA if it's written to prioritize low latency over service unavailability.
This is it's own negative feedback loop, but change from sending RSTs to silently dropping and you no longer have the feedback.
> Sometimes, you just need a decentralized solution, rather than a distributed one
Agreed
> Loose coupling is your friend.
Until it isn't? :)
> add an extra layer of indirection, but you probably need to pay more attention to cache invalidation
Fixes for additional layers tend to increase system complexity compared to fixes for fewer layers.
> Throughput probably matters more than latency.
Until it doesn't :)
> Reducing the size/number of writes will probably help more than trying to speed them up.
Depending on 20 different things... You really have to account for all the system's limits (and business use cases) and find the solution that matches the implementation needs.
> there is probably a huge business for multi-tenancy-as-a-service
Sure, it's called EKS :-) Just build more clusters... Don't worry, we'll bill you...
> Don't overthink it
Yes and no; Yes, in that there will always be unknowns. But no, in that often improvements in communication will provide better solutions without extra work. Think smarter, not harder!
for the uninitiated,
The Dunning–Kruger effect is a cognitive bias in which people with low ability at a task overestimate their ability.
Congratulations, you are better than 95% of the people that I've interviewed out there saying they are experienced building distributed systems. Including system/solution architects.
The conditions where the monolith must be converted to a scalable system should be defined as early as possible.
Start off with a monolith if it makes most sense (in terms of simplicity / MVP) but from the start keep future scalability in mind and plan for it.
I like this.
It’s important to keep mind that SOA is about scaling teams first, code second and not really about throughout per se. A share-nothing web tier plus a couple judiciously applied databases and background job queues can effectively scale a huge proportion of applications without the overhead of a full SOA.
In today's cloud deployments where cost is not a mega concern (as compared to having physical servers in a DC), this approach is very much feasible.
You could run two systems in parallel (monolith v/s scalable) and ditch the monolith once you have the necessary cutover logistics in place.
Of course, the sooner you do this, the better and it is non trivial to figure out the 'when'
It's indeed pretty rare to see services that crumble under the exponential increase in traffic. Especially if it's not a spike from publicity (I don't know how you would say "being slashdotted" these days). It happens, but it's rare. But because it's usually newsworthy or at least interesting it gets more attention and it will feel much more of a danger than it actually is. Both for software developers and entrepreneurs (or even for the general public) it feels lame. Ha! They should have expected this! But the truth is that most of the time you should not prepare for this, because it's pretty rare, it's not even necessary for success, not even for startup scale success at the beginning.
It's pretty easy to see actually: (almost) all of today's successful services provided by startups started as monoliths (or maybe more realistically some kind of SOA, because monolith vs. microservices is really a false dichotomy). It did work a decade ago with weaker hardware, why wouldn't it work now?
It's more important to be aware of overload, and provide feedback to users (or potential users) and some form of load shedding. And test that all. For a lot of applications, being able to quickly plop together a second monolith could work to address spikes. Or switching to a bigger machine: Epyc 2-socket systems get you up to 128 cores (256 SMT threads), and I think 8TB of ram. You can do an awful lot with that.
This is really important wisdom being shared.
And with a single monolith you can comfortably handle several thousand concurrent users, not just several hundred.
From my humble experience I've found that it is also relatively easier to migrate an established monolith to a semi-distributed system, than it is to build and scale a distributed system from the ground up whilst at the same time trying to figure out all its kinks before it too becomes an established system.
That depends on what kind of service you offer, of course.
When the time comes, cleave these chunks off and wrap them in RPC/REST/graphql.
It's good form even for monoliths since it makes for an easy interface to test.
I guess the trick when building new monoliths is to identify these fault lines up front - as opposed to hopefully finding them later.
Will certainly guide my development going forward, thanks.
So basically some monoliths playing together.
I've always been admiring that design and how much it could do and handle with so little resources.
Most times for a whole editorial newsroom a single 48c/128gb machine can be more than enough.
They were particularly smart on avoiding relational databases and sticking to an object database (Versant oodbms, another mostly unknown marvel of software engineering).
That software, even with its shortcomings, was really a great piece of software engineering.
If you can get away with a monolith, it's the way to go, as long as you leave room to quickly grow horizontally as well, if needed.
They do paper reviews and I pick the papers/topics which seem the most interesting (rather than doing them in order).
They are great if you've already gone over the basics using "Designing Data Intensive Applications" or "Distributed Systems for Fun and Profit".
https://www.youtube.com/playlist?list=PLrw6a1wE39_tb2fErI4-W...
Dr. Kleppman also provides notes for that specific lecture [0]; they contain all the slides with the corresponding detailed text, which is really awesome!
[0] https://www.cl.cam.ac.uk/teaching/2021/ConcDisSys/dist-sys-n...
Alex Petrov’s Database Internals: A Deep Dive Into How Distributed Data Systems Work (2019) is another essential recent reference that should be here. Not as broad as Kleppmann but dives a lot deeper into certain topics.
Looks like it was last updated in 2013: https://github.com/dancres/Pages/commits/master (there's a commit from 2018, but it doesn't touch the actual page)
EDIT: (Ah, I see you are a TLA+ promoter, that's why you made a comment like that)
I wouldn't particularly say I'm a TLA+ promoter (it's a FOSS project), any more than anyone who has a great fascination with a language/framework/algorithm/viewpoint is a promoter. We're all promoters of the memes that live inside our heads!
Real world systems need to handle all the things you mention. TLA+ helps you consider all these issues with spelling them out individually. There is no point in building a complex system if you haven't taken the time to validate the correctness of the target system in the first place.
I don't think this is correct.
What TLA+ allows us to do is be more creative in our design and choice of algorithms, while allowing the computer to help us reason about whether the choices we're making still result in a system that is correct. "Correct" in in this context means two things: "safe" as in it doesn't lose or corrupt data, and "live" as in it eventually makes progress without deadlock or other blockers. That doesn't capture "meets the SLA" or "fast enough for real use" or even "tolerates gray failures". All of those are critical properties indeed - but unless you have fundamental safety and liveness you're never going to get those properties anyway. You might think you have them, but then you'll have a bad time eventually.
So TLA+ (and similar tools) aren't a complete solution to the problem, but they are an exceptionally useful one. Fundamentally, they're useful because distributed and concurrent protocols, even very simple ones like 2PC, are wickedly difficult to reason about clearly. Computers can help us reason, and specification languages can help us communicate clearly about our reasoning.
You want correctness first, and performance second. But these two are very much intertwined. And knowing exactly where the boundary is will help you co-design them.
Some AWS engineers have said the following[1]:
"TLA+ [...] giving us enough understanding and confidence to make aggressive performance optimizations without sacrificing correctness."
[1] https://blog.acolyer.org/2014/11/24/use-of-formal-methods-at...
Having used TLA+ for years, I would say it's the exact opposite. All bugs happen due to someone believing something is simple enough to work out in their head while it actually isn't. So if your judgment about what you can keep in your head is good, you never have any bugs and you really don't need TLA+. But if you do happen to have bugs occasionally, then your belief about how much you can keep in your head is sometimes wrong. TLA+ is a very quick way to write down what's in your head so you can think about it more rigorously. Surely, if you can truly keep it in your head, it should be easy for you to write it down precisely. And just in case you're wrong, there are tools that can check if you're right, just to be extra sure. In short, TLA+ helps if you ever have bugs. It doesn't help if you never do.
It's a fantastic visual walk through of the Raft consensus protocol (used in many modern distributed systems like Consul).
Not saying that Distributed Systems is easy, doing research in Distributed Systems is very challenging, but I don’t think reading and understanding these materials should be that difficult for an average software engineer.
I'm taking graph theory next semester and wondering what else might be useful. I took a distributed systems course already and we used no math at all.
You don't have to read DDIA front to back. Just picking a topic (for instance "Distributed Transactions") is enough to get you started building an intuition about these issues.