4 ms·
This algorithm seems to resemble HyperLogLog (and all its variants), which is also cited in the research paper. Using the same insight of the estimation value o
by pixelmonkey 2y ago
This algorithm seems to resemble HyperLogLog (and all its variants), which is also cited in the research paper. Using the same insight of the estimation value of tracking whether we've hit a "run" of heads or tails, but flipping the idea on its head (heh), it leads to the simpler algorithm described, which is about discarding memorized values on the basis of runs of heads/tails.
This also works especially well (that is, efficiently) in the streaming case, allowing you to keep something resembling a "counter" for the distinct elements, albeit with a error rate.
The benefit of HyperLogLog is that it behaves similarly to a hash set in some respects -- you can add items, count distinct them, and, importantly, merge two HLLs together (union), all the while keeping memory fixed to mere kilobytes even for billion-item sets. In distributed data stores, this is the trick behind Elasticsearch/OpenSearch cardinality agg, as well as behind Redis/Redict with its PFADD/PFMERGE/PFCOUNT.
I am not exactly sure how this CVM algorithm compares to HLL, but they got Knuth to review it, and they claim an undergrad can implement it easily, so it must be pretty good!
- hmottestad 2y agoIt’s also possible to use HLL to estimate the cardinality of joins since it’s possible to estimate both the union and the intersection of two HLLs. http://oertl.github.io/hyperloglog-sketch-estimation-paper/ http://oertl.github.io/hyperloglog-sketch-estimation-paper/
- j-pb 2y agoIt's a really interesting open problem to get the cost of these down so that they can be used to heuristically select the variable order for worst case optimal joins during evaluation. It's somewhere on the back of my todo list, and I have the hunch that it would enable instance optimal join algorithms. I've dubbed these the Atreides Family of Joins: - Jessicas Join: The cost of each variable is based on the smallest number of rows that might be proposed for that variable by each joined relation. - Pauls join: The cost of each variable is based on the smallest number of distinct values that will actually be proposed for that variable from each joined relation. - Letos join: The cost of each variable is based on the actual size of the intersection. In a sense each of the variants can look further into the future. I'm using the first and the second in a triplestore I build in Rust [1] and it's a lot faster than Oxigraph. But I suspect that the constant factors would make the third infeasable (yet). 1: https://github.com/triblespace/tribles-rust/blob/master/src/query.rs https://github.com/triblespace/tribles-rust/blob/master/src/...
- refset 2y agoHaving read something vaguely related recently [0] I believe "Lookahead Information Passing" is the common term for this general idea. That paper discusses the use of bloom filters (not HLL) in the context of typical binary join trees. > Letos join God-Emperor Join has a nice ring to it. [0] "Simple Adaptive Query Processing vs. Learned Query Optimizers: Observations and Analysis" - https://www.vldb.org/pvldb/vol16/p2962-zhang.pdf https://www.vldb.org/pvldb/vol16/p2962-zhang.pdf
- j-pb 2y agoThanks for the interesting paper! We now formally define our _God-Emperor Join_ henceforth denoted join_ge... Nice work with TXDB btw, it's funny how much impact Clojure, Datomic and Datascript had outside their own ecosystem! Let me return the favour with an interesting paper [1] that should be especially relevant to the columnar data layout of TXDB. I'm currently building a succinct on-disk format with it [2], but you might be able to simply add some auxiliary structures to your arrow columns instead. 1: https://aidanhogan.com/docs/ring-graph-wco.pdf https://aidanhogan.com/docs/ring-graph-wco.pdf 2: https://github.com/triblespace/tribles-rust/blob/archive/src/triblearchive/succinctarchive/succinctarchiveconstraint.rs https://github.com/triblespace/tribles-rust/blob/archive/src...
- refset 2y ago> Nice work with TXDB btw It's X.T. (as in 'Cross-Time' / https://xtdb.com https://xtdb.com), but thank you! :) > 1: https://aidanhogan.com/docs/ring-graph-wco.pdf https://aidanhogan.com/docs/ring-graph-wco.pdf Oh nice, I recall skimming this team's precursor paper "Worst-Case Optimal Graph Joins in Almost No Space" (2021) - seems like they've done a lot more work since though, so definitely looking forward to reading it: > The conference version presented the ring in terms of the Burrows–Wheeler transform. We present a new formulation of the ring in terms of stable sorting on column databases, which we hope will be more accessible to a broader audience not familiar with text indexing
- 2y ago
- jalk 2y agoIirc intersection requires the HLLs to have similar cardinality, otherwise the result is way off.
- willvarfar 2y agoJust curious, dusting off my distant school memories :) How do the HLL and CVM that I hear about relate to reservoir sampling which I remember learning? I once had a job at a hospital (back when 'whiz kids' were being hired by pretty much every business) where I used reservoir sampling to make small subsets of records that were stored on DAT tapes.
- michaelmior 2y agoI guess there is a connection in the sense that with reservoir sampling, each sample observed has an equal chance of remaining when you're done. However, if you have duplicates in your samples, traditional algorithms for reservoir sampling do not do anything special with duplicates. So you can end up with duplicates in your output with some probability. I guess maybe it's more interesting to look at the other way. How is the set of samples you're left with at the end of CVM related to the set of samples you get with reservoir sampling?
- _a_a_a_ 2y agoWas wondering the same, thanks for an answer.
- krackers 2y agoKnuth's presentation of this [1] seems very very similar to the heap-version (top-k on a uniform deviate) of reservoir sampling as mentioned in [2]. The difference is in how duplicates are handled. I wouldn't be surprised if this algorithm was in fact already in use somewhere! [1] https://cs.stanford.edu/~knuth/papers/cvm-note.pdf https://cs.stanford.edu/~knuth/papers/cvm-note.pdf [2] https://florian.github.io/reservoir-sampling/ https://florian.github.io/reservoir-sampling/ Edit: Another commenter [3] brought up the BJKST algorithm which seems to be similar procedure except using a suitably "uniform" hash function (pairwise independence) as the deviate instead of a random number. [3] https://news.ycombinator.com/item?id=40389178 https://news.ycombinator.com/item?id=40389178
- snewman 2y agoYou could merge these data structures as well. If the two instances to be merged are not at the same "round", take the one that's at an earlier round and advance it (by discarding half the entries at random) by the difference in rounds. Then just insert the values from one list to the other, ignoring duplicates; if the result is too large, discard half at random and increment the round number. I implemented exactly this algorithm at my previous employer, except that alongside each value, we stored an estimate of the number of times that value appeared. This allowed us to generate an approximate list of the most common values and estimated count for each value.
- finnh 2y agoWas there, reviewed the PR, can confirm. Hi Steve! Since then we've also tuned it up in a couple ways, in particular adding "skip" logic similar to fast reservoir sampling to trade some accuracy for the ability to not even look at the next N {M,G,T}B if you've already seen many many many matches. For non-selective searches over PB of data it's a good tradeoff, despite introducing some search-order bias.
- gwillen 2y agoMerging like that doesn't work -- it will tend to overestimate the number of distinct elements. This is fairly easy to see, if you consider a stream with some N distinct elements, with the same elements in both the first and second halves of the stream. Then, supposing that p is 0.5, the first instance will result in a set with about N/2 of the elements, and the second instance will also. But they won't be the same set; on average their overlap will be about N/4. So when you combine them, you will have about 3N/4 elements in the resulting set, but with p still 0.5, so you will estimate 3N/2 instead of N for the final answer. I have a thought about how to fix this, but the error bounds end up very large, so I don't know that it's viable.