ableegoldman commented on a change in pull request #9157:
URL: https://github.com/apache/kafka/pull/9157#discussion_r483305501
##########
File path:
streams/src/main/java/org/apache/kafka/streams/kstream/internals/KStreamSlidingWindowAggregate.java
##########
@@ -192,29 +213,125 @@ public void processInOrder(final K key, final V value,
final long timestamp) {
//create left window for new record
if (!leftWinAlreadyCreated) {
final ValueAndTimestamp<Agg> valueAndTime;
- //there's a right window that the new record could create -->
new record's left window is not empty
- if (latestLeftTypeWindow != null) {
+ // if there's a right window that the new record could create
&& previous record falls within left window -> new record's left window is not
empty
+ if (previousRecordTimestamp != null &&
leftWindowNotEmpty(previousRecordTimestamp, timestamp)) {
valueAndTime = ValueAndTimestamp.make(leftWinAgg.value(),
timestamp);
} else {
valueAndTime = ValueAndTimestamp.make(initializer.apply(),
timestamp);
}
final TimeWindow window = new TimeWindow(timestamp -
windows.timeDifferenceMs(), timestamp);
putAndForward(window, valueAndTime, key, value, closeTime,
timestamp);
}
- //create right window for new record
if (!rightWinAlreadyCreated && rightWindowIsNotEmpty(rightWinAgg,
timestamp)) {
- final TimeWindow window = new TimeWindow(timestamp + 1,
timestamp + 1 + windows.timeDifferenceMs());
- final ValueAndTimestamp<Agg> valueAndTime =
ValueAndTimestamp.make(getValueOrNull(rightWinAgg),
Math.max(rightWinAgg.timestamp(), timestamp));
+ createCurrentRecordRightWindow(timestamp, rightWinAgg, key);
+ }
+ }
+
+ /**
+ * Created to handle records where 0 < timestamp < timeDifferenceMs.
These records would create
+ * windows with negative start times, which is not supported. Instead,
they will fall within the [0, timeDifferenceMs]
+ * window, and we will update or create their right windows as new
records come in later
+ */
+ private void processEarly(final K key, final V value, final long
timestamp, final long closeTime) {
+ // A window from [0, timeDifferenceMs] that holds all early records
+ KeyValue<Windowed<K>, ValueAndTimestamp<Agg>> combinedWindow =
null;
+ ValueAndTimestamp<Agg> rightWinAgg = null;
+ boolean rightWinAlreadyCreated = false;
+ final Set<Long> windowStartTimes = new HashSet<>();
+
+ Long previousRecordTimestamp = null;
+
+ try (
+ final KeyValueIterator<Windowed<K>, ValueAndTimestamp<Agg>>
iterator = windowStore.fetch(
+ key,
+ key,
+ Math.max(0, timestamp - 2 * windows.timeDifferenceMs()),
+ // to catch the current record's right window, if it
exists, without more calls to the store
+ timestamp + 1)
+ ) {
+ KeyValue<Windowed<K>, ValueAndTimestamp<Agg>> next;
+ while (iterator.hasNext()) {
+ next = iterator.next();
+ windowStartTimes.add(next.key.window().start());
+ final long startTime = next.key.window().start();
+ final long windowMaxRecordTimestamp =
next.value.timestamp();
+
+ if (startTime == 0) {
+ combinedWindow = next;
+ if (windowMaxRecordTimestamp < timestamp) {
+ // If maxRecordTimestamp > timestamp, the current
record is out-of-order, meaning that the
+ // previous record's right window would have been
created already by other records. This
+ // will always be true for early records, as they
all fall within [0, timeDifferenceMs].
Review comment:
I think it means, for a generic out-of-order record, it's _possible_
that the previous record's right window will have already been created (by
whatever record (s) are later than the current one). But for an early record,
if `maxRecordTimestamp > timestamp`, then we _know_ that the previous record's
right window must have already been created (by whatever record(s) are within
the combined window but later than the current record).
This is relevant to setting `previousRecordTimestamp` because if
`maxRecordTimestamp >= timestamp`, the previous record's right window has
already been created. And if that's the case, we don't have to create it
ourselves and thus we don't care about the `previousRecordTimestamp`
Does that sound right Leah?
----------------------------------------------------------------
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.
For queries about this service, please contact Infrastructure at:
[email protected]