Class RowSetUnionBatcher

java.lang.Object
io.deephaven.engine.rowset.RowSetUnionBatcher
All Implemented Interfaces:
SafeCloseable, AutoCloseable

public final class RowSetUnionBatcher extends Object implements SafeCloseable
Accumulates the union of row sets handed over one at a time, merging them a batch at a time rather than inserting each one into a growing result.

Inserting each row set into the result separately costs a pass over the result every time, which is quadratic when the inputs are disjoint and arrive in an order unrelated to their keys; a batch costs one union, which sorts itself first. The batch size is what bounds that: a caller passes how many row sets it expects to produce, and one no larger than maxBatchSize merges the whole input at once.

Merging every batch into one result would reintroduce the same problem one level up, at one pass per batch rather than one per row set. Instead the entries are held in two regions of a list of 2 * batchSize slots. A full batch collapses into a single row set that stays where it is, so the front of the list fills with collapsed groups while the back gathers the next batch. The groups are allowed to fill their half of the list; the batch that would need a slot past it folds everything into one instead. A result is therefore merged into again once per batchSize batches instead of once per batch, which is the same tree the multi-pass merge inside union builds, one level up.

Two things are handled here rather than at the merge. Empty row sets are dropped as they arrive, so a batch never spends a slot on one. A row set that falls clear of the last entry is folded into it in place instead of taking a slot of its own: past the end always, since the row set implementations splice rather than merge range by range, and below the start when the entry it would be folded into is no larger than it is. Input that arrives in ascending order therefore collapses into a single row set that build() hands over as it stands, with no union at all, as does a descending stream of row sets each at least as large as what it has already gathered.

build() is what releases the result to the caller. Everything this batcher is still holding at close() is closed, so a caller drives it from a try-with-resources block and a traversal that throws part way through abandons what it had gathered rather than handing back half a union.

  • Field Summary

    Fields
    Modifier and Type
    Field
    Description
    static final int
    Default for maxBatchSize.
    static int
    The most row sets gathered into one batch, whatever count a caller asks for: large enough that each merge amortizes the pass it costs, small enough that input driven by data rather than by the shape of the query cannot make this hold an unbounded number of row sets.
  • Constructor Summary

    Constructors
    Constructor
    Description
    RowSetUnionBatcher(long batchSize)
     
  • Method Summary

    Modifier and Type
    Method
    Description
    void
    add(@NotNull WritableRowSet rowSet)
    Add rowSet to the union, merging the gathered batch if it is now full.
    Merge whatever is outstanding and hand over the union of everything added since this batcher was constructed or last built.
    void
    Close everything not yet handed over by build().

    Methods inherited from class java.lang.Object

    clone, equals, finalize, getClass, hashCode, notify, notifyAll, toString, wait, wait, wait
  • Field Details

    • DEFAULT_MAX_BATCH_SIZE

      public static final int DEFAULT_MAX_BATCH_SIZE
      Default for maxBatchSize.
      See Also:
    • maxBatchSize

      @VisibleForTesting public static int maxBatchSize
      The most row sets gathered into one batch, whatever count a caller asks for: large enough that each merge amortizes the pass it costs, small enough that input driven by data rather than by the shape of the query cannot make this hold an unbounded number of row sets. This caps the batch rather than every merge: build() hands union the collapsed groups as well as the batch, at most 2 * batchSize - 1 row sets.

      Read from the RowSetUnionBatcher.maxBatchSize configuration property, default DEFAULT_MAX_BATCH_SIZE. The union builds a batch of small row sets in one linear pass, so a larger cap hands it more at once and leaves fewer batch results to merge afterwards; the cap is what bounds how many row sets are held while the batch is gathered.

  • Constructor Details

    • RowSetUnionBatcher

      public RowSetUnionBatcher(long batchSize)
      Parameters:
      batchSize - The number of row sets to gather before merging, which a caller passes as the number it expects to produce. Clamped to [1, maxBatchSize]: under the cap the whole input merges at once, and over it, or where the count is only an upper bound or no bound at all, the cap takes over and the count costs nothing to have passed. Taken as a long so that a caller counting rows rather than objects has nothing to narrow and no reason to know the cap. The cap is taken as configured.
      Throws:
      IllegalArgumentException - If the resulting batch size is not positive, or twice it would not fit an array, since the entries list must be able to hold the batch and the collapsed groups together
  • Method Details

    • add

      public void add(@NotNull @NotNull WritableRowSet rowSet)
      Add rowSet to the union, merging the gathered batch if it is now full.

      Ownership of rowSet passes here: it may be closed before this call returns, and the caller must not use or close it afterwards. A caller that keeps its row set hands over a copy, which is a copy-on-write reference that later mutation of the original does not disturb.

      Parameters:
      rowSet - The row set to add; ownership passes here
    • build

      public WritableRowSet build()
      Merge whatever is outstanding and hand over the union of everything added since this batcher was constructed or last built. Ownership passes to the caller, and this batcher is left empty and ready for more.
      Returns:
      A WritableRowSet containing every row key added, which the caller owns. Not necessarily a newly constructed one: a single outstanding row set is handed back as it stands.
    • close

      public void close()
      Close everything not yet handed over by build().
      Specified by:
      close in interface AutoCloseable
      Specified by:
      close in interface SafeCloseable