bvarghese1 commented on code in PR #29410:
URL: https://github.com/apache/flink/pull/29410#discussion_r4235478367
##########
flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/runtime/operators/over/AbstractNonTimeUnboundedPrecedingOver.java:
##########
@@ -422,6 +436,36 @@ RowData setAccumulatorAndGetValue(RowData accumulator)
throws Exception {
return aggFuncs.getValue();
}
+ /**
+ * Reads the accumulator stored for a sort key, as a copy the caller owns.
+ *
+ * <p>An accumulator can hold mutable content, either a data view or a
field such as a bitmap,
+ * an array or a {@code byte[]}. The aggregate functions change that
content in place, and the
+ * heap state backend hands back the very object it stores. Without a
copy, accumulating into
+ * one sort key's accumulator would also change the one state holds for
another.
+ *
+ * @param accKey the sort key to read the accumulator of
+ * @return a private copy of the stored accumulator, or null if there is
none
+ */
+ RowData getAccFromState(RowData accKey) throws Exception {
+ final RowData acc = accMapState.get(accKey);
+ return acc == null ? null : accSerializer.copy(acc);
Review Comment:
This affects the hot path. Should we copy only on the reads that need it
i.e. from `setAccumulatorOfPrevRow/setAccumulatorOfPrevId` ?
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]