Uber made an early strategic decision to invest in on-premise infrastructure due to fears that either Amazon or Google would enter the on-demand market as competitors and bring their cloud infrastructure to bear and potentially squeeze us for costs. Azure wasn’t much of an option during this time. This decision limited our adoption of cloud native solutions like SpannerDB and DynamoDB. We ended up doing a lot of sharded MySQL in our own data centers instead.
This on-prem decision led to a lot challenges internally where we would adopt OSS and then have difficulty scaling it to our needs. For some tech like Kafka it worked out, and we hired Kafka contributors who helped us scale it. For other tech like Cassandra it was a pretty epic failure. I am sure more of these war stories exist that I wasn’t privy to myself.
Coupled with the fact that we were early adopters into Golang which had its own OSS ecosystem, we found that writing a lot of our own infrastructure solutions was the only viable option at our scale.
What you are seeing now is a lot of that home grown infrastructure being open sourced in big way as people who have left Uber continue to see value in investing in the tech that they worked so hard to build. There is probably a nontrivial amount of work to scale the Uber OSS down for smaller use cases but some startups are emerging to make that happen.
Source: I worked at Uber from 2015-2019 on product and platform teams and had several close colleagues in infra.
[0]: https://netflixtechblog.com/scaling-time-series-data-storage... [1]: https://www.datastax.com/resources/video/cassandra-netflix-a...
When people are openly and stupidly incentivized like this, expect those people to behave in a predictable way. People started building new services to get promotions instead of “toiling” at supporting their fellow engineers.
It affected most of engineering but especially in teams like Cassandra, where you needed guidance and support to properly use it effectively, it was a disaster. There should have been open office hours to help people with questions and to ensure that teams were using it properly but there wasn’t. Instead people were left to do what they wanted with no structure or guidance and Cassandra was completely misused. Productions problems ensued, people left the team because they didn’t want to be oncall fixing fires all the time, and eventually it came to the point where they decided to stop supporting it altogether. It was a complete disaster caused by very poor engineering management.
We all knew that Netflix and Facebook use it without issues, but because of stupid management, it failed at Uber.
Ok, but I am fairly confident Netflix also is at that kind of scale.
Netflix has a section on Atlas's documentation about how they get around this: https://github.com/Netflix/atlas/wiki/Overview#cost
They also did this nice video that outlines their entire operation including how they do rollups: https://www.youtube.com/watch?v=4RG2DUK03_0
This is how they do the rollup but keep their tails accurate to parts per million and the middle to be parts per hundred: https://github.com/tdunning/t-digest
A few of my thoughts on this, and this has come up before. Firstly Netflix self-identifies it is expensive to run an in-memory TSDB for metrics - for instance Roy's talk on Atlas mentions this as such[0] at the 37min mark of his Operations Engineering talk "It scales kind of efficiently. I'd love to say efficiently instead of efficiently-ish however that's hard to claim when my platform until this last quarter cost Netflix more than any other element of the cloud ecosystem ... Atlas and the associated telemetry costs Netflix 100s of thousands of dollars a week". At Uber M3 cost a significant amount to run as well at first and that is why M3DB was born to drive down that cost as much as it could and still provide a ton of instrumentation to engineers. Either way, giving engineers tons of room to instrument their code will result in a high cost no matter what since it will be viewed as a free lunch, that is why squeezing the economics on this matters since you want to provide as much instrumentation as possible at the lowest cost.
Regarding your points about their documentation on cost:
1) Yes reducing cardinality by dropping node dimension on metrics, etc is possible to save cost - but also keeping things on disk is an alternate and great way to save cost too and keep the data at high fidelity. The challenge is making on disk lookup fast too, which with M3DB is what we were focused on doing.
2) Dropping replication of the data to a single replica is another way to save cost, however also comes with operational complexity as now you need to do backup/restore if you lose data and lose the ability to query that data in the meantime. This is why M3DB always is recommended (as per documentation) to run at RF=3 with quorum reads and writes so losing a single machine does not impact the availability of your operational monitoring and alerting platform.
3) Regarding rollups and tail solutions accurate, we always push for people to use histograms as that can be aggregated over any arbitrary time window and across time series. T-Digests are much more expensive to store raw and aggregate later. Bjorn talked about histograms, their use in Prometheus at FOSDEM[1] and why they're more desirable than t-digests or other similar aggregations.
[0]: https://www.infoq.com/presentations/netflix-monitoring-syste... (video, quote is at 37minutes in)
[1]: https://fosdem.org/2020/schedule/event/histograms/ (slides and videos)
Maybe Bjorn's talk has the answers, but would you mind explaining how histograms are easy to aggregate? Don't you need either fixed buckets or raw data to produce a new histogram over a different dataset? (I know there are tricks to get great estimates, but naturally every re-aggregation would add larger and larger +/- intervals, no?)
Think of OpenTSDB and Prometheus. Or for a better comparison think of Thanos https://thanos.io/
As to whether they could fulfil Uber's needs, the thing about scale (real massive scale - I work at Cloudflare) is that everything breaks in weird ways according to your specific uses of a technology. The things listed above work for companies, until they don't. There's few things that seem to truly work at every scale, Kafka and ClickHouse come to mind for wholly different use cases than a time series database.
150m-200m events/minute and about 20-30 trillion (10^12) events stored. Doubling about every 12-18 months or so.
While it's true that things start to creak at scale, this has worked remarkably well for us so far. I doubt M3DB is somehow magical in this regard.
For us the cost savings vs OpenTSDB (millions of dollars of hardware), the faster query time and the reduction in oncall overhead was worthwhile.
ClickHouse works fine as a TSDB if you don't mind getting a little dirty
[1] https://github.com/VictoriaMetrics/VictoriaMetrics/blob/mast...
[2] https://medium.com/@valyala/how-victoriametrics-makes-instan...
Not sure I want that in a TSDB!
"Reducing disk space usage by deleting unneded time series. This doesn't work as expected, since the deleted time series occupy disk space until the next merge operation, which can never occur."
Ouch. But ok, disk space is cheap.
The killer point: it seems to be purely json based- no SQL of any kind. I'm not sure about that. A lot of code would have to be changed to fit that model.
Raw Prometheus: Isn't able to hold my data.
Thanos: I liked the project, it's architecture and ease of deployment, but after spending a non-trivial amount of time with it I wasn't able to setup any long-term caching. Thanos uses the prometheus storage format. So whenever I was querying one metric, it was downloading all metrics which were in the same block (all metrics basically afaik), this resulted in gigabytes/s of network traffic where it definitely wasn't necessary, and fairly long query times. (I used it with ceph) Though I know the maintainers were planning to add some kind of caching so this may be fixed. By using the native prometheus data format you also don't get storage space savings over it.
Cortex: Didn't spend any time on it, as I expected similar problems as with Thanos, so left it out for the end (which didn't came after all). I know it does contain a caching element.
Victoria Metrics: As far as I know it's very well engineered and performs great. But I see only one active maintainer so am afraid to use it.
M3DB: Requires a non-trivial amount of memory (I have 3 machines, each 128GB RAM to handle 70k writes/s each (though 1 was able to handle 120 and be stable)). However, with all machines on bunches of raid 0 ssd's, querying is quite snappy. You can set it up with different storage resolutions, so you get detailed data for recent queries, but also fast long range queries. It also uses a magnitude of storage space less than raw prometheus. The documentation is lacking in my opinion in terms of performance tuning, however, the code is well written, so I've just spent a while reading it and it exports very good metrics for itself. Network traffic between the m3coordinator (prometheus remote write gateway) and m3db nodes is kinda huge (5-10x the traffic prometheus->gateway) but that wasn't an issue. Another bonus is that it handles statsd metrics, though I haven’t yet tried that.
For anybody afraid of it operationally, I’ve had no problems. It mostly worked as is.
m3query did have some inconsistencies compared to prometheus in how queries using intervals evaluated. (and sometimes didn't return any data because of it).
Having stepped through the code I don't remember the reason why that was, but I ended up using prometheus instances using remote_read from the m3coordinator (gateway), working like a charm.
(I am a Cortex maintainer)
> It also uses a magnitude of storage space less than raw prometheus
AFAIK, Prometheus compression is about 1.2-3 bytes per datapoint. A magnitude less is 0.12-0.3 bytes - are these numbers correct?
I admit I’ve exaggerated a bit as Prometheus doesn’t support downsampling, in m3db I only keep 2 weeks of data at full resolution, 2 months at lower, and 5 years at even lower.
If you still want to be able to quickly graph and view old data, downsampling is the only way to keep your queries interactive.
Take for instance 30s data vs 10min data. 20x more computation, network exchange and everything else of that nature needs to happen.
Also if you want to keep only a subset of your data for a very long time, you need to have retention policies - otherwise you end up storing all that extra data forever.
At large numbers (terabytes to petabytes) this stuff is impactful, at smaller numbers (gigabytes) my points here are far less relevant.
It does work in m3db.
Following are the on-disk sizes of one replica:
2 weeks at 15s res: 90G
2 months at 1m res: 160G
3 months at 5m res: 60G
EDIT: The only thing is that m3db doesn't really downsample. You just create one namespace (table) for each resolution, and set the m3coordinator up so that it writes to each at the wanted interval, then set up a different retention for it. (this way you have duplicates in recent data)
M3 namespaces aren't set up for a specific resolution. The writer decides and you can write various series in different resolutions to one namespace in theory.
As someone who has auditioned it, briefly, let me assure you that it is certainly not the former, and only appears to be the latter due to a lot of cut corners and spec-violating implementations.
It’s hard to justify using tens of millions of dollars more of hardware more to run Cassandra.
This makes it tough to use solely either ScyllaDB, ClickHouse or Cassandra for that matter for metrics workloads at scale since they need to find a needle in a haystack - a few thousand time series amongst a set of millions to billions, where users only specify a subset of the dimensions on the metrics in any order they want to. This is hard to do without an inverted index.
The first reason is that open source platforms struggle beyond a certain scale due to architectural weaknesses, which becomes an ongoing operational headache. Most companies just deal with it but it gets worse as the workload grows.
The second reason is that it is expensive to run the open source platforms due to their very low efficiency. I've seen companies reduce their hardware footprint by a factor of 10 by rolling their own metrics/time-series implementations due solely to superior software design. When you are running a petabyte of metrics per day through these systems, that adds up to a lot of money.
tl;dr: it is technically straightforward for a company to design their own metrics infrastructure that massively outperforms the open source tooling, and the limitations of the open source implementations are often painful enough as the data models scale up that many companies do.