Operating a large distributed system in a reliable way: practices I learned
blog.pragmaticengineer.com
blog.pragmaticengineer.com
In other words, when person is in ER, doctors are looking for heartbeat, temperature, ... , not for some low level metric, like how many grams of oxygen is consumed by some specific cell.
A well designed distributed system is going to be able handle a certain failure rate in its components because requests will be retried automatically. If a component's failure rate increases from 0.01% to 0.1%, there will probably not be any user-visible impact... but if you can detect that increase, you might be able to correct the underlying issue before that component's failure rate increases to 1% or 10% or 100% -- at which point no amount of retrying will avoid problems.
And yes, distributed systems are hard, because it is inherently hard to reason about what will happen when something changes/fails.
I'm not from Uber but I recently worked (was responsible for a huge chunk of infra) at a company that is bigger and has more products running and I saw some hilarious failures. And what is different, it was rarely a bad code push, and when it was, due to the nature of the business, sometimes it was really hard to roll back.
I'd phrase your advice differently as well: monitor on contact points. Monitor where rubber meets road. Monitor where two components interesect, be it between teams working on frontend and backend components, your application and it's database, or low level metrics - the system API. I've found many interesting problems and averted several potential crises because I saw that the metrics things like requests looked different between the client and server, the system metrics and my applications, etc. Contact points are what's important, and it just happens that every (Figuratively) application has a contact point with the OS.
To put it unscientifically, something always goes wrong, but that does not mean SLA is violated. As a corrollary, you will be alerted constantly about something that might go wrong because of some signal that someone thought might lead to some outage.
At face value it seems proactive, but it's really not. Better spend time actively making the system more reliable (e.g. look at your or someone else's postmortems, or do premortem exercises).
Without those "irrelevant" low level metrics for a large growing distributed system, you're not only flying blind but crashing increasingly more and more often.
Crashing more often is caught by an alert that looks at total capacity is s service. Doesn't matter if it's random bit flips or OOMing nodes at first. Longer term these metrics can be useful to increase efficiency again (should i first fix random bit flips or OOMing tasks). I would not base SLA relevant alerts on too low level alerts.
Systems that only monitor the highest levels appear to function fantastically well right up to the moment they crash spectacularly. Then the forensics will show you that at lower levels there were plenty of warning signs telling you that the system was headed for the cliffs and a large number of those signs will be apparent while there was still time to do something about them. Reliable systems engineering is not something you can do just at the highest levels because of the build in resilience in intermediary levels. This is counter-intuitive but born out in countless examples of systems that look robust but aren't versus those systems that really are robust.
Priority is always the customer. The focus can change given the circumstances and the context of each issue. Sometimes the focus must be high-level, sometime it requires a low-level focus, but the priority is always the customer.
While you're monitoring for traffic/errors/latency throw in minimum success rate. Make a good estimate of how many successful operations monitored systems will do per minute on the slowest hour of the year and put in an alert if the throughput drops below that. You'd be surprised how many 0 errors/0 successes faults happen in a complex system, and a minimum throughput alarm will catch them.
It makes communication more efficient and enjoyable when all parties are aware of this.
Basically a high level guide through modern architectures, frameworks, and database designs. So far, my takeaway has been learning what tool would be useful for certain types of data engineering, not the details of how to write code with it.
Edit - link: https://news.ycombinator.com/item?id=20417801
The best source for distributed computing I found is Facebook's engineering blog.
- https://github.com/mxssl/sre-interview-prep-guide
- https://github.com/theanalyst/awesome-distributed-systems
One linked resource for example:
1) Maybe not everyone belongs to the 1% that can pass the interviews?
2) Maybe not everyone lives in a place that has a Google office where development is happening?
Sorry but you sound terribly arrogant here, pretending that it's just an option that everyone has and that the whole world lives in the same place where you live.
On normal complicated control systems, such as in an airplane, you do have plenty of time to complete your task in the given timeframe, but in F1 there's not much time left, and you have to optimize everything, HW, protocols and SW. Much harder than in gaming or at scale at Google/Facebook/Amazon.
Eg you have to convert your database rows from write-optimized to read-optimized on demand. A simple fast database cannot do both fast enough. A normal firewire protocol is not fast enough, you have split it up into trees. A normal CAN protocol ditto, you have to multiplex it. You have to load logic into the sensors to compress the data and filter it out, to help transmitting the huge amount of data.
[1] http://shop.oreilly.com/product/0636920025986.do- Active monitoring
- Chaos testing
- Cold start testing
Clearly never worked for a hospital. Hospitals need good engineers (and often don’t have them). Our ‘nines’ are embarrassing...
> The Five Whys, as it’s commonly presented, will lead us to believe that not only is just one condition sufficient, but that condition is a canonical one, to the exclusion of all others.
Five whys presents itself as a way to dig deep but promotes doing so linearly and getting to a singular thing you can fix, hiding a lot of potential learnings along the way. Thinking about contributing factors is a much more powerful framework.
Really great content, but was really taken back by “I” used everywhere. Maybe it’s a new thing that I am not hip on that I ought to try - “I built and ran transaction processing software for Bloomberg! This is what I learned!”
But perhaps you really did all that by yourself, in that case sorry that i doubted you, looks like it’s a lot.
If you've designed software where the whole service can degrade based on the CPU consumption of a single machine, that right there is your problem and no amount of alerting can help you.
Unless it's your database.
For example: if you have a fleet of hosts handling jobs with retries, a bad job could end up being passed host to host killing each host / locking up each one as it gets passed along. And that could happen in minutes while replacing and deploying and bootstrapping a new host takes longer. So by the time your automated system detects, removes, and spins up a new host everything is on fire.
I stand by my beef with this article. The statement that "I've talked with engineers at Google [and concluded that a thing Google wouldn't tolerate is a must-have]" doesn't make sense. What I get from this article is you can talk with engineers at Google without learning anything.
Things like CPU hopefully shouldn't be your key/gold service-up metric, but paradoxically, the more mature your system the more CPU can tell you; you can catch problems before they happen. It can help notice things like bad CPUs.
Memory stays pretty important in my experience; even more than CPU.
And in addition to all the other responses there are also different levels of pages: Some are page me at 5am, some can wait till morning, and some can wait till Monday. FAANG is more likely to have their own hardware so you actually get deeper/more diverse monitoring needs than a shop on AWS or something.
Source: FAANG-ish tier infra work
I wouldn’t alert on a single machine having CPU issues, but I’m definitely interested in a small collection of individual machines all having CPU issues at the same time.
Software deliveries/releases can often realistically be non-perfect. (Don't have direct experience with Canary releases TBH though)
In case anything goes wrong any objective evidence which helps to reconstruct the failure scenario is valuable.
Also... Murphy's Law.
>If you've designed software where the whole service can degrade based on the CPU consumption of a single machine.
Typically, if such software is indeed released, I think it will be several CPUs on several hosts.