Conversation
After the last two changes the wide aggregation's tail is fold, sort, build, and the sort is the only one still on a single thread. On the hundred thousand group bench at eight workers that is 2.7 ms of a 28 ms query, and it is 2.7 ms of chasing a random word per compare while the pool sits idle. Split first, sort after. The groups go into as many ascending key ranges as there are hands, by a binary search against pivots sampled every nth group, and then each hand sorts its own range and builds its rows before the ranges are laid end to end. The split pass is a few compares a group against a pivot list that stays in cache, against the seventeen a full sort costs on a hundred thousand groups, and every compare after it is on a hand. Groups sit in first seen order before any of this, which is the order their keys turned up in the scan, so a stride through them samples the keys fairly whether the column arrived shuffled or already sorted. An uneven range costs a hand some idle time and nothing else.
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Follow on to #766 and #771. Those two took the fold and the row build of a keyed aggregation off the single thread that used to do them. What was left in the tail was the ordering, and on the hundred thousand group bench at eight workers that is a couple of milliseconds of chasing a random word per compare while the pool sits idle.
Split first, sort after. The groups are cut into as many ascending key ranges as there are hands, and each hand then sorts its own range and builds its rows before the ranges are laid end to end. The cut is a binary search against a handful of pivots per group, against the seventeen compares a full sort costs on a hundred thousand groups, and the pivot list is small enough to stay in cache, so every compare that is actually chasing a key word ends up on a hand.
The pivots come from every nth group, taken before anything is ordered. Groups sit in first seen order at that point, which is the order their keys turned up in the scan, so a stride through them samples the keys fairly whether the column arrived shuffled or already sorted. An uneven range costs a hand some idle time and nothing else, so eight samples a hand is enough.
Same hand rule as before: one per worker that ran, not one per core, and under four thousand groups a hand it stays where it is.
Measured on gamingpc, quiet. The tail alone, meaning the sort and the row build and not the fold, timed in place over twenty four runs of each of two separately built binaries, read 3.53 ms against 2.07 at eight workers and 4.59 against 4.81 at one. The whole query over ten paired alternating runs went 28.3 ms to 25.8 at eight workers, 354 to 388 M rows/s, faster in eight of the ten pairs, and read 90.9 against 91.6 at one worker, which is unmoved. Five percent of a whole query is under what that box resolves in six pairs, which is why the tail is timed in place as well.
That shape scales 3.6x from one worker to eight now, against 3.2x before this and 1.7x three changes ago. The rest of the gap to the twelve group shape's 6.4x is the scan rather than the tail. Eight workers each build their own hundred thousand group table, and the random probe into thirty two megabytes of them is the part that does not scale. That is the next thing to look at, not this one.
Tests: the parallel answer has to equal the one hand answer exactly, on a shuffled key, on an already sorted key, and on the counting table whose groups live in the slots rather than the key words. The split has its own test that every group lands in exactly one range, that the ranges ascend, and that none of them takes half the groups.
Part of P2, zu#75.