Hi Xinqi and Yuan,

Thank you for sharing the design and the detailed review comments. I am 
following this proposal closely, particularly the ordered-input property, 
distributed intermediate-state handling, and memory-management approach.

This is an important addition to the IoTDB table model, and I look forward to 
the next revision and implementation progress. I would also be happy to 
participate in further discussion and review.

Best regards,
Hongyin Zhang


> 2026年7月24日 17:02,Yuan Tian <[email protected]> 写道:
> 
> 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