1.1B Taxi Rides Using OmniSciDB and a MacBook Pro
tech.marksblogg.com
tech.marksblogg.com
Take this query:
SELECT cab_type,
count(*)
FROM trips
GROUP BY cab_type;
This is just counting occurrences of distinct values from a bag of total values sized @ 1.1B.He's got 8 cores @ 2.7GHz, which presumably can clock up for short bursts at least a bit even when they're all running all out. Let's say 3B cycles/core/second. So in .134 seconds (the best measured time) he's burning ~3.2B cycles to aggregate 1.1B values, or about 3 cycles/value.
While that's ridiculously efficient for a traditional row-oriented database, for a columnar scheme as I'm sure OmniSciDB is using, it's less efficient than I might have expected.
Presumably the # of distinct cab types is relatively small, and you could dictionary-encode all possible values in a byte at worst. I'd expect opportunities both for computationally friendly compact encoding ("yellow" is presumably a dominant outlier and could make RLE quite profitable) and SIMD data parallel approaches that should let you roll through 4,8,16 values in a cycle or two.
Even adding LZ4 should only cost you about a cycle a byte.
That's not to denigrate OmniSciDB: They're already several orders of magnitude better than traditional database solutions, and plumbing all the way down from high-level SQL to bit twiddling SIMD is no small feat. More that there's substantial headroom to make systems like this even faster, at least until you hit the memory bandwidth wall.
I think this is a good point. On GPUs, SIMT is effectively automatic vectorization, so our focus has been on the memory bandwidth wall (we make use of cuda shared memory in nvidia GPU mode for aggregates like the above query). Non-random access compression on GPUs also has been a nonstarter, at least historically. With more recent GPUs and more recent versions of CUDA, perhaps this is changing. But on CPUs, we have started looking into vectorization. There is a tradeoff, though -- the vectorization LLVM passes do add time to the compilation phase, and at subsecond query speeds that time isn't always worth it.
There are also a few other tricks to get closer to roofline performance. If you sort the input data on the key you're grouping by you can see small performance improvements, mostly from better cache locality. But, part of the "magic" of OmniSciDB is that you can group on any key and get good performance without ingesting, reindexing, etc.
SELECT cab_type, count(*)
FROM trips
GROUP BY cab_type;
From execution time it seems to me that this is a straight sum() of 32-bit integers. "cab_type" has two distinct values and if stored 32bit value for "green" is 0 and "yellow" is 1, straight sum of these integers will produce the desired outcome and explain performance. That said the same performance will not extend to key that has three or more distinct values.- cab_type has very few distinct values. so you can encode those values from 1 .. N and use an array of size N instead of the hash table
- you can build a “parallel scan”: split the rows evenly across many threads and each thread processes it’s on chunk
- the operation per row is very basic: you need to look up in the array and increment a value. so you can use SIMD to perform operations on multiple rows at the same time
- using some bit manipulation magic you can do the above on “encoded values”: you never need to convert cab_type bit represetation to an integer from 1..N
I have written this piece of code: https://github.com/questdb/questdb/blob/master/core/src/main...
This sums 64bit values and using AVX2 it will sum 1Bn in 0.26s. Incrementing conditionally will not be as fast and will throw vectorization out of the window too.
Yes. Branching will absolutely hurt. Good old x100 paper teaches how to avoid branching: http://cidrdb.org/cidr2005/papers/P19.pdf.
And of course there is no branching in MemSQL for this use case. And also no hashing b/c number of groups is small and you can use an array and not a hashtable.
Finally if you compress data rather than do the sum on an uncompressed array you will have a lot more compact data representation which would allow you not hit the memory bandwidth ceiling this quickly (4 threads)
When you say an array and not a hash table, do you just mean a simple perfect hash table indexed by the offset of the dictionary id? We use this fairly extensively for inputs of bounded domain (i.e. dictionary-encoded strings, moderately-sized integer ranges, even binned values, numeric or timestamp), but call it a perfect hashing. Assume we're talking about the same thing but wanted to clarify.
I’m still of an opinion that it’s important to demonstrate performance on more complex queries with joins, subqueries, subselects, and clustered data movements. The count(*), group by query is a very very simple case.
For any particular implementation of a tight inner loop like this, you could measure IPC via internal counters and it would be quite consistent, but it’s really hard to estimate without them.
And does it matter? Ultimately you care about how many cores at how many GHz you need to get the job done.
Something like `count(*)` needs to work well where you have no idea at all about the data.
You can also install the open source version of OmniSciDB, either via tar/deb/rpm/Docker for Linux (https://www.omnisci.com/platform/downloads/open-source) or by following the build instructions for Mac in our git repo: https://github.com/omnisci/omniscidb (hopefully will have standalone builds for Mac up soon). You can also run a Dockerized version on your Mac, but as a disclaimer the performance, particularly around storage access, lags a bare metal install.
Is this on your roadmap?
We'll plan to update that issue with the above.
The ClickHouse team was (obviously!) very interested in Mark's result and tried out OmniSciDB on the standard analytics benchmark that CH uses to check performance. Results are here: https://presentations.clickhouse.tech/original_website/bench...]
Anyway, really intriguing results from Mark. Looking forward to learning more about the source of the differences.
Disclaimer: I work at Altinity, which supports ClickHouse.
Edit: Fixed bad link
https://www.vertica.com/docs/9.2.x/HTML/Content/Authoring/Gl...
https://www.vertica.com/docs/9.2.x/HTML/Content/Authoring/Ad...
Yes, he does have the Intel GPU he mentioned, but if he paid $200 to upgrade the GPU as he claims, he would also have a dedicated AMD Radeon Pro 5500M 8GB.
Still, that could have been clarified in the article.
Consultants have a rule: "Never say No. Always say: this is how much it will cost".
It's a "no" but with a threshold.
People can elect to work on things voluntarily for personal reasons, but if someone asks them to fix or do something for them, then instead of saying no, they might say "how much will you pay me?" It's their time and they're under no obligation to offer it to you for free (though they can voluntarily choose to if they want).
His articles are extremely well documented, so there's nothing stopping anyone from performing the same benchmarks.
He lost me here...I get that it doesn't matter, but come on, if you don't know that your computer has a GPU other than the integrated graphics (that you admit you paid more to upgrade) then what are you really doing...
If he ran his benchmarks and then checked his hardware while writing this article then the author might've gotten confused by that.
COPY trips
FROM '/Users/mark/taxi_csv/*.gz'
WITH (HEADER='false');
> The above managed to complete in 31 minutes and 40 seconds. The resulting import produced 294 GB of data in OmniSciDB's internal format.I’m really curious how a simple import (no indexes or data typing) into SQLite would compare. But I don’t have 700GB of free SSD space to spare.
[1] https://tech.marksblogg.com/benchmarks.html
[2] https://tech.marksblogg.com/billion-nyc-taxi-rides-sqlite-pa...
https://tech.marksblogg.com/benchmarks.html
Caveat: these benchmarks only test the simplest of operations like aggregation (GROUP BY, COUNT, AVG) and sorts (ORDER BY). No JOINs or window operations are performed. Even basic filtering (WHERE) doesn't seem to have been tested. YMMV.
Disclaimer: No disrespect to ClickHouse here, it's an amazing system that I'm sure beats out OmniSci for certain workflows.
mkdir -p application
tar --strip=1 -C application -xf archive.tar
For example installing NodeJS from archive when composed with curl: OS=$(uname -s)-x64
VER=12.16.1
curl -#L "https://nodejs.org/dist/v${VER}/node-v${VER}-${OS,,}.tar.gz" \
| sudo tar --strip=1 -xzC /usr/local
IMHO thats probably the reason why most apps are having proprietary installers :)Is this what a keyword-stuffed URL looks like? This is terrible! This does nothing at all to communicate semantic meaning about the post's content