Uber Engineering Tricks of the Trade: Tuning JVM Memory for Large-Scale Services
eng.uber.com
eng.uber.com
Yeah if Uber was born today and didn't want to use any on-prem resources they'd use S3, but HDFS is the best alternative to object storage out there, and there's a huge ecosystem of tools around it. If you're running your own datacenter, there's not any serious alternative.
It's a little bit sad that they did not hire an expert.
And i'm far from being an expert in GCs, but the basics are just wrong.
GCs before ~2010 have been optimized for throughput more than for latency. Since then, you can choose.
Since 10 years, G1 let you set your maximum pause time and you get a huge warning if it miss that target, ZGC or Shenandoah are from the beginning latency first, throughput second.
Like your handle, too.
It is also quite ironic that they paid for Zing C4, and then came up with "P90 average RPC queue time" as the metric, ignoring Gil Tene's wisdom that you should never average percentiles.
Another problematic statement is this:
> which in turn decreases request latency and increases throughput
Gil Tene wants you to measure p99.99 or pmax, so that you will pay for Zing C4 to get lower tail latency, while sacrificing some throughput. If you end up getting both lower tail latency and more throughput, then something is already very wrong with the existing setup (i.e. the full GC).
I am surprised how much legacy inefficient crap is lingering around in companies like Uber.
Edit: instant downvotes. Okay, S3/GCS/Azure is the typical answer (egress costs be damned)
- if on-prem is a must there are multiple options, generally something with erasure codes (it is a game changer for storage)
So far I have been using enterprise storage (that has some potential problems when mounted as nfs volumes), works for petabytes, already decouples storage from compute.
More recently I was experimenting with MinIO. No conclusion so far.
The problems are with Hadoop:
- unfortunate design choices (namenode??)
- extremely unfortunate implementation (I probably spent more time in the Hadoop codebase than any other, found many bugs, some I could fix, most I couldn’t)
I think I have migrated away from Hadoop 10 PB worth of data infra in the last 5 years, mostly to AWS, some to Azure. Average cost saving is between 10-30% yoy.
Some comments point out the network cost. The reality is most companies collect a giant amount of data (ingress) and publish dashboards (egress). It makes cloud pretty viable.
S3 is beating the shit out of HDFS in reliability and cost, even though most Hadoop shops spread the fud that it is slow. Same way these companies used to spread the fud that snappy is best for data compression.
As of 2021 even the latest adopters (banks and insurance companies) use cloud. Maybe extremely few dogmatic companies remain in the onprem crowd. Even those will eventually give up.
Per the article Uber has hundreds of petabytes of data.
Well in fairness, have you ever seen the S3 codebase? I mean honestly it could be a fork of HDFS for all we know.
S3 has a really good architecture and a great implementation.
HDFS has a meh architecture with a bad implementation.
There were obvious signs. I remember when Twitter decided to investigate why HDFS was slow and they figured out some details about how Hadoop guys decided to implement their own dictionary for configuration that had a much worse time complexity than the default dictionary in Java. There might be a video about this somewhere.
And there are more things like that. I used to have 5-10 years old HDFS Jira tickets open. I just gave up.
Here is a video:
https://www.youtube.com/watch?v=jupArYWxoq0
Hadoop is full of these things.
One more thing:
https://lamport.azurewebsites.net/tla/formal-methods-amazon....
I would love to see similar approach to Hadoop.
I've seen distributed file systems on S3 - can it also be done with Azure Blob storage?
https://doc.dataiku.com/dss/latest/hadoop/hadoop-fs-connecti...
https://devblogs.microsoft.com/cse/2016/05/22/access-azure-b...
AWS is incredibly expensive in comparison. Uber is /not/ a small company technology-wise, for better or for worse.
Maybe there are some JVM specific libraries they need for things like mapping that don't exist on GoLang? From reading their tech blog, it sounds like they're mostly a GoLang shop and they were apparently Python before that. So it seems like they're probably forced into using Java because some of the libraries they use aren't worth rewriting in-house to avoid having to use the JVM.
Do you consider monetary cost as part of efficiency? The scale that these companies operate at make it worthwhile.
People constantly underestimate the fact that >80% of the cost of software is maintenance. Upgrades and patches and redesigns/refactors and following modern conventions costs 4x more than the cost it took to build it. But nobody puts that in their budget. And then later people find this huge pile of legacy shit, and think, wow, this company really screwed up. But in fact this is totally normal.
By the way, this is the case because of how we develop software and hardware, not because they are some inscrutable element that can't possibly be built once and maintained for a lifetime (like, say, a building). We just don't build it to last.
With that, STW pause time does not depend on the heap size or on the root set size (the stack size). In practice, it means pause < 1ms, at that point the OS becomes the bottleneck, not the GC.
So latency is good but throughput can be reduced by 30%.
That might be what they meant.
Also the go GC is behind the JVM GCs if you want to really tune it IMHO