MetricsDB: TimeSeries Database for storing metrics at Twitter
blog.twitter.com
blog.twitter.com
I thought Circonus’ IronDB was open source but they’ve either brought it back closed source or I was mistaken.
In any case it seems weird that everyone is just rehashing the same space with metrics dbs and not moving to where Google/Circonus have been for a while.
Associated with their Materialized view and Liveview features, you can achieve observability [1].
[0]https://clickhouse-docs.readthedocs.io/en/latest
[1]https://www.altinity.com/blog/2019/11/13/making-data-come-to...
Edit: fixed broken link.
The aggregation functionality is best done in a separate layer so that you have the flexibility to change how you create your distributions without having to change your storage backend.
This is different than say how Prometheus stores histograms where it treats the bucket boundaries as a dimension in a fixed bucket histogram and the cell as the count of items in that bucket.
[0] https://www.circonus.com/2018/05/effective-management-of-hig...
> The servers checkpoint in-memory data every two hours to durable storage, Blobstore. We are using Blobstore as durable storage so that our process can be run on our shared compute platform with lower management overhead.
Companies like Twitter, Uber etc can spend many millions and won't notice so they can build tailored databases for themselves.
Will someone else use it? Probably not (sure unless this is something new or that differs a lot from existing solutions and gives a lot of benefits)
Because TimescaleDB is built on Postgres, it is also a relational time-series database. Ie you can store metadata / business data alongside your time-series data.
(Note: I work on TimescaleDB.)
In particular, I'd love to know if theres anything major that generic RDBMS's could do better here.
For example, double delta compression, or Gorilla for floating point numbers. For more, take a look at our open source VictoriaMetrics database, which uses all of such tricks.
- Availability preferred over consistency (you want your metrics when bad things are happening, even if they may not be 100% accurate)
- Extremely write heavy, with most data never being read (they mentioned that only 2% of data is ever read. in my experience at another large company, it was way less than 2%)
- For the data that is read, most of reads are for most recent points (mostly alarms, but some dashboards as well)
- Different SLAs for queries based on usage - ie. queries for alarms must be fast. Dashboards, and trend analysis - not so much.
If I didn't want to use MySQL or Postgres, I'd rank Prometheus #1 and Clickhouse #2.
The killer thing for Clickhouse is that Percona supports it, so if you want to outsource the installation, mgmt. and support, you can just write a check and get good results.
Also, Clickhouse is a column store with SQL, so you could use an instance for monitoring and another to replace Vertica or Greenplum or whatever so long as it has the client libraries you need.
I'm concerned with there being a support issue, as well as a smaller user base also affecting support in the long term.
Well, everybody with experience outsources monitoring now since it's a non-core cost center, unless there's a compelling scale or secrecy issue.
If RAM and CPU were free, I'd use MySQL or Postgres w/partitions because of their mgmt. features, tested replication and SQL.
But Prometheus or Clickhouse are 10-25x more efficient in terms of space, and often have much faster queries. The tradeoffs are bizarre HA gaps, lack of trained people, and ops groups are stuck supporting it.
I would never recommend monitoring with anything based on HDFS (OpenTSDB), written in Java (Cassandra), or in-memory for large clusters (InfluxDB.)
For monitoring under 200 nodes, anything will work.
If you only have a day to do something, just install Nagios and you'll get 99% of what you really need.
Source: DBA.
That has not been my experience. Quite the opposite every place I’ve been that outsourced monitoring ended up bringing back significant portions of their observability stack either for needing more control, different feature sets or because the outsourced solution was cost prohibitive.
I think it is true that lots of teams continue to outsource storage of metrics data but the outsourced vendors are not incentivized to make it easy to do retention/filtering/aggregation well.
Source: Have run observability stacks across a variety of domains.
In this scenario I find it odd that there are so many opensource projects with traction and success (one name above all, Prometheus) if they are addressed only to specialized companies. What you say applies to small companies or very big ones with lot of money to spend. Mid-sized companies in my experience prefer to spend money on in-house solution because at their volumes outsourcing is really costly. But that's just my small experience.
What "time-series" databases offer is more functionality around time being a primary component of the data. For example, automatic partitoning/sharding on time, various date handling techniques, better bucketing/gap filling/smoothing functions, data retention policies based on time, automatic rollups and aggregations, etc.
Some have custom data stores, some use key/value stores like Hbase or Cassandra, and others use relational databases. Using relational foundations offers more flexibility (like Timescale on top of Postgres) than the others like InfluxDB or OpenTSDB.
Time-series databases make specific architectural decisions and introduce advanced capabilities that enable orders of magnitude better insert/query performance, while also reducing storage cost via compression.
For example, TimescaleDB (where I work), is a relational time-series database built on top of Postgres, which means it includes all of the goodness within Postgres, but also achieves:
* 96%+ compression (ie only uses ~4% of the storage of Postgres) [1]
* 100x-1000x faster queries than Mongo [2]
* 10x higher inserts and 50-1000x faster queries than Cassandra [3]
If your RDBMS is good enough - then please keep using it. :-) But if query latency, insert performance, or cost are becoming concerns, then I'd suggest looking at a relational time-series database.
[1] https://blog.timescale.com/blog/building-columnar-compressio...
[2] https://blog.timescale.com/blog/how-to-store-time-series-dat...
[3] https://blog.timescale.com/blog/time-series-data-cassandra-v...
Time series DBs (and OLAP dbs in general) have very different trade-offs/needs than transactional DBs.
Not sure what you mean by this. TimescaleDB outperforms InfluxDB, another purpose-built time-series DB, on most workloads [0], especially on ones with high-cardinality [1].
In particular our key insight, which some may still find heretical and hard-to-believe, is that it is quite possible to produce best-in-class performance characteristics for time-series using a relational database. In particular, we have been able to add columnar compression to our row-oriented format resulting in 96%+ compression rates [2], and multi-node scale-out resulting in 10M+ inserts per second [3].
TimescaleDB is not designed for workloads where most queries touch all data points (ie full table scans). But that's not what time-series workloads look like.
[0] https://blog.timescale.com/blog/timescaledb-vs-influxdb-for-...
[1] https://blog.timescale.com/blog/what-is-high-cardinality-how...
[2] https://blog.timescale.com/blog/building-columnar-compressio...
[3] https://blog.timescale.com/blog/building-a-distributed-time-...
In some sense, the organization is making an expected value of information estimate, where there's a complex interaction between the data you have and actions you can take, especially to preempt or resolve issues.
1.) Down-sample data to aggregates. Keep the aggregated but drop the source data after a certain period of time. You can see historical trends down to the sampling interval but cannot drill down to individual observations.
2.) Use tiered storage. Keep recent data on NVMe SSD and migrate to high density HDD or object storage over time. This allows you to keep more data. Not all DBMS support this out of the box.
You can compress an arithmetic series down to a handful of bytes...
E.g., Gorilla uses delta-delta-encoding for time-stamps, and a novel XOR-based compression scheme for floats.
For more, here is a primer on Gorilla and other leading time-series algorithms:
https://blog.timescale.com/blog/time-series-compression-algo...
At 1.5 petabytes, they could just dump this in Google's BigQuery for the cost of an engineer and have full SQL power.
SQL also lets you do a lot of aggregations in the database instead of pulling back raw metrics. And BQ has flat-rate pricing and discounts for large clients.
Anyways the point is that it's well within the financial means of Twitter, especially when compared to developing and operating this proprietary system.
And then do that 10s of 1000s of times per minute, every minute, 24 hours a day.
You can see how the query costs add up...
I know it's easy to come on HN and play armchair architect but this is one area where you are just incorrect in your assumptions. Did Twitter need to build a(nother) in-house database? I don't know for sure. But I do know that BQ is not well suited to observability use cases.
The details in the article (1.5PB logical data, 3x replication factor, only 2% of metrics ever read) and the fact that Twitter already runs infra in Google Cloud seem to align well with BQ in my experience (10+ years in adtech building even more complex backends).
Please tell me what assumptions are incorrect so I can reconsider.
Yes, only 2% of metrics ever read. That was the same back then too. The kicker is that you can't be sure which will be read, so the systems are built to be able to read any of the metrics with the same SLA. This is especially critical during an outage where engineers will need to quickly read metrics that in many cases were only written only a few seconds ago that they wouldn't otherwise read and are not configured as part of an ongoing alert.
Additionally, the alerting infrastructure that runs on top of the TSDB is configured with 10k+ queries that run every minute. So even at 2% read, you're querying them over and over and over again because you always need the latest data plus whatever trailing data is needed to fulfill the alerting needs (trailing 10m, hour, day, month, etc). This also makes it a particularly hard caching problem.
I can't speak to the state of GCP at Twitter. I know they had started to migrate some things but when I was there it was all colo.
Could they use BQ? Sure it could probably be tooled with partitions/caching/etc to work I guess, but at the end of the day it's not what BQ is designed for which is data warehousing and BI. Could they have used something that wasn't a custom built in-house database? Yea almost certainly (No secret that there is some NIH going on at Twitter well before this)
This happens alll the time on HN so I'm not here to call you out specifically...but it's really easy to read a company's blog post about their infrastructure decisions and immediately scoff and say "psh why didn't they just use X?"...as if you now have the same context that the engineering team has by reading a single blog post.
https://aws.amazon.com/about-aws/whats-new/2018/07/amazon-s3...
Also, if they're planning to do alerting, S3 would not help with that.
PUTs are 1 cent per 2000:
When you work for a large company, you have 2 problems:
1) Every has-been manager bikesheds over metrics retention
2) you can't say no. See #1.
The cost of puts alone would kill you.
https://htx-pub.s3.amazonaws.com/Screen+Shot+2020-06-17+at+1...
If you think it might be useful and you want to try it out drop me an email michael @ heartex.net