Ask HN: What DB to use for huge time series?
thanks in advance!
thanks in advance!
Some (constantly-growing) timeseries can be stored on a per-row basis, while other (static or older) timeseries can be stored in a packed form (e.g. an array column).
I find that most of the time, "Big Data" isn't really all that big for modern hardware, and so going through all of the extra software work for specialized data stores isn't really all that necessary. YMMV, of course, depending on the nature of your queries.
HN discussion: https://news.ycombinator.com/item?id=7809819
I totally agree. Most of useful "big data" is time-series data, and they aren't all that huge compared to images/videos/etc.
That being said, I think the reason to adopt something like Hadoop/MPP engines is not for storage but ease of querying: while Postgres can handle storing terabytes of data, joining two terabyte-scale tables can get a little iffy. This gets even more complex if you start packing data into array columns for space efficiency.
There is an argument to be made that historical/archival data aren't all that useful and thus do not need to be analyzed: that was definitely my assumption coming from finance. However, I've been surprised how far back some of our customers at Treasure Data go to mine insights from data.
If you have less than 10 billion items, Postgres will be fine, and is easier to manage IMO.
If you do use postgres, you should vertically partition the table. This will help keep indexes smaller, improve the the cache hit rate, vastly improve the ease with which you can drop older data, and make various other admin tasks easier.
I've done this in the past with a compound primary key of (topic_id, t) where t was a microseconds-past-the-epoch timestamp (bigint) unique within a topic. Then set up a parent table: CREATE TABLE events (topic_id, t, data_fields..) and "CREATE TABLE .. INHERITS events" from it into multiple subtables, named based on the timespan they will hold, like events_2013, events_2014.
Depending on how much data you have, either partition by day/month/year/etc. I partitioned every million seconds (~11 days), since that kept the resulting table sizes a bit more manageable (gigs not TBs).
Add a CHECK CONSTRAINT to each sub-table to constrain the timespan (ie, WHERE t BETWEEN ?? and ??).
When you do a SELECT * FROM events WHERE topic_id = 1 AND t BETWEEN $x AND $y ORDER BY t DESC; the query planner knows which sub-table(s) to query, and doesn't touch the other tables at all.
You can also add a BEFORE INSERT trigger to the parent table that inserts into the correct sub-table, otherwise get clients to compute the correct table name when inserting.
- What do the writes look like? If they are coming in a stream how many writes per second do you need to support? If they are a bulk load how large and frequent are the batches? Simple numerical values?
- What do the reads look like? How many queries per second do you need to support? How much data per query? How fast do the queries need to be? Will your queries be simple aggregations? Dimensional queries? Unique dimension value counts? Are approximations tolerated?
- How much history do you need to keep?
- What are your requirements for availability?
- What are your requirements for consistency?
- How fast does new data have to show up in reads?
Without more detail, you're going to get dozens of suggestions which may each be right for a particular case.
Writes: not totally sure in terms of how the data is being packaged before being sent yet, but it'll probably be more than 10 writes a second but less than 1000 initially(?). Not sure yet if we're aggregating and batching before sending or if we are, to what degree.
Availability: If it has brief breaks where it just misses some data (<3seconds?) probably not the worst thing, but really trying to avoid big gaps in the data.
Reads will likely be grabbing the last n records of a given set of sensors maybe with some light math on it if the query language supports it, though there might be an easier way to cache recent history and then only need to go to the big list for responding to a longer-term issue. Also the nature of reads is very subject to change since there's a bunch of use-cases for the data being kicked around and I haven't gone through what each use's reads would look like yet.
New data needs to show up in reads in soft-real time. The napkin-estimate indicates that we might be looking at asking for about 6-80MB returned per query as a generally large but perhaps not max query, bigger operations that dealt with legitimately huge amounts of data will probably be scheduled around lighter periods/put on different machines (not sure how adding more machines reading would impact since I don't know what db it will be yet).
Ideally keep as much history as humanly possible, possibly moving them to physical archival at some point (1yr+?).
I have no affiliation, other than being a customer. Its as close to a standard as you can find in finance.
There are many useful tutorials out there that let you try it out and you can usually get an eval version to try before you buy.
http://code.kx.com/wiki/Startingkdbplus/contents If you find something that is comparable in terms of performance and features, but cheaper, please mail me!! I would be very grateful.
(c) http://code.kx.com/wsvn/code/kx/kdb%2B/c/c/k.h (c#) http://code.kx.com/wsvn/code/kx/kdb%2B/c/c.cs
This guy is truly depraved.
k5 isn't that big (about 9 C files)
Specially loved this comment:
// remove more clutter
#define O printf
#define R return
#define Z staticComing from that background, C and especially C# must seem extremely verbose.
For example (from Wikipedia):
In K, finding the prime numbers from 1 to R is done with [0]:
(!R)@&{&/x!/:2_!x}'!R
And APL[1]: (~R∊R∘.×R)/R←1↓ιR
Its truly awesome stuff.[0]: http://en.wikipedia.org/wiki/K_(programming_language) [1]: http://en.wikipedia.org/wiki/APL_(programming_language)
I find my code very clear and readable in it.
Even the java driver is written similar to the C code :) http://code.kx.com/wsvn/code/kx/kdb%2B/c/jdbc.java
If you are interested in using HDF5 and PyTables to store time series data, check out this little library that I created: http://andyfiedler.com/projects/tstables-store-high-frequenc...
If you have mega huge data http://opentsdb.net/ seems pretty decent, however I have not tried it out.
I like InfluxDB and still use it.
A typical production instance of the time series database is based on four distinct Cassandra clusters, each responsible for a different dimension (real-time, historical, aggregate, index) due to different performance constraints. These clusters are amongst the largest Cassandra clusters deployed in production today and account for over 500 million individual metric writes per minute. Archival data is stored at a lower resolution for trending and long term analysis, whereas higher resolution data is periodically expired.
https://blog.twitter.com/2014/manhattan-our-real-time-multi-...
The upside is that OpenTSDB scales really well with hadoop cluster size, so you can just scale it up to handle more load.
The downsides are that their data schema and query format are optimized for data efficiency, not speed or flexibility. It's really easy to refine a search for a particular metric by filtering on tags, but it's really hard to do any sort of analysis across metrics, so you have to write your own glue on top of that which fetches the datapoints for the metrics you care about, and does its own aggregation.
That said, I've used Cassandra in the past for timeseries data as one of the useful queries that can be made is a range query (if the composite key is set up correctly)
Bad news is picking a system means understanding access patterns -- reading, not writing. Do you only need to look within a single user? That's much easier. If you have to query across users, or do things like (and I have no idea what your problem domain is, but if it's utility usage, things like average usage by zip or block; if it's wearables, activity by city, etc), stuff gets much harder. How granular do you need to be able to query, and how far back? What is the sla on a query: are results calculated in batch mode or on demand for a website? You often have to duplicate data in order to optimize one set for throughput access and the other set for minimal random query time. Can you get away with logarithmic granularity for queries, ie every sample is available for 1 month, every 3rd for the next month, every 10th for a couple months after that, etc. What windowing functions do you need to run, and how frequently do they need to be updated? What is the ratio of writes to reads? If you have to access random data quickly, eg for a site, can you calculate > 1 day back in batch mode, cache those results, and add the last 24h of data at runtime? etc etc etc.
You need to have some conversations with the data consumers.
Edit: and I've assumed these data are read-only; if you can update them, then there's far more difficulty.
Cassandra has a nice and simple architecture (every node is identical, no zookeeper roles etc), high write performance and scalability [1], and is fairly robust. My main piece of advice is to get the tables correctly set up. You need to know exactly what queries you want to make and design a table around that query (Cassandra only allows performant queries to be made, unless you go out of your way to set a flag). Whether a query is possible or performant depends on the key of the rows for the table, which may be a composite key. Take a look at the cassandra documentation for more details.
1. http://techblog.netflix.com/2011/11/benchmarking-cassandra-s...
On a side note, you can hook a Hadoop cluster up to SQL Server if you're into that kind of thing for storage.
Elasticsearch might seem like a strange option at first since it's historically a text search engine, but it's main datastructure is a compressed bit array which is ideal for OLAP processing.
We had to build our own Time-Series streaming / storage / query so we could handle millions of points per second and years of retention.
(we love ElasticSearch, though)
As another person mentioned, you're going to be looking at columnar databases (few/one rows, with a very large amount of columns) if you have truly large storage requirements. Since my data is still small, I'm sticking with Postgres for now.
I've seen a couple people mention OpenTSDB; another alternative to that is KairosDB[1], which adds Cassandra support and focuses on data purity[2] (OpenTSDB will interpolate values if there are holes).
And to echo another person, just forget about Graphite/Whisper. It uses a simple pre-allocated block format that will eventually cause problems when you want to change time windows.
[1]: https://github.com/influxdb/influxdb/issues/68
[2]: http://influxdb.com/docs/v0.8/api/continuous_queries.html
http://www.rackspace.com/blog/cloud-metrics-working-toward-a...
(Disclaimer: I am the Product Manager on that project)
I'm working on a timeseries database aimed at replacing graphite. It's just getting started, so it probably won't work immediately, but contributions are welcome. Currently the write performance is already better than graphite [1].
https://github.com/stucchio/timeserieszen
[1] This was one of the design goals. Whenever graphite receives a data point a disk seek is incurred - the data point must be appended to the timeseries file. Timeserieszen uses a WAL - data flowing in is immediately written, and periodically the WAL is rolled over into permanent storage.
I commented on some other graphite replacement projects at https://news.ycombinator.com/item?id=8368689
https://github.com/soundcloud/roshi
Roshi is basically a high-performance index for timestamped data. It's designed to sit in the critical (request) path of your application or service. The originating use case is the SoundCloud stream; see this blog post for details.
If your format is cast in stone you may also be able to get away with using flat-files. If you implement the List interface or something similar it would be very easy to integrate into your application. (Normally I wouldn't recommend flat-files for anything, but for time series it can be not a bad option, as much as that makes me cringe).
If there are any questions we can answer to help you make a more informed decision, drop us a line at support@influxdb.com or reach out to the community: https://groups.google.com/d/forum/influxdb
If in doubt, start with a traditional RDBMS. And ONLY after you profile your application and see exactly where your pain points are, do something crazy.
Have fun!
Here is a toy hand crafted time series storage design:
Say you are storing tuples of {<timestamp>,<datablob>}. Then querying it by timestamp.
Writer can store it in two files,open in append only mode only. One is the data file one is the index file. Data might look like:
<datablob1><datablob2>...
And an index file, it stores timestamps and offsets into the data files where the blobs are:
<timestamp1><offset1><timestamp2><offset2>...
If you need rolling fall-off. Then create new pairs of files every day (hour, week, month). And delete old ones as you go.
Then if you can ensure that your have time synchronization set up and timestamp are in increasing order (this might be hard). You can do binary searching. If you use rolling fall-offs. Then you can discard whole files periods based on the query range when you search.
All this would go into a directory. Reader and writer could be different processes. Your timestamp and offset sizes should be fixed length. Writer first appends to the data file and then writes the index. Reader knows how to find the last valid record by looking at the size of the file.
But sometimes depending on the requirements a file is enough. If you intimately know the and control the bytes that get written it is easier to understand and reason about your systems (that means optimizing it, scaling it, making it fault tolerant).
Also one way to make a resilient and fault tolerant database is to have less code running. Sometimes the base libc and unix offer a good and stable base on which it is easy to build. If you append the file in read or append only mode. You can rely on certain behavior now.
People in the past have bought into marketing crap and got stuff like MongoDB which would throw data over the fence and pray that it would be synced eventually (by default!). But heck it was WebScale(tm).
The talk is posted here. https://www.youtube.com/watch?v=ovMo5pIMj8M
The free plan allows 10M records per month with a maximum capacity of 150M.
Full disclosure: I work there.
So "massive" -- why not prototype on Postgres, and then migrate when you actually have projections on size.
Different orders of magnitude change the technology you work with. Additionally, the latency with which you need to access the metrics (real time, report based).
Cassandra is a pretty solid choice, Influx is really new to the game but is promising.
Druid is trusted by a lot of people, Metamarkets (the author) among them, but may or may not be what you need.
I'd spend some time talking to the people in #druid-dev on Freenode, they're friendly and can help guide you.
If accuracy doesn't have to be 100%, a number of options open up.
Anyway, no one has mentioned RRD tool yet: http://oss.oetiker.ch/rrdtool/
"RRDtool is the OpenSource industry standard, high performance data logging and graphing system for time series data. RRDtool can be easily integrated in shell scripts, perl, python, ruby, lua or tcl applications."
---
Data Acquisition
When monitoring the state of a system, it is convenient to have the data available at a constant time interval. Unfortunately, you may not always be able to fetch data at exactly the time you want to. Therefore RRDtool lets you update the log file at any time you want. It will automatically interpolate the value of the data-source (DS) at the latest official time-slot (interval) and write this interpolated value to the log. The original value you have supplied is stored as well and is also taken into account when interpolating the next log entry.
Consolidation
You may log data at a 1 minute interval, but you might also be interested to know the development of the data over the last year. You could do this by simply storing the data in 1 minute intervals for the whole year. While this would take considerable disk space it would also take a lot of time to analyze the data when you wanted to create a graph covering the whole year. RRDtool offers a solution to this problem through its data consolidation feature. When setting up an Round Robin Database (RRD), you can define at which interval this consolidation should occur, and what consolidation function (CF) (average, minimum, maximum, last) should be used to build the consolidated values (see rrdcreate). You can define any number of different consolidation setups within one RRD. They will all be maintained on the fly when new data is loaded into the RRD.
Round Robin Archives
Data values of the same consolidation setup are stored into Round Robin Archives (RRA). This is a very efficient manner to store data for a certain amount of time, while using a known and constant amount of storage space.
It works like this: If you want to store 1000 values in 5 minute interval, RRDtool will allocate space for 1000 data values and a header area. In the header it will store a pointer telling which slots (value) in the storage area was last written to. New values are written to the Round Robin Archive in, you guessed it, a round robin manner. This automatically limits the history to the last 1000 values (in our example). Because you can define several RRAs within a single RRD, you can setup another one, for storing 750 data values at a 2 hour interval, for example, and thus keep a log for the last two months at a lower resolution.
The use of RRAs guarantees that the RRD does not grow over time and that old data is automatically eliminated. By using the consolidation feature, you can still keep data for a very long time, while gradually reducing the resolution of the data along the time axis.
Using different consolidation functions (CF) allows you to store exactly the type of information that actually interests you: the maximum one minute traffic on the LAN, the minimum temperature of your wine cellar, ... etc.
Historians like Pi etc will 'compress' time series by only storing data points where data has changed by some threshold. If you look back 5 years all the resolution is still there.
Here are some resources:
Webinar: the Elk Stack in a Devops Environment http://www.elasticsearch.org/webinars/elk-stack-devops-envir...
Webinar: An Introduction to the ELK Stack http://www.elasticsearch.org/webinars/introduction-elk-stack...
[1] https://github.com/imperialwicket/postgresql-time-series-tab...
Disclaimer: I'm a former core contributor to blueflood.
* I've done it in PostgreSQL using triggers and table inheritance. With this technique trimming old data is as simple as dropping old tables.
* Logstash folks use daily indices on ElasticSearch to store log data which is time series by nature.
* I have heard from quite a few people that Cassandra works really well with this data model too.
I'm now glad I never made the jump... in the meantime, pgsql is still on my list
It is also much more highly regarded as a primary data store than Elasticsearch.
For data storage Cyanite [0] speaks the graphite protocol and stores the data in Cassandra. Alternately, InfluxDB [1] speaks the graphite protocol and stores the data in itself
To get the data back out, there's graphite-api [2] which can be hooked up to cyanite [3] or influxdb [4]. You can then connect any graphite dashboard you like, such as grafana [5], to it.
[0] https://github.com/pyr/cyanite [1] http://influxdb.com [2] https://github.com/brutasse/graphite-api [3] https://github.com/brutasse/graphite-cyanite [4] https://github.com/vimeo/graphite-influxdb [5] http://grafana.org
Last I looked at Graphite I balked at the data store design (very I/O heavy) and the awful front ends (very limited graphing and reporting capabilities). But I haven't discovered a good alternative that has traciton. Diamond seems like the thing to use for collecting metrics (instead of collectd), though.
Edit: Grafana looks good, actually.
I have personally seen millions of records saved per minute on a top end SSD server.
We're using it in production... it's still early but there are about 1-2 dozen moderate sized installs (like 10 box installs).
We're pretty happy with it so far..
http://blog.foundationdb.com/designing-a-schema-for-time-ser...
One of our largest customer installations is for this purpose.
I'm not affiliated with them, I just met them once.
SciDB: http://www.scidb.org/
Paradigm4: http://www.paradigm4.com/
... anyone have experience of using SciDB?
We offer cyclic time warp search, time warp search and any other metric of pseudo-metric you can come up with :-)
As for the customers on TempoDB, we are working them to transition to TempoIQ if the switch makes sense or offering to guide them in a transition to another time-series database like InfluxDB.
I am currently working on a project analyzing massive amounts of options data and have found this approach to be both quite easy as well as flexible to work with... and as my project matures I may move select parts of it into a database.
What is "massive" for you? I was under impression you can't use R or pandas for anything that doesn't fit into memory.
oracle too. Just did something relatively small with that (< 10MM rows), but it's pretty solid.