-
Notifications
You must be signed in to change notification settings - Fork 1.8k
feat: Support recursive queries with a distinct 'UNION' #18254
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Conversation
Rely on aggregate GroupValues abstraction to build a hash table of the emitted rows that is used to deduplicate We might make things a bit more efficient by rewriting a hash table wrapper just for deduplication, but this implementation should give a fair baseline
tobixdev
left a comment
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
From my perspective this is a very nice and concise solution to the problem.
Furthermore, from my understanding this should also correctly terminate the recursion as only each unique row is pushed into the WorkTable and at some point (as it can be seen in the closure example) this will reach a fix point.
What I am also thinking about is test coverage. My gut feeling says there should be some test cases in the SQLite test suite that cover distinct recursion. Would this cause the extended test suite to fail? Ideally, this solution passes all these test cases now! 🥳 However, I am a bit unsure how this is setup currently.
Thank you!
CAVEAT: I am by no means a DataFusion (nor recurisve query) expert so take my comments with a grain of salt.
| } | ||
|
|
||
| /// Return a mask, each element true if the value is greater than all previous ones and greater or equal than the min_value | ||
| fn are_increasing_mask(values: &[usize], mut min_value: usize) -> BooleanArray { |
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
I think I understood what this function does, but I had a hard time with min_value. Maybe we can be more explicit here. Just some suggestions:
input parameter: min_value -> highest_group_id
// Always update the min_value to do de-duplication within a record batch.
let mut min_value = highet_group_id;May the integrating the comment in the doc comment for are_increasing_mask is also more than enough.
I think this assumes that the group ids are assigned in-order within the record batch but I think this is a valid assumption. Maybe someone more familiar with the aggregation infrastructure has more information on that.
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
I think this assumes that the group ids are assigned in-order within the record batch
yes, this is part of the GroupValues trait documentation.
I have rephrased the doc comment. I hope it's clearer now.
I have not renamed min_value to highest_group_id, the function does not depends on any specific semantic outside of creating the mask from its inputs. But happy to do the rename if you feel strongly about it.
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
That's perfectly fine. Just a suggestion 👍
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
I also found this confusing. Some suggestions:
- Rename the function to
new_groups_maskto reflect what it does - Rename
min_valuetomax_seen_group_idormax_emittedor something like that.
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Done! 0af5648
|
Sorry -- this PR hasn't been on my radar. I will put it on my review list and try and get it in the next few days |
6cc4434 to
48e8e33
Compare
alamb
left a comment
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
| } | ||
|
|
||
| /// Return a mask, each element true if the value is greater than all previous ones and greater or equal than the min_value | ||
| fn are_increasing_mask(values: &[usize], mut min_value: usize) -> BooleanArray { |
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
I also found this confusing. Some suggestions:
- Rename the function to
new_groups_maskto reflect what it does - Rename
min_valuetomax_seen_group_idormax_emittedor something like that.
|
@alamb Thank you! Suggestions applied |
|
Thanks again @Tpt |
|
Thank you! |
Thank you for your patience |
Rely on aggregate GroupValues abstraction to build a hash table of the emitted rows that is used to deduplicate
We might make things a bit more efficient by rewriting a hash table wrapper just for deduplication, but this implementation should give a fair baseline
Which issue does this PR close?
UNIONin recursive CTE #18140.Rationale for this change
Implements deduplicating recursive CTE (i.e.
UNIONinside ofWITH RECURSIVE) using a hash table. I reuse the one from aggregates to avoid rebuilding a full wrapper and specialization for types. Each time a batch is returned by the static or the recursive terms of the CTE, the hash table is used to remove already seen rows before emitting the rows and keeping them in memory for the next recursion step.What changes are included in this PR?
Reusing
GroupValuestrait implementations inside ofRecursiveQueryExecto get deduplication working.Are these changes tested?
Yes, some sqllogictests have been added, including ones that would lead to infinite recursion is deduplication where disabled.
Are there any user-facing changes?
No