Hi xinqi,

Do you think about pushing these four functions down to TableScan and
combining AggNode with TableaScan into an AggTableScanNode. If so, why do
you choose to not do that?

Best regards,
---------------------
Yuan Tian

On Fri, Jul 24, 2026 at 10:33 AM Xinqi Zhao <[email protected]> wrote:

> Hi Yuan and all,
>
> Thank you for the detailed review and suggestions. I have revised the
> technical design accordingly.
>
> Related issue:
> https://github.com/apache/iotdb/issues/17976
>
> The main changes are summarized below.
>
> Separate Ordered and Naive accumulator implementations
>
> The Ordered and Naive modes have been split into different concrete
> classes because they maintain different states and have different lifecycle
> requirements.
>
> For each function, an abstract base class now contains the shared
> validation, window semantics, and final calculation logic, while the
> concrete subclasses implement the corresponding state management:
>
> AbstractRateTableAccumulator
>
> OrderedRateAccumulator
> NaiveRateAccumulator
>
> The same structure is used for increase(), irate(), and delta(), as well
> as for the GroupedAccumulator path.
>
> The accumulator path and the algorithm are treated as two independent
> dimensions:
>
> The operator selects TableAccumulator or GroupedAccumulator.
> The aggregation step and ordering property select Ordered or Naive.
>
> Ordered implementations are used only for SINGLE aggregation and keep O(1)
> state per group. Naive implementations preserve all samples and support
> SINGLE, PARTIAL, INTERMEDIATE, and FINAL stages.
>
> Explicit derivation of the ordered-input property
>
> The updated design defines inputOrderedByTimeAscending as a
> per-aggregation physical property. It is derived in
> TableDistributedPlanGenerator.visitAggregation() after the final physical
> child of the AggregationNode has been established.
>
> The proof uses the final child’s OrderingScheme. For a particular
> aggregation:
>
> The second argument must be a directly addressable time_col symbol.
> Before time_col appears in the ordering prefix, only SQL grouping keys are
> allowed.
> time_col must use ascending order.
> Ordering symbols after time_col do not affect the proof.
> If the ordering is missing, uses descending time, contains a non-grouping
> symbol before time_col, or otherwise cannot be proven, the property is
> false.
> Arbitrary time expressions conservatively use the Naive implementation.
>
> This also means that Project, Exchange, Join, Union, HOP/TUMBLE, and other
> operators are handled through the ordering property of the finalized
> physical child. If they do not preserve and expose a sufficient
> OrderingScheme, the Ordered implementation is not selected.
>
> Only the following condition enables the Ordered implementation:
>
> step == SINGLE && inputOrderedByTimeAscending
>
> PARTIAL, INTERMEDIATE, and FINAL stages always use the Naive
> implementation.
>
> The current version also disables aggregation pushdown to
> AggregationTableScanNode and ExternalTsFileAggregationScanNode for these
> functions. The explicit time_col may differ from the physical TIME column,
> and the functions require raw samples rather than file, chunk, or page
> statistics.
>
> Common runtime validation flow
>
> A stateless RateFunctionValidation utility is now shared by Ordered/Naive,
> Table/Grouped, and the different aggregation stages.
>
> The validation order and behavior are explicitly defined:
>
> Rows with NULL value_col are ignored.
> If value_col is non-NULL, the required time and window arguments must be
> non-NULL.
> NaN and Infinity are rejected.
> Negative values are rejected for rate(), increase(), and irate(), while
> delta() accepts negative values.
> window_start must be less than window_end.
> Samples must satisfy window_start <= time_col < window_end.
> Window boundaries must be consistent within the same SQL group.
> Duplicate timestamps are rejected.
> Ordered implementations also reject an unexpected decrease in timestamps.
> After filtering, fewer than two valid samples produce NULL.
>
> Regarding NULL handling, I retained the requirement-level contract that
> only a NULL value_col causes a row to be ignored. A non-NULL value with a
> NULL time or required window boundary is treated as invalid input. This
> rule is now applied consistently in every execution path.
>
> Versioned Intermediate State protocol
>
> The exact BLOB formats are now specified as follows:
>
> rate(), increase(), and delta():
>
> version:int32
> windowStart:int64
> windowEnd:int64
> sampleCount:int32
> repeated { timestamp:int64, value:double }
>
> irate():
>
> version:int32
> sampleCount:int32
> repeated { timestamp:int64, value:double }
>
> The initial protocol version is 1.
>
> The state rules are now explicit:
>
> A partial state with no valid samples is represented as NULL.
> A one-sample state is preserved and encoded normally.
> NULL intermediate positions are skipped.
> Window boundaries are checked for consistency when states are merged.
> Samples from different states are merged without relying on state arrival
> order.
> Duplicate timestamps are checked after all states have been merged and
> sorted.
> The codec validates the version, header and total length, sample count,
> window values, timestamps, and sample values before creating the decoded
> buffer.
> Serialized-size calculations use exact arithmetic and must fit the
> ByteBuffer integer capacity.
>
> I also corrected the TimeValueBuffer.merge() issue identified in the
> previous pseudocode: the original size is retained as the copy offset, and
> the size is updated only once after copying.
>
> Incremental grouped memory accounting
>
> Grouped Naive accumulators no longer traverse every non-empty group after
> each input batch.
>
> TimeValueBufferBigArray maintains buffersRetainedBytes incrementally. Only
> buffers modified by the current operation, together with any ObjectBigArray
> capacity change, contribute to the memory delta.
>
> As a result:
>
> Appending one sample has amortized O(1) accounting cost.
> Merging a state is proportional to the number of merged samples.
> Processing a batch is proportional to the positions or groups actually
> visited.
> There is no additional O(blockCount × groupCount) scan.
>
> Final sorting is performed in place using introsort on the parallel
> timestamp and value arrays. Its worst-case time complexity is O(n log n),
> with O(log n) auxiliary space.
>
> Numerical edge cases
>
> The numerical rules have also been clarified:
>
> Timestamp subtraction uses BigInteger before conversion to seconds,
> avoiding long overflow while preserving exact tick differences.
> Window and sample intervals are validated before applying the increase ==
> 0 shortcut.
> Intermediate and final arithmetic results are checked for finiteness.
> A non-finite durationToZero does not cause zero-point truncation and is
> not itself treated as an error.
> sampleCount uses int and is updated with exact overflow checks. The
> serialized state is also limited by the maximum ByteBuffer size.
> All input values are represented internally as double. Therefore, INT64
> values whose absolute value exceeds 2^53 follow IEEE 754 double semantics,
> and unit-level differences are not guaranteed to remain exact. This
> limitation is now documented explicitly.
>
> The revised design continues to preserve complete samples during
> distributed aggregation and performs sorting, duplicate detection,
> counter-reset correction, and extrapolation only after all states have been
> merged.
>
> Best regards,
> Xinqi Zhao
>
>
>
>
>
> 原始邮件
> 发件人:Yuan Tian <[email protected]>
> 发件时间:2026年7月17日 11:57
> 收件人:dev <[email protected]>
> 主题:Re: [DISCUSS] Technical design for rate(), irate(), increase(), and
> delta() aggregate functions
>
>
> Hi Xinqi,
>
> Thanks for sharing the technical design. I reviewed the detailed design
> document and left several inline comments. Overall, reusing the existing
> table-model aggregation framework and preserving complete samples during
> distributed aggregation looks reasonable. Before implementation, I think
> the following points should be clarified or corrected.
>
>    1. Accumulator organization
>
> The ordered and buffered modes maintain substantially different states and
> follow different processing paths. I suggest reconsidering whether both
> modes should live in the same concrete accumulator class. A common abstract
> base class could contain shared validation and calculation logic, while the
> ordered and buffered implementations could be separate subclasses. This may
>
> reduce mode-specific branching and make their respective invariants clearer.
>
>    2. Derivation of the ordered-input property
>
> The design says that the ordered implementation is selected when the
> planner can prove that samples are ordered by time_col, but it does not yet
> define how this proof is made.
>
> The design should specify:
>
>    - At which physical-planning stage the property is derived.
>    - How per-group ordering, rather than only global ordering, is proven.
>    - How Project, Exchange, Join, Union, HOP/TUMBLE, and other operators
>    affect the property.
>    - Whether only a direct time_col symbol is supported; an arbitrary time
>
>    expression should conservatively fall back to the buffered implementation.
>    - That the property must be false whenever ordering cannot be proven.
>
> Although inputOrderedByTimeAscending is not meaningful for every
> aggregation function, keeping it in AggregationNode.Aggregation as a
> serialized per-aggregation physical hint seems acceptable if its scope and
> exact meaning are clearly documented. The derivation should happen after
> the physical input is finalized rather than during logical function
> analysis.
>
>    3. Runtime validation contract
>
> The design should define a common validation flow shared by
> ordered/buffered, Table/Grouped, and SINGLE/PARTIAL/FINAL paths. In
> particular:
>
>    - Ignore samples whose value or time is NULL before validating valid
>    samples.
>    - Reject NaN and Infinity.
>    - Reject negative counter values for rate(), increase(), and irate();
>    delta() may accept negative values.
>    - Require non-NULL window bounds and window_start < window_end.
>    - Require sample timestamps to be within [window_start, window_end).
>    - Require consistent window bounds within the same group.
>    - Reject duplicate timestamps.
>    - Return NULL when fewer than two valid samples remain.
>
> Without one explicitly defined common entry point and validation order, the
> different accumulator paths may produce inconsistent behavior.
>
>    4. Intermediate-state format and merge behavior
>
> The exact binary format should be specified, for example:
>
> rate/increase/delta:
> version | windowStart | windowEnd | sampleCount | samples
>
> irate:
> version | sampleCount | samples
>
> The design should also define:
>
>
>    - Whether a partial state with no valid samples is NULL or has sampleCount
>    = 0.
>    - How a one-sample partial state is preserved.
>    - How addIntermediate() handles NULL positions.
>    - How inconsistent window bounds and duplicate timestamps across states
>    are handled.
>    - How unknown versions, truncated states, trailing bytes, and oversized
>    BLOBs are handled.
>
> There is also a correctness issue in the current TimeValueBuffer.merge()
> pseudocode: it assigns size = mergedSize before copying and then copies
> from offset size, which will cause an out-of-bounds access. The old size
> should be retained as the copy offset, and size should be updated only once
> after the copy.
>
>    5. Grouped memory accounting
>
> The proposed grouped implementation recalculates retained memory by
> traversing every non-empty group after each input or intermediate batch.
> This adds approximately O(blockCount × groupCount) work and may approach
> O(n²) in unfavorable cases.
>
> It would be better to maintain the retained-size delta incrementally for
> only the buffers changed by the current batch, including any ObjectBigArray
> capacity change.
>
>    6. Numerical edge cases
>
> The extrapolation design should also address the following cases:
>
>    - Math.subtractExact(later, earlier) may overflow for two otherwise
>    valid INT64 timestamps.
>    - The final division performed by rate() or irate() may still produce
>    Infinity even if the extrapolated intermediate result is finite.
>    - The increase == 0 shortcut should occur only after validating the
>    window and sample interval.
>    - Converting INT64 values greater than 2^53 to double can lose unit
>    increments; the intended semantics should be documented.
>    - The upper bound implied by using an int for sampleCount should be
>    defined.
>
> These points do not change the overall direction of the design, but they
> affect correctness, consistency, and performance. I suggest making them
> explicit in the technical design before implementation begins.
>
> Best regards,
> ---------------------
>
> Yuan Tian
>
> On Fri, Jul 17, 2026 at 11:02 AM Xinqi Zhao <[email protected]> wrote:
>
> > Hi IoTDB community,
> >
>
> > Following the previous requirements discussion, I have prepared an initial
> > technical design for adding the Prometheus-like rate(), irate(),
> > increase(), and delta() aggregate functions to the IoTDB table model.
> >
> > Related issue:
> > https://github.com/apache/iotdb/issues/17976
> >
> > The main design is summarized below.
> >
> > Integration with the existing aggregation framework
> >
> > The four functions will be implemented entirely within the existing
> > table-model aggregation framework. No new SQL execution layer or
> > third-party dependency will be introduced.
> >
> > Both execution paths will be supported:
> >
> > TableAccumulator for global aggregation and streaming aggregation when
> > groups are fully ordered.
> >
> > GroupedAccumulator for HashAggregationOperator and
> > StreamingHashAggregationOperator.
> >
> > Each function will provide both TableAccumulator and GroupedAccumulator
> > implementations.
> >
> > The signatures are:
> >
> > rate(value_col, time_col, window_start, window_end)
> > increase(value_col, time_col, window_start, window_end)
> > irate(value_col, time_col)
> > delta(value_col, time_col, window_start, window_end)
> >
> > The value column supports INT32, INT64, FLOAT, and DOUBLE. Values are
>
> > converted to double internally, and all four functions return DOUBLE. Time
> > arguments support TIMESTAMP and INT64.
> >
> > Ordered and unordered implementations
> >
> > Each accumulator will support two internal modes rather than introducing
> > separate Ordered and Naive accumulator classes.
> >
>
> > When the aggregation step is SINGLE and the planner can prove that samples
> > within each group are ordered by time_col, the accumulator will use a
> > streaming implementation with fixed-size state:
> >
> > Time complexity: O(n)
> >
> > Additional space per group: O(1)
> >
> > Otherwise, the accumulator will buffer all valid samples and sort them
> > before the final calculation:
> >
> > Time complexity: O(n log n)
> >
> > Additional space per group: O(n)
> >
> > The two modes will have identical function semantics, validation rules,
> > and results.
> >
> > A new inputOrderedByTimeAscending property will be added to
> > AggregationNode.Aggregation. The ordered implementation will be selected
> > only when:
> >
> > step == SINGLE &amp;&amp; inputOrderedByTimeAscending
> >
> > PARTIAL, INTERMEDIATE, and FINAL aggregation stages will always use the
> > buffered implementation because they need to generate, merge, or consume
> > intermediate states.
> >
> > The ordered mode will still perform runtime checks for duplicate
> > timestamps and unexpected out-of-order input.
> >
> > Accumulator organization
> >
> > rate() and increase() will share common base classes because they use the
> > same counter-reset correction and boundary-extrapolation logic.
> >
> > The main TableAccumulator classes will include:
> >
> > AbstractRateIncreaseAccumulator
> >
> > RateAccumulator
> >
> > IncreaseAccumulator
> >
> > IrateAccumulator
> >
> > DeltaAccumulator
> >
> > Equivalent GroupedAccumulator classes will maintain state independently
> > for each groupId.
> >
> > In ordered mode, the accumulators retain only fixed state such as the
> > first, previous, and last samples, sample count, corrected increase, and
> > window boundaries.
> >
> > In unordered mode, samples will be stored in a new TimeValueBuffer.
> > Grouped accumulators will use TimeValueBufferBigArray to maintain one
> > buffer per groupId.
> >
> > Intermediate-state handling
> >
> > For distributed aggregation, PARTIAL, INTERMEDIATE, and FINAL stages must
> > preserve the complete sample sequence because counter-reset detection,
>
> > duplicate-timestamp validation, and the final result depend on global time
> > ordering.
> >
> > The intermediate state will contain:
> >
> > A state version
> >
> > window_start and window_end, except for irate()
> >
> > The number of valid samples
> >
> > All timestamp/value pairs
> >
> > PARTIAL stages collect and serialize their samples without calculating a
> > final result.
> >
> > INTERMEDIATE stages deserialize and merge sample buffers without sorting
> > them.
> >
> > The FINAL stage merges all states, sorts the complete sample set by
> > timestamp, checks for duplicate timestamps, and then performs the final
> > calculation.
> >
> > Deserialization will validate the sample count and serialized length to
>
> > reject corrupted intermediate states and prevent invalid memory allocation.
> >
> > Memory management
> >
> > The unordered implementation will integrate with the existing query
> > memory-management mechanism, following a design similar to the exact
> > percentile() accumulator.
> >
> > TimeValueBuffer and TimeValueBufferBigArray will report their retained
> > memory. Accumulators will use MemoryReservationManager to reserve or
> > release the difference whenever buffers grow, merge, or shrink.
> >
>
> > reset() will clear the calculation state, shrink oversized internal arrays
> > to their initial capacity, and release the corresponding query memory.
> >
> > Shared extrapolation logic
> >
> > A shared ExtrapolationUtil will implement the boundary-extrapolation
> > algorithm used by rate(), increase(), and delta(). It will be responsible
> > for:
> >
> > Converting timestamp differences to seconds according to the current
> > timestamp_precision
> >
> > Calculating the average sample interval
> >
> > Applying the 1.1-times extrapolation threshold
> >
> > Limiting extrapolation to half an average interval when a boundary is too
> > far away
> >
> > Applying counter zero-point protection for rate() and increase()
> >
> > Calculating and validating the final extrapolation factor
> >
> > Counter-reset detection and sample traversal will remain in the
> > accumulators rather than in ExtrapolationUtil.
> >
> > rate() will divide the extrapolated increase by the complete window
> > duration. increase() will return the extrapolated increase directly.
> > delta() will extrapolate the raw difference between the last and first
> > values without counter zero-point protection. irate() will use only the
> > final two samples and will not perform boundary extrapolation.
> >
> > All time conversions will respect the configured ms, us, or ns timestamp
> > precision, while rate() and irate() will always return values per second.
> >
> > Please review this initial design, especially the ordered-input property,
>
> > the complete-sample intermediate state, and the memory-management approach.
> > Any comments or alternative suggestions are welcome.
> >
> > Best regards,
> > Xinqi Zhao
>
>
>
>

Reply via email to