Skip to content

[Status Update] Simplifing Streams to be more textbook-like and have less state while keeping the same perf #23974

Description

@rluvaton

If you wanna look at before and after example look at:

Status:

Background

So, I start seeing that some code that on the surface should be dead simple in textbook form in reality the main entry point and the whole state handling is really complex

For example SortMergeJoin, on the surface the algorithm is simple

Algorithm

Taken from Sort-Merge Joins

Image

But due to all the following reasons:

  1. child return pending and have to keep state between calls
  2. programming language as opposed to pseudo code
  3. Flow is not linear as pseudo code since you can return only 1 value at a time
  4. needing to handle spill
  5. performance optimization

the actual implementation is really complex which means that doing any sort of PR there is hard to review or to understand

Tradeoffs

So I started with moving some stuff to async generators that solve some of the problems while adding others.
problem that it is solving:

  1. Child return pending and have to keep state between calls
  2. Flow is now linear since we can yield from the code and have the state and code in mind kept
  3. Fewer lines - since less code needed for holding and managing the state
  4. In some cases improve performance since going back to where you was in the state by making the function idompotent could be expensive

Problems that it is creating:

  1. Timing is more annoying, you should pause the elapsed_compute timer between
    i. yields - since you shouldn't count the time that the parent has done work between calling you
    ii. awaits on child streams - because you don't want to count the child time
    iii. awaits on stuff that you do async like reading a spill file - because even though the on the surface this is your work, you might overcount the elapsed time and the spill reading finish earlier but you did not woke up (up for debate)
  2. Reserving memory is less intuitive since the code look linear but you hand off control between awaits so the data that you hold should be reserved for
  3. You need to wrap the async generator with ObservedStream to track end time and output batches/rows/etc
  4. can introduce performance problems since we are adding async machinery if not used with caution

Alternative solutions

We already have RecordBatchReceiverStreamBuilder but the flow of data is different - i.e. it is push based and not pull based which means:

  1. more memory is being in the buffer (have at least 1 item pending
  2. You produce data even if not needed
  3. In case you are spawning blocked task to do all your work you can get to something like what was fixed here: Remove waits from blocking threads reading spill files. #15654

Initial Progress

I started the first PR in

and @pepijnve kindly extracted a trimmed down version for the async generators from genawaiter crate that I previously used in that PR and tokio async-stream crate in:

All the discussion for why we have our own implementation and also why not having macro implementation can be found in:

the first rewrite PR allowed for having the code rewritten to match simpler form:

Activity

  1. buraksenn commented on Jul 29, 2026

    @buraksenn
    Contributor

    I would like to try one of them to learn more about streams while trying to simplify it. Are you open to help on this?

  2. rluvaton commented on Jul 29, 2026

    @rluvaton
    MemberAuthor

    @buraksenn sure, try to pick something simple but complex enough that the added cost of async won't be too much

  3. buraksenn commented on Jul 29, 2026

    @buraksenn
    Contributor

    Thanks. I've looked into this a bit and will start working on SingleHashAggregateStream if you've no objections. Otherwise I can pick up what you advise too.

  4. saadtajwar commented on Aug 3, 2026

    @saadtajwar
    Contributor

    Hey @rluvaton - could I take a stab at CrossJoinStream?

  5. rluvaton commented on Aug 6, 2026

    @rluvaton
    MemberAuthor

    Sure

  6. saadtajwar commented on Aug 7, 2026

    @saadtajwar
    Contributor

    @rluvaton thanks! I'll have a PR up for it in the next week or so :)

  7. saadtajwar commented on Aug 12, 2026

    @saadtajwar
    Contributor

    @rluvaton PR up #24291 ! Thanks!

  8. saadtajwar commented on Aug 18, 2026

    @saadtajwar
    Contributor

    @rluvaton / @2010YOUY01 - going to take NestedLoopJoinStream next if that's OK!

  9. 2010YOUY01 commented on Aug 19, 2026

    @2010YOUY01
    Contributor

    @rluvaton / @2010YOUY01 - going to take NestedLoopJoinStream next if that's OK!

    Thank you! I have a few thoughts you could consider. I’m only around 70% confident in them, so please only use them as suggestion.

    • For the complexity in NLJ, we might want to keep explicit state management, with the current state represented as an enum, and let that coexist with the generator pattern. My reasoning is:
      a) We probably need a state-transition diagram to fully understand the implementation anyway, and structuring the code similarly makes it easier to reason about.
      b) Explicit states make the entry and exit conditions for each state visible, which may make the implementation safer.
    • We could split the regular joins (inner, left, right, full) and the semi/anti joins into separate streams. They are really different relational operations, and different optimizations tend to apply to each. The current approach combines them and therefore needs several flags/configurations to route the internal logic, which adds complexity. We have already made a similar split in sort-merge join.

    This second point may be slightly outside the scope of this project, but since we are already doing a fairly large refactor, I wanted to mention it briefly.

  10. saadtajwar commented on Aug 20, 2026

    @saadtajwar
    Contributor

    @rluvaton / @2010YOUY01 - going to take NestedLoopJoinStream next if that's OK!

    Thank you! I have a few thoughts you could consider. I’m only around 70% confident in them, so please only use them as suggestion.

    • For the complexity in NLJ, we might want to keep explicit state management, with the current state represented as an enum, and let that coexist with the generator pattern. My reasoning is:
      a) We probably need a state-transition diagram to fully understand the implementation anyway, and structuring the code similarly makes it easier to reason about.
      b) Explicit states make the entry and exit conditions for each state visible, which may make the implementation safer.
    • We could split the regular joins (inner, left, right, full) and the semi/anti joins into separate streams. They are really different relational operations, and different optimizations tend to apply to each. The current approach combines them and therefore needs several flags/configurations to route the internal logic, which adds complexity. We have already made a similar split in sort-merge join.

    This second point may be slightly outside the scope of this project, but since we are already doing a fairly large refactor, I wanted to mention it briefly.

    Thanks for the detailed thoughts @2010YOUY01 ! This definitely makes sense to me - it seems like it might be a bit easier to split this into separate PRs of:

    • First separating regular joins from semi/anti/mark joins (preserving the existing polling behavior/spilling/state machines)
    • Then converting the streams to generators independently

    If that makes sense to you then I'll go ahead and start on implementation :)

  11. 2010YOUY01 commented on Aug 21, 2026

    @2010YOUY01
    Contributor
    • First separating regular joins from semi/anti/mark joins (preserving the existing polling behavior/spilling/state machines)
    • Then converting the streams to generators independently

    If you decide to split the semi/anti/mark joins, perhaps we can

    PR1: Keep existing implementation, and implement a new stream for semi/anti/mark joins directly with the generator pattern

    if !standard_join:
        SemiAntiStream
    else:
        NestedLoopStream
    

    PR2: Implement the standard join similarly, and delete the legacy implementation

    If the generator pattern produces overall simpler solution, then directly implement it this way might be easier 🤔, but anyway this is still a full rewrite, can be a little tricky, we can try and see how to split the work down to smaller PRs further.

    BTW, we could open a new issue and continue the discussion there.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions