How we achieved write speeds of 1.4M rows per second
questdb.io
questdb.io
The process that calls the FLUSH simply appends more data to a file, relying on tunable mechanics the OS provides (i.e. the commit interval).
The reader of the data that is responsible for querying the data based on user requests would then have to use some kind of an index to locate the correct range of data from the appended-to files. A radix seems to be an excellent choice here.
Any data that arrives more than 10 seconds late is discarded. This would of course have or be a tunable parameter.
Isn’t that all you need from a system like this?
Using index would have been much simpler but also slower.
Once you flush your input buffer, everything becomes sequential so after your 10 second window you no longer have this problem. And even for the latest 10 seconds reading an index won’t be terrible as far as performance because your input buffer is relatively small (compared to your huge dataset).
As far as compaction, I suspect for the kinds of workloads you are using this kind of storage for you won’t get much in terms of deduplication. IoT devices sending sensor readings are unlikely to produce loads of stable readings and the time stamps will be ever increasing.
The devil, however, is in the details:
* What kind of in-memory buffer to use ? QuestDB doesn't want to use LSM or BTree, but its own data structure.
* How to flush the buffer ? Based on how you maintain your in-mem buffer, memcpy might not be good enough.
* When to flush the buffer ? You describe a tumbling window of 10 seconds, yet a sliding window might be more appropriate when dealing with late data, etc...
* Do you need background compaction ? Many write-optimize DB need background compaction to optimize for the read path.
* Do you need MVCC, ACID, etc... ? I'm not sure if QuestDB provides any of these, but it'll also affect how you design a DB storage engine.
You flush it by resetting the root node of your index. You don’t need to actually zero out the whole memory structure. You can actually initialize two different buffers and when the writer says FLUSH it’s simply handed the buffer to read while the input processor uses the other one to collect the next 10 seconds of data. Since writing the data to disk should take less than 10 seconds (or the whole thing doesn’t work anyways), this would again mean no memcpy at all.
Based on the problem description in TFA a tumbling window is perfectly appropriate here.
Compaction/deduplication wasn’t discussed in TFA, so I can’t tell if it is something that was a part of the 65k lines of code or not but would assume it wasn’t. You can do compaction within your input buffer easily enough and depending on the kind of data you store it might not be all that useful anyways (vehicle tracking comes to mind where compaction is basically useless because you are storing continuous values and time stamps will keep changing too).
You don’t need MVCC because this is essentially an append-only log file with an index. There aren’t multiple versions, it’s an event log (unless of course we are getting signals from the multiverse :)). You might or might not want ACID. The design QuestDB outlines, if I’m reading it correctly, basically drops inserts that are too old. That means that if you lose your input buffer because of a power outage, or simply restart the process, you can just wait for future updates from your inputs. Again, vehicle tracking comes to mind. Losing a late arriving packet from a minute ago doesn’t matter if you’ve gotten position updates since then. This means that the D part of ACID can be relaxed in that you’ll lose at most 10 seconds of data (or whatever tunable value you set), if you lose the buffer. The other properties are built into the design by the fact that your data structure is append-only: it’s atomic because you don’t have duplicate arrivals of the inserts (and if you do, that’s fine), it’s consistent because there is only one source of truth for it all at any given time and the whole thing is just a append-only file, it’s isolated because as soon as a write is inserted into the input buffer it’s available for the reader process, and it’s as durable as just writing to a file with OS primitives with an built-in in-memory buffer enabled.
Elasticsearch, IIRC, does this.
See also: https://www.youtube-nocookie.com/embed/b6SI8VbcT4w
I had a question about this part:
> We found out that this model does not fit all data acquisition use cases, such as out-of-order data
I know the whole blog post is about how you're now able to handle out-of-order data, but I'm curious - what use cases is QuestDB best for? Can you provide some concrete examples of where customers have ordered data and need QuestDB?
> QuestDB is the fastest open source time series database
Time series data is any type of data that always has a timestamp attached to it, such as logging, metrics, telemetry, etc.
The advantage advertised by QuestDB is that it is optimized specifically for this common use-case. The disadvantage of this being that it only allows clustered indexes for timestamp columns.
Another would be a manufacturing firm with thousands of sensors sending data to the database continuously, latency and delivery mechanism mean that timestamps often arrive out of order. There is a case study we did with CNC machines manufacturer DATRON [3].
Tracking moving vehicles (say cars, ships, drones, planes) on a map where the two primary axis are time and space - if you want to see the position of those vehicles in the past, you want to have the timestamps ordered.
Finally you can think of typical devops/monitoring/alerting use cases. Verizon uses QuestDB to monitor metrics for autoscaling decision to their Vespa engine that provides search, recommendation, and personalization to their users [4]. Here you want the series to be ordered by time.
[1] https://levelup.gitconnected.com/tracking-multiple-cryptocur...
[2] https://www.tradersinsight.news/ibkr-quant-news/optimizing-t...
Writing in batches of sizes dynamically determined by the backpressure mechanism seems to bring incredible performance benefits. Both SSDs and CPUs benefit hugely from any technique that can batch transactions at any level.
To the user of your database, what's the difference between 1-10ms best case, or 5-20ms average? To your CPU and flash, that 10ms extra batch width can mean the difference between 5000 and 500k inserts per second.
You must send the write ACK after your transaction has been written, which could increase latency in the worst case but you will still have much higher throughput.
Are you only evaluating running OPs rates compared to a theoretical maximum?
I am not attempting to do anything like this. The back pressure is an inherent property of the ring buffer abstraction I have selected for serializing transactions.
Batch size is effectively determined by how long the previous batch takes to complete. So, if one batch takes a little while, but is not at RingBufferMaximum, the next batch will be a bit larger (assuming constant system load), thus extracting a little bit more uplift from the benefits of running everything on hot cache lines. The maximum size of the ring buffer prevents things from getting out of control. You simply test for the best batch size and set that as the max for the buffer.
Best case scenario, your system is fully loaded and slamming through 4k+ batch sizes all day long. Worst case, you are doing a transaction per IO. I.e., what everyone else on earth is doing right now.
Any resources you would recommend about this kind of low level I/O programming and tuning?