Notes on Distributed Systems for Young Bloods
somethingsimilar.com
somethingsimilar.com
1. The metrics comment could not be more true. You cannot think hard enough about the actual problems you will encounter and what metrics you will have. Have a common way to pull information about each service, and to analyze it in some common place.
2. Standardize, standardize, standardize. When someone is answering a page at 2 AM and has to deal with a component of your system that they don't really know, the more it resembles other components, the better. Think hard about how you can get every service to be written in such a way that the same critical information is available. Where is it documented? Who do I page? Where is the monitoring? Yes, random third party components won't follow your standards, and will be hard to integrate. But standardization is a good thing.
3. You need to be able to send canary requests through that trace through your whole infrastructure. You should be able to flag a front end request, and have every single request that it generates through your entire system be logged somewhere so that you can see a breakdown of what happened. Hide this so that nobody can use it to take your site down, but build the capacity for yourself. (Standardization will make such a system much easier to build!)
4. Randomly canary a small fraction of your traffic. There is tremendous value in having a random sample of traced traffic. When you're trying to understand how things work, there is nothing like taking an actual request going through a complex system, and seeing what it did. Furthermore if you've got intermittent problems for a small fraction of users, being able to look at a random slow request really, really, really helps you track down issues that otherwise would be virtually impossible to replicate.
I've apologized to him before, but after having worked on Minefold for the last 2 years, I feel like I need to apologize again. The work that Jeff and the others at Twitter have done has been amazing. I was also a massive cock.
Also, that was perhaps the biggest apologetic act I have ever witnessed...good job.
I refer to that as the Night of a Thousand Australians. Good times.
Once you have a Javascript application with state speaking over one or more APIs to your backend services, you're in the domain of distributed system design. Especially if you use the application cache and support offline operation. I frequently see people who are used to more traditional web applications underestimate the system design challenges this causes.
(Those challenges are totally worth solving, because the model is very powerful.)
I wrote up some notes on this, though they're about a year old now and the browser UIs continue to evolve: http://tech-blog.clericare.com/2012/02/choosing-right-browse...
For someone building a consumer web app that needs to painlessly convert users in volume, I can see where it would still be a pain. For me, it's not so bad -- my app is B2B, and having an install step as simple as clicking "Allow this app to use up to 500MB on your computer?" is actually a huge improvement over the "enterprise" junk we compete with.
https://github.com/daleharvey/pouchdb
Since Couch has a master-master replication protocol the database in the browser can just sync with the one on the server. Then the network gets disconnected but that's alright. When it gets re-connected continuous updates pick up where they left off, conflicts are resolved and it sort of works a distributed system. Where a browser is just another node.
Shameless plug It's validating that we mentioned over half of these points in an interview with startup founder Eric Lindvall on his Papertrailapp.com log management app: https://peepcode.com/products/up-lindvall
I'm working on a mobile app that does a key exchange with a server before allowing a server-based registration or login. It's nowhere near as complex as your average distributed system.
That said, I've run into a scary amount of the things mentioned in this article in my tiny little use case. Just trying to ensure a decent user experience (timing out a comms check after two seconds rather than waiting up to 60 seconds when the phone switches from networked to disconnected) in an async message exchange needs some crazy orchestration. Keeping the code clean means refactoring stuff I thought I had nailed two months ago.
I've been programming for a long time, but this stuff humbles me. And happily, I love it.
It isn't just load testing but more that the whole system should be considered suspect. If you don't act defensively at all steps you will be hosed by something you thought will never happen. Just had a good talk about this last weekend. Memory, TCP, and all other rock-solid things can and will have issues in large systems.
This is perhaps a point one cannot get to theoretically just sort of thinking about. This realization comes after observing the effects in practice. It is interesting to observe distributed systems, especially the ones that rely on asynchronous messages, thing go wrong and queues start filling up. Or even weirder machines start to synchronize for some reason. Like the size of the queues will oscillate in resonance of some sorts.
One way out is to reduce asynchronicity but that comes with serious performance penalties.
If I synchronously hit some service then I need to build facilities into the service to say that it's too heavily loaded, and perhaps handle those responses on the client.
If I pop some job onto a message queue or asynchronous endpoint, things will just slow down for a while. Which is normally as good a way as any as soaking up load.
> If I synchronously hit some service then I need to build facilities into the service to say that it's too heavily loaded, and perhaps handle those responses on the client.
Most of the time, you have to have timeouts when there is a synchronous API across the network.
In a large system ideally you'd want to reflect the level of loading back to the input source so the source can slow down sending the data. Think of TCP, tcp works this way on a small scale. One way to fix the problem is to leverage that to actually open a TCP stream and send data that way. The sender will slow down accordingly.
Now if you know that your load is not constant, so there are periods of high activity when inputs are generated then you can try and absorb (and amortize) those high volume peaks using asynchronous queues.
- A website using a load balancer to offset service to multiple machines
- A hadoop cluster running a series of mapreduce tasks
- A botnet setup to DDOS some target
My Desktop has like 40 chrome processes, is running GeoServer on Tomcat (currently idle, but whatever), and the entire Unity desktop environment. I'm currently at 2.7GB. What kind of server OS needs 4-5GB?
0. naming things
1. cache invalidation
2. off-by-one errors
[1] [http://techblog.netflix.com/2012/07/chaos-monkey-released-in...]
Brilliant, made me smile :)
God, I wish I'd known this sooner.
Intellectual honesty at its best...