M3DB, a distributed timeseries database
m3db.io
m3db.io
Slides https://fosdem.org/2020/schedule/event/m3db/attachments/audi...
These open-sourcings seem a bit like PR pieces with no guarantees of any support or evolution after being published.
Uber also uses M3DB extensively internally and the project is nowhere near being abandoned or on life support: https://github.com/m3db/m3/commits/master
A pretty common story these days:
1. I built an X to solve Unicorn U’s problem.
2. Open source it, give talks on it
3. Leave Unicorn U and start a company based on X
4. ...
5. Profit(???)
Also, debatable whether or not Unicorn U actually needed a freshly built X instead of using existing tools/tech.
What is the actual gap that is not addressed with one or all of the following?
- Prometheus / Grafana [1]
- Datadog [2]
- Cloudwatch [3]
1. https://docs.docker.com/config/thirdparty/prometheus/
2. https://www.datadoghq.com/blog/introducing-live-container-mo...
3. https://docs.aws.amazon.com/AmazonCloudWatch/latest/monitori...
Datadog gets very expensive beyond a certain point; much bigger than most startups, but much smaller than Uber.
Cloudwatch is usually used as a source of metrics rather than a destination. It doesn’t have as rich a data model as Prometheus for custom metrics, and has a lot of quite limiting restrictions (like only 10 tags per metric).
Unless Uber is actively blocking contributions, it's not Uber's fault if no community formed around something they opensourced.
As for this being a PR piece, they could have achieved the same with just a detailed blog post and no code. It looks like a expensive PR piece if they have to opensource work that took probably hundreds of development hours.
But without any certainty around the roadmap, support, and longterm commitment by Uber to maintain these projects, they're nothing more than interesting repos amongst a sea of interesting repos.
The way Uber brands them suggests that they're suitable for use in production environments, but so far that hasn't been the case with anything they open sourced outside a narrow envelope that resembles their own operating model. Maybe this project will set a new trend, but so far nothing they put out gained any traction or became suitable for general purpose production use. In that regard, H3 and their other projects have remained at the level of decent 'show HN' pieces rather than something you'd ever use professionally. In other words, marketing.
Based on previous news coming out e.g. https://news.ycombinator.com/item?id=20931644
It seems like Uber had too big of an engineering department with too little work to do, so they started reinventing wheels. Which is cool if they're willing to support them in the long term, but so far that hasn't proven to be the case.
Protects that don't do that are therefore unlikely to remain interesting for long.
Former Uber engineer here. I can assure you that while our engineering team was massive, there was anything but too little work. If anything most engineers were massively overtaxed. Whether or not the work we were undertaking was meritious and valuable is an entire branch of philosophy I'm pretty sure.
Part of the struggle at big companies is that a lot of existing solutions just don't work. Let me use an example with chat. A few years ago Slack was evaluated as a replacement for HipChat, since Atlassian's outages had finally started affecting us during our own outages.
Everybody wanted to go to Slack, but the cost of Slack was tremendously prohibitive and the state of the service then (as I was told) was such that it could not support a company of Uber's size. Tremendous effort would have been undertaken by Slack to support Uber and they didn't want to expend that effort for a single customer. This was late 2015 early 2016.
There were tons of options, but ultimately an in-house chat software was created. At the time it seemed required to make our own highly reliable chat, considering how distributed engineering teams were. I think if you talk to anybody without the background of how chat evolved at Uber they would think the in-house chat project would have been a boondoggle.
Not all over-scoped engineering projects are actually so noble. There was certainly a ton of "reinventing wheels" going on. There was significantly more "these problems are really hard and I only have bad solutions."
Though if the result is ultimately, "nothing more than interesting repos amongst a sea of interesting repos" sign me up.
Are you affiliated with Uber?
You drank too much kool-aid. uChat was just a reskinned Mattermost.
I have heard that the team that put it together actually tried to hide that fact from the company (for the glory, I guess). But that could be apocryphal.
Mattermost didn't work out of the box, and certainly not the way and at the scale Uber needed it to. I'm not overly familiar with the technical details, but one thing in particular stands out as an example. There was a Town Hall channel that every user had to be a member of. This unfortunately did not scale, and not enough ACLs were available to limit all the ways users could use this universal room. Eventually they really fixed the problem, but it was a tremendous pain point for a long time. There were a lot of fundamentally "less than great" things about Mattermost that had to get updated to work for Uber.
There was the amusing time employees found out anybody could change the topic in the room, even if our chat permissions had been disabled. It was absolute chaos for at least an hour, I can't remember if it actually negatively impacted the deployment though it sounds vaguely familiar.
It's pretty telling of employees that badmouth the uChat team. That team ultimately was trying to do what they thought was best for the company, even if at the time it seemed like they bit off more than they could chew. There was no other engineering team so directly visible and exposed to the entire company internally like they were. People dismissive of their efforts are generally not used to the difficulty of making so many very vocal customers happy all at the same time, and could be more sympathetic.
https://news.ycombinator.com/item?id=19101617
The fact that uChat is commonly mentioned by Uber employees as being the "custom chat solution we built in-house". It leads me to believe that comment that the team tried to hide that it was built on open source.
I don't doubt for a second that scaling Mattermost for a huge organization like Uber was a big undertaking. But it seems disingenuous for people to always mention that Uber built uChat when it should be more like "Uber put a lot of work into Mattermost to scale it up."
To be honest, I don't much like Slack because I feel the desktop app doesn't feel like a real macOS app. And I don't use all these features. So in the end, IRC would be fine for me. But it wouldn't for the rest of the company.
Not a contradiction. Many of these tools are suitable for use in production, almost by definition, since they are being used in production, at Uber. They might or might not work in your environment out of the box. But they are certainly often likely a better starting point than an empty editor, even when they do not. Most of the ones I am familiar with, are happy to get PRs generalizing them to more varied environments.
> they're nothing more than interesting repos amongst a sea of interesting repos.
As someone who has open-sourced on GitHub: research prototypes hacked together for a research paper deadline in grad school, class projects, for-fun hacks, and also production tooling I built as a paid engineer, I'd say there is a big difference! :) And there would still be a big difference even if the later were somehow never touched again after the first "we are open-sourcing this!" commit.
That said, we do try to maintain the things we open-source. Standards of support vary because individuals maintaining these projects, and their situations, vary. This is true for non-OSS internal tools too. In my experience, having gone through the Uber OSS process twice, and having started it a third time and decided against releasing (yet?), Uber does try to make reasonably sure that it's open-sourcing stuff that will be useful and is planned to be maintained. At the same time, they have to balance it with making it easy to open-source tools, otherwise too many useful things would remain internal only.
Also, note, some of these tools have exactly one developer internally as the maintainer, and not even as their full time job. For example, I am the sole internal maintainer[1] for https://github.com/uber/NullAway and also have 3-4 other projects internally on my plate, most of which are in earlier stages and need more frequent attention[2]. If and when said developer leaves, effort is made to find a new owner. This is not always successful, particularly if the tool has become non-critical internally. Sometimes, leaving owners retain admin rights on the repos and keep working on the tool (Manu, NullAway's original author, co-maintains it), but I don't think anyone is suggesting that that should be an obligation.
Finally, obviously, nothing here is the official Uber position on anything, just my own personal observations. This doesn't represent my employer, and so on. I am also pretty sure most of this is not even Uber specific :)
[1] Not the only internal contributor! Also, there is one external maintainer, as mentioned a few sentences later. But in terms of this being anyone's actual responsibility...
[2] Just to clarify, I think between Manu's interest, my own, and it being relatively critical tooling at Uber, NullAway is pretty well maintained. But I can understand why that isn't always a given for all projects.
First, "guarantee" is the wrong word to take too literally here. Depending on how you want to look at it, there are no guarantees, even with guarantees.
But looked at more loosely, answering that is really Uber's responsibility. Why did they release it? If it is just a PR release, fire-and-forget works fine for that.
If they want to see wider adoption outside of their firm, there are some fairly obvious things they should do to foster that. Sometimes you release just the right thing at just the right time and everyone else does your evangelism and support work for you, but it is much more normal for your next great thing to take a while to build a user base.
Based on your specific complaints, it sounds like your opinion doesn't matter in this case; you're not the audience. You want support and a predictable future: you're looking for a product, not for technology. This isn't a product, and it's not a platform.
If instead you represented another company looking into solving this same problem yourself, and are looking at starting points, then you're the perfect audience. In that case, you'd have time and motivation to contact the developers directly rather than gripe on HN. You'd be less interested in whether there was an organized community, and more interested in how to directly influence the roadmap. You'd care about what the code looks like, how they solved Problem X and Problem Y, that kind of thing.
It looks to me like maybe their engineers internally are fans of the idea of "open source", and the PR department is happy to try and get some good press out of it, but the company culture isn't really set up to develop in public or maintain these things they've nominally "released".
Sadly, this isn't unusual among tech companies, but it'd be more obvious what's happening if they just put up a bare-bones FTP with a README: "Here's some code under <LICENSE>. Use it at your own risk."
A company should first put their own system through hell and decide that "yeah, this is good, we are sticking with this", before luring people to use it.
But they are probably in the business of selling themselves as a software/ tech company. People value tech companies, so you'd better be one, even if you do taxi services, produce and distribute tv series, or rent office space.
The test for open source is if it keeps getting maintained and supported for years. That only happens when the project is a core business effort, has some direct means of support (e.g. dual licensing or SaaS), or happens to be one of the few genuinely volunteer driven large scale open source projects.
We're also basically done with a new Python wrapper written in Cython. https://github.com/uber/h3-py/tree/cython
We could probably use some help with the last step of packaging, if anyone is interested!
Beringei, a TSDB, (https://github.com/facebookarchive/beringei) in particular with what you are saying was a PR piece (https://engineering.fb.com/core-data/beringei-a-high-perform...) since it was never really used by anyone outside of FB.
I really don't see the negative part of free code that you can learn from and/or incorporate at all.
When you rely on the systems themselves, and do so with expectation of support from the originating company, your expectations will almost certainly be broken.
I think that's the simplest takeaway - not to run away from any open-sourced project, but to take into proper consideration if/how they plan on supporting the tool, and how much you would be capable of adapting and owning yourself if the worst happened.
Chronosphere is the SaaS part for M3DB in this case. The negativity around someone open sourcing code for PR is nuts especially when all of the code is available. I love reading the code and getting ideas about how things work.
Open sourcing something is naturally more expensive than not. It’s seldom the case that impact to the community triggers contributions that outweigh that cost.
The fallacy we hold is that companies will prop up software that is open source for everyone to use despite the lacking community contribution.
We as engineers should push ourselves to contribute when we find issues - rather than simply create tickets that represent work were want to have done for free. This is how open source software dies.
There is a minority that does this.
I get your frustration, but everyone should remember there are never any promises of support with open source software, regardless of how well supported it is at a particular time.
https://github.com/richardartoul/tsdb-layer
The README does a really good job explaining the internals and motivation.
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.
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.
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.
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...
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?)
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.
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.
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.
> 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.
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)
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.
In the beginning, we didn't really have metrics, we had logs. Lots of logs. We tried to use Splunk to get some insight from those. It kinda worked and their sales team initially quoted a high-but-reasonable price for licensing. When we were ready to move forward, the price of the license doubled because they had missed the deadline for their end of quarter sales quota. So we kicked Splunk to the curb.
Having seen that the bulk of our log volume was noise and that we really only cared about a few small numbers, I looked for a metrics solution at this point, not a logs solution. I'd operated RRDtool based systems at previous companies, and that worked okay, but I didn't love the idea of doing that again. I had seen Etsy's blog about statsd and setup a statsd+carbon+graphite instance on a single server just to try out and get feedback from the rest of the engineering team. The team very quickly took to Graphite and started instrumenting various codebases and systems to feed metrics into statsd.
statsd hit capacity problems first, as it was a single threaded nodejs process and used UDP for ingest, so once it approached 100% CPU utilization, events got dropped. We switched to statsite, which is pretty much a drop-in replacement written in C.
The next issue was disk I/O. This was not a surprise. Carbon (Graphite's storage daemon) stores each metric in a separate file in the whisper format, which is similar to RRDtool's files, but implemented in pure Python and generally a bit easier to interact with. We'd expected that a large volume of random write ops on a spinning disk would eventually be a problem. We ordered some SSDs. This worked okay for a while.
At this point, the dispatch system was instrumented to store metrics under keys with a lot of dimensions, so that we could generate per-city, per-process, per-handler charts for debugging and performance optimization. While very useful for drilling down to the cause of an issue, this led to an almost exponential growth in the number of unique metrics we were ingesting. I setup carbon-relay to shard the storage across a few servers- I think there were three, but it was a long time ago. We never really got carbon-relay working well. It didn't handle backend outages and network interruptions very well, and would sometimes start leaking memory and crash, seemingly without reason. It limped along for a while, but wasn't going to be a long-term solution.
We started looking for alternatives to carbon, as we wanted to get away from whisper files... SSDs were still fairly expensive, and we believed that we should be able to store an append-only dataset on spinning disks and do batch sequential writes. The infrastructure team was still fairly small and we didn't have the resources to properly maintain a HBase cluster for OpenTSDB or a Cassandra cluster, which would've required adapting carbon- I understand that Cassandra is a supported backend these days, but it was just an idea on a mailing list at that point.
InfluxDB looked like exactly what we wanted, but it was still in a very early state, as the company had just been formed weeks earlier. I submitted some bug reports but was eventually told by one of the maintainers that it wasn't ready yet and I should quit bugging them so they could get to MVP.
Right around this time, we started having serious availability issues with metrics, both on the storage side- I estimated we were dropping about 60% of incoming statsd events, and on the query side- Graphite would take seconds-to-minutes to render some charts and occasionally would just time out. We had also built an ad-hoc system for generating Nagios checks that would poll Graphite every minute to trigger threshold-based alerts, which would make noise if Graphite was down and the monitored system was not. This led to on-call fatigue, which made everybody unhappy.
We started running an instance of statsite on every server which would aggregate the individual events for that server into 10 second buckets with the server's hostname as a key prefix, then pushed those to carbon-relay. This solved the dropped packets issue, but carbon-relay was still unreliable.
We were pretty entrenched in the statsd+graphite way of doing things at this point, so switching to OpenTSDB wasn't really an option and we'd exhausted all of the existing carbon alternatives, so we started thinking about modifying carbon to use another datastore. The scope of this project was large enough that it wasn't going to get built in a matter of days or weeks, so we needed a stopgap solution to buy time and keep the metrics flowing while we engineered a long term solution.
I hacked together statsrelay, which is basically a re-implementation of carbon-relay in C, using libev. At this point, I was burned out and handed off the metrics infrastructure to a few teammates that ran with statsrelay and turned it into a production quality piece of code. Right around the same time, we'd begun hiring for an engineering team in NYC that would take over responsibility for metrics infrastructure. These are the people that eventually designed and built M3DB.
TimescaleDB is a more versatile time-series database. It supports a variety of datatypes (text, ints, floats, arrays, json), allows for out-of-order writes and backfilling of old data, supports full SQL, JOINs between tables (eg for metadata), flexible continuous aggregates, native compression, and is backed by the reliability of Postgres. [0]
M3DB seems much more limited in scope [1]:
"Current Limitations
Due to the nature of the requirements for the project, which are primarily to reduce the cost of ingesting and storing billions of timeseries and providing fast scalable reads, there are a few limitations currently that make M3DB not suitable for use as a general purpose time series database.
The project has aimed to avoid compactions when at all possible, currently the only compactions M3DB performs are in-memory for the mutable compressed time series window (default configured at 2 hours). As such out of order writes are limited to the size of a single compressed time series window. Consequently backfilling large amounts of data is not currently possible.
The project has also optimized the storage and retrieval of float64 values, as such there is no way to use it as a general time series database of arbitrary data structures just yet."
If you're building it with a specific application and a concrete schema you can create which will result in fast queries and don't have requirements for arbitrary dimensions being specified for lookup, then yes it's great as a TSDB.
Prometheus, M3DB, etc all use an inverted index alongside the column store TSDB to help with metrics workloads.
With regards to trie vs inverted index for Graphite data, I'd actually still be inclined to say inverted index is better based on the amount of queries I saw at Uber with Graphite where people did `servers.*.disk.bytes-used` type queries which is way faster to do using an inverted index since you have a postings list for each part of the dot-separated metric name, rather than traversing a trie with thousands to tens of thousands of entries in index 1 host part of the Graphite name. This is what M3DB does[0].
[0]: https://github.com/m3db/m3/blob/b2f5b55e8313eb48f023e08f6d53...
Regarding the auto-rebalance feature, I cannot much more agree with you. It's something that clickhouse definitely need to handle internally.
I'm assuming this is an out of process inverted index used alongside ClickHouse? Or is it more of a secondary table contained by ClickHouse which can be searched to find the metrics, then the data is looked up?
The latter scales not as well with billions of unique metrics since it's always a scan across the unique metrics stored in the time window your query searches for (since any arbitrary dimensions can be specified, all must be evaluated). This is the drawback of PromHouse which is an implementation of Prometheus remote storage on top of ClickHouse - and the major reason why PromHouse was only ever a proof of concept rather than a production offering.
I agree that python implementation of graphite was not particularly fast but there was faster implementation in C that companies used first to significantly increase performance. Then coordination of storage backend becomes complex when you try to scale the initial design. This is where clickhouse really shine. It provide out of the box distributed storage with compaction, rollup and fast querying. The other layers are stateless, which means that they'll scale with your computing ressources.
M3DB is roughly doing the same thing as clickhouse but clickhouse is much more advanced database that has proven records of running at petabyte scale without a sweat. For example they now have tiered storage which means that you can store recent event in nvme and rollup to standard HDD...
[1]https://medium.com/avitotech/metrics-storage-how-we-migrated...
Any database can handle it, and columnstore RDBMS are designed to store and query trillion-row tables with full SQL functionality. The only advantage a "time-series" database gives you is some time-based query operators (like gap filling, last value, smoothing, etc). Those are now being added to SQL support for RDBMS so there's really nothing to be gained from a time-series database anymore.
That's a pretty big simplification. It's like saying one could go running in dress shoes. (Yes, it's possible, but don't you want to use the right tool for the job?)
As one time-series database example, because TimescaleDB [0] is focused on time-series data, its users benefit from [1]:
* 40-50x compression for metrics data (so storage costs for compressed data are 2-2.5% what they would normally be)
* Versatile continuous aggregate policies
* Variable data retention policies
* Overall much more efficient compute and memory utilization (because of faster insert and query rates)
* And yes, also time-based query operators for gap filling, first/last, LOCF, etc
(And I'm sure roskilli could describe M3DB's own advantages over non-time-series DBs.)
[0] Disclaimer to other readers, I'm a co-founder (although OP already knows this, as we've jousted on HN before :-) )
[1] All of our benchmarks and other engineering notes are published here: https://blog.timescale.com/tag/engineering/
At the very least, I think we both agree that the relational/SQL options work just fine compared to limited time-series databases like Influx. And for the record, we do use timescale so you've won me over on the PG usability front.
Great! Didn't realize that.
And yes, I would agree that a relational/SQL time-series database like TimescaleDB can work quite well. :-)
For anyone who's never heard of M3DB, and lives in a place where Uber doesn't operatore or is even banned (and so isn't part of daily life or conversation) "Ubers" might just as easily be some db researcher affiliated with the university of who knows where showing off something they came up with last summer and got a grant for.
- the car icon doesn't move as location updates come in (as evidenced by the route line getting shorter)
- the on-screen keyboard does not allow me to type anything after the first message (no other app on my phone has this problem)
- after I rate a driver, the app shows a map and none of the UI. I have to kill the app and restart it in order for it to be usable again. Picture this: I'm trying to get a ride, I open the app, I get nagged to rate the last driver. I agree just to be nice to the driver. (I should skip instead and get back to my task of getting a ride). After putting up with the nagware, the app fails 100% and I cannot complete the original task!
Alternatively, if they demand the collection of so many metrics, why don't they collect the metrics that would show them just how pathetically broken their Android app is?
Opentsdb for instance was built on top of Hbase because implementing one naively in hbase hits tons of performance issues.
Someone could probably elaborate on this a massive amount. I'm sure there is some nuance and a lot of relevant details around how that optimization is done.
Uber didn't invent the term, there are a lot of existing products in market.
Question is why none of them worked for them. I have used OpenTSDB and it worked great at Mastercard scale. What issues did Uber had?
Also with a fast inverted index we were able to achieve much faster query times than OpenTSDB at scale.
[0]: https://softwareengineeringdaily.com/2019/08/21/time-series-...
Teradata Vantage also supports both. And you're absolutely right, it's important not to conflate "temporal" and "time series" support.
Your regular RDBMS is going to be either write-heavy or read-heavy. You can pretty easily[ß] optimise the database for one of these utilisation patterns. But a TSDB basically combines the worst of both worlds: telemetry at any scale is important, and monitoring reliability in an always-online system is not optional.
TSDBs are written to very frequently; even at a reasonably low scale we could be talking about couple of hundred thousand writes every few seconds. But because they are also used for system-wide monitoring, they are read from all the time.
ß: a read-heavy regular DB has the ratio of reads:writes in thousands, perhaps millions; a write-heavy DB can be read from a couple of times every few seconds, but can be written to at a rate of tens of thousands of entries per second. You - or your expensive DBA - can optimise the DB for one of these patterns, but not for both. TSDBs have to support both patterns at the same time, so their internals have been geared to this one specific domain.
And this is the billion dollars mistake of the current devop culture. I believe we are doing monitoring wrong. Real time monitoring need no persistent storage. Troubleshooting does need persistent storage, but not monitoring, and unless your infrastructure is broken all the time then querying past data must occur only rarely.
From what i have seen, this mistake seems to stem from the web culture that tends to favor designs centered on a database. Whereas you want your real-time monitoring to be centered on stream processing, with one output to the persistent store for later retrieval.
There is no good reason for your dashboards nor your alerts to hit a database every few minutes, this design is just wrong.
I'm actually working on the prototype of a stream processor tailored at small scale network monitoring and would welcome any discussion/criticism on this topic. Notice that "small scale" for a stream processor is much larger than "small scale" for a database, and that i believe a good stream processor capable of running arbitrary persistent queries for monitoring on the order of a million data points per second should fit a single server and be more than enough to monitor an above the average sized business infrastructure.
This is similar to the insight you mentioned: I don't actually need to store a separate series per hostname in perpetuity just to know when the worst one is out of control. I can ask it to persist the top five individually, and then an average.
But when responding to a PagerDuty alert, the first thing I want to see is a plot of the alerting series. The next thing I want to see is a plot of every series about that service (bonus points if you can find the ones that have discontinuities with similar timing). So you're really not getting out of the "fast serving" business, just reducing the load on it.
When you do receive a page, you can then hit a database. This happens rarely enough (hopefully!) so that your database can be optimised for writes only, not for writes and reads, as $parent suggested.
I would even argue that this is a similar use case than dashboarding: when paged you want to see the last 3 days or so of accurate data + some longer term averages/trends/baseline. All of this could fit in RAM. At least in the current prototype I mentioned above the ambition is to be able to serve recent data for such dashboards out of RAM without bothering with a DB, and have the DB completely out of the way for the monitoring use case (while still having a persistent store for long term capacity planning / business analyses).
Assuming doubles are compressed and a collection rate of 1/10s, 3 days of data is <100KiB per timeseries, ie 0.5GiB for 10k source timeseries.