How We Built a Vectorized SQL Engine
cockroachlabs.com
cockroachlabs.com
> The Go templating engine allows us to write a code template that, with a bit of work, we can trick our editor into treating as a regular Go file. We have to use the templating engine because the version of Go we are currently using does not have support for generic types.
Go’s lack of an expressive type system continues to disappoint :(
How do they deal with GC pause? I'm sure that they have an answer, but they'd have so much more latitude if they had the facility to reason about this from the ground up. They've sort of painted themselves into a corner now.
C, C++, or Rust would have been a better match and given me more faith in their product.
We're generally pretty happy with our choice to use Go for the database, even when it causes us pain like the kind you see in this post. Rust wasn't really ready when we started working on CockroachDB, and we think that we wouldn't have been able to make progress as quickly as we did if we had used C(++).
It is possible to write Go programs which efficiently use memory using object pooling just as you would in C. Go additionally can stack allocate many variables that would require a heap allocation in Java, where the GC is a much bigger issue.
According to the Go Wikipedia article, as of 2017:
> Garbage collection pauses should be significantly shorter than they were in Go 1.7, usually under 100 microseconds and often as low as 10 microseconds.
According to [1], "A trivial SELECT can take in the order of 0.1ms to execute server-side", i.e., 100 microseconds. Any I/O will of course significantly increase that.
This sounds like an entirely reasonable price to pay. Writing high-level functionality in C in 2019 would not give me much faith in their product.
[1]: https://www.2ndquadrant.com/en/blog/postgresql-latency-pipel...
Databases make extensive use of scratch space for undo logs and redo logs and especially while processing JOINs. It's tremendous pressure on the GC that results in frequent performance dips. With C/C++/Rust, you'd pipeline the query stages to malloc/free that space outside the hot path to minimize query response times.
I don't understand why you seem to assert that GC'd languages would force you to put allocations on the hot path. Surely in Go you can also allocate that space elsewhere?
Yes, that is a possibility. Typically there are two ways to deal with this in a GC'd language: (a) You know that GC is only triggered by certain operations that might allocate, and you avoid those operations in your critical hot section, or (b) you know how to disable the GC temporarily, while you are in that critical hot section. For Go the latter is apparently done by calling this function: https://golang.org/pkg/runtime/debug/#SetGCPercent
I'm not a Go developer myself, but I imagine people developing such low-latency software in Go know about this.
This only matters to databases that actually implement these optimizations. If a database does not then a GC has a much smaller impact. But it also means that, all other things being equal, the database will never be competitive with a non-GC design since those optimizations can provide an integer factor improvement in performance.
Also, many database engines aren't doing malloc() at runtime. 100 microseconds of CPU time is an integer factor larger than many atomic database operations in a fast database engine. An operation taking an order of magnitude longer than the scheduler would expect based on runtime state has real adverse consequences.
I haven't ever touched it myself, but, if the benchmarks that Cockroach Labs posts on their blog are to be believed, it's pretty respectable for what it is.
However, as data velocity and volume grow, optimization of absolute performance and hardware efficiency has an increasingly large impact on the cost of operating a database and the kinds of applications you can economically run. Databases that sacrifice major optimizations at an architectural level will have no chance to be competitive with databases that do not over the long term as average data volume and velocity grows. At the scale of database operations many companies are using today, these optimizations literally save them tens of millions of dollars per year on infrastructure.
CockroachDB produces what looks like a fine product and has every right to make whatever technical decisions they wish. That the architecture has made substantial performance sacrifices is not theoretical though, and they don't pretend to operate in markets that require any kind of absolute operational efficiency.
So, given that CockroachDB is a distributed system, I'd want to start by seeing some measurements to demonstrate that GC overhead is having a bigger impact than, say, network latency.
What's important to notice though is that what really matters about latency and predictability is _relative_ performance.
Now there is a category of databases that are built for speed of response first, some of which have some pretty spaced out architecture under the hood. Aerospike comes to my mind, together with ScyllaDB or Redis maybe. They trade off this speed with massive compromises in complexity, consistency and scalability.
GC would be disastrous for these, so they're all written in C(++).
The main premise of CockroachDB is geographic replication, durability and scalability while maintaining full consistence and serialization.
In order just to keep these promises, a large part of the design of Cockroach is that the DB needs to keep _waiting_ most of the time of any request until all geo replicated shards are consistent. We're talking dozens to hundreds of milliseconds here. And even if that wasn't the case, a full GC of 100 microseconds probably isn't even noticeable compared with the weight of complex SQL query planning and execution.
You may stink a lot of valid points about Go, but this isn't one of them.
The trick with Go seems to be that programs don't seem to generate much garbage to collect, so I don't know if there is much use comparing with Java apps. But garbage collection overhead is certainly a major issue to be wary of for any application needing consistently low response times.
Especially if you're working at the level where you want/need vectorization.
[1] https://medium.com/@valyala/measuring-vertical-scalability-f...
[2] https://medium.com/@valyala/insert-benchmarks-with-inch-infl...
[3] https://medium.com/@valyala/high-cardinality-tsdb-benchmarks...
https://github.com/cockroachdb/cockroach/blob/master/pkg/col...
After doing all this work, is it your opinion that (hypothetically) a columnar on-disk representation would be a bigger or smaller win than a fully-built-our vectorized execution engine?
IO is still a bottleneck without column-stores. The selectivity, compression and encoding alone can generate massive speedups because there's less data to process, and less to move through the execution pipeline.
But batch/vectorized processing on row-stores is gaining adoption. MemSQL and SQL Server also use similar techniques.
The term "columnar" covers a diverse set of architectures with very different operational characteristics. For OLTP (like CockroachDB), using a classic DSM-style columnar representation is going to offer poor write performance no matter what you do with the execution engine. On the other hand, if you are using one of the newer vectorized page representations (VSM) that are popular for mixed workloads, which are quasi-columnar but not DSM (nor one of the intra-page DSM hybrids like PAX), the loss of write performance may be minimal versus a classic row store (NSM) but with much faster query processing (faster than DSM for some types of queries). The execution engine design you would attach to any of these models is pretty different.
Note also that a practical limit on this kind of optimization is code complexity. While VSM-style on-disk representation and matching execution engine sounds like a nearly optimal hybrid of both NSM for write performance and DSM for query performance, an implementation that is general purpose and performs well across a diverse set of data models is massively more difficult and complex to build in practice so most database designers avoid it at all costs due to the engineering overhead. These tend to be more common when the set of supported data models are very limited at design time. It is a research area that still offers a lot of opportunity -- the literature mostly ignores parts of the productive design space that intrinsically have extremely high implementation complexity (too difficult to produce code in support of paper publication).
On another note: Kdb has figured out how to optimize the heck out of columnar data stores and vector processing. They would be my benchmark for these types of optimizations.
For queries with very high selectivity it seems like any gains from increased cache-friendliness or SIMD would be erased by still needing to make random forward strides through the data.
The only solution I can really think of is to make the index point to a block of N elements and then take advantage of vector processing on the matching blocks. Have you run into this issue?
I haven't been able to find any examples of column SQL engine discussions that don't require whole-table scans.
People just like complaining about names.
People like complaining about names that they wish things didn't have.