ableegoldman commented on a change in pull request #9239:
URL: https://github.com/apache/kafka/pull/9239#discussion_r486619734
##########
File path:
streams/src/main/java/org/apache/kafka/streams/kstream/internals/KStreamSlidingWindowAggregate.java
##########
@@ -205,32 +221,67 @@ public void processInOrder(final K key, final V value,
final long inputRecordTim
}
}
}
+ createWindows(key, value, inputRecordTimestamp, closeTime,
windowStartTimes, rightWinAgg, leftWinAgg, leftWinAlreadyCreated,
rightWinAlreadyCreated, previousRecordTimestamp);
+ }
- //create right window for previous record
- if (previousRecordTimestamp != null) {
- final long previousRightWinStart = previousRecordTimestamp + 1;
- if (rightWindowNecessaryAndPossible(windowStartTimes,
previousRightWinStart, inputRecordTimestamp)) {
- final TimeWindow window = new
TimeWindow(previousRightWinStart, previousRightWinStart +
windows.timeDifferenceMs());
- final ValueAndTimestamp<Agg> valueAndTime =
ValueAndTimestamp.make(initializer.apply(), inputRecordTimestamp);
- updateWindowAndForward(window, valueAndTime, key, value,
closeTime, inputRecordTimestamp);
- }
- }
+ public void processReverse(final K key, final V value, final long
inputRecordTimestamp, final long closeTime) {
- //create left window for new record
- if (!leftWinAlreadyCreated) {
- final ValueAndTimestamp<Agg> valueAndTime;
- // 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 (leftWindowNotEmpty(previousRecordTimestamp,
inputRecordTimestamp)) {
- valueAndTime = ValueAndTimestamp.make(leftWinAgg.value(),
inputRecordTimestamp);
- } else {
- valueAndTime = ValueAndTimestamp.make(initializer.apply(),
inputRecordTimestamp);
+ final Set<Long> windowStartTimes = new HashSet<>();
+
+ // aggregate that will go in the current record’s left/right
window (if needed)
+ ValueAndTimestamp<Agg> leftWinAgg = null;
+ ValueAndTimestamp<Agg> rightWinAgg = null;
+
+ //if current record's left/right windows already exist
+ boolean leftWinAlreadyCreated = false;
+ boolean rightWinAlreadyCreated = false;
+
+ Long previousRecordTimestamp = null;
+
+ try (
+ final KeyValueIterator<Windowed<K>, ValueAndTimestamp<Agg>>
iterator = windowStore.backwardFetch(
+ key,
+ key,
+ Math.max(0, inputRecordTimestamp - 2 *
windows.timeDifferenceMs()),
+ // to catch the current record's right window, if it
exists, without more calls to the store
Review comment:
nit: `add 1 to upper bound to catch the current records'...` here and
elsewhere, it's not totally clear what this comment is referring to
##########
File path:
streams/src/main/java/org/apache/kafka/streams/kstream/internals/KStreamSlidingWindowAggregate.java
##########
@@ -205,32 +221,67 @@ public void processInOrder(final K key, final V value,
final long inputRecordTim
}
}
}
+ createWindows(key, value, inputRecordTimestamp, closeTime,
windowStartTimes, rightWinAgg, leftWinAgg, leftWinAlreadyCreated,
rightWinAlreadyCreated, previousRecordTimestamp);
+ }
- //create right window for previous record
- if (previousRecordTimestamp != null) {
- final long previousRightWinStart = previousRecordTimestamp + 1;
- if (rightWindowNecessaryAndPossible(windowStartTimes,
previousRightWinStart, inputRecordTimestamp)) {
- final TimeWindow window = new
TimeWindow(previousRightWinStart, previousRightWinStart +
windows.timeDifferenceMs());
- final ValueAndTimestamp<Agg> valueAndTime =
ValueAndTimestamp.make(initializer.apply(), inputRecordTimestamp);
- updateWindowAndForward(window, valueAndTime, key, value,
closeTime, inputRecordTimestamp);
- }
- }
+ public void processReverse(final K key, final V value, final long
inputRecordTimestamp, final long closeTime) {
- //create left window for new record
- if (!leftWinAlreadyCreated) {
- final ValueAndTimestamp<Agg> valueAndTime;
- // 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 (leftWindowNotEmpty(previousRecordTimestamp,
inputRecordTimestamp)) {
- valueAndTime = ValueAndTimestamp.make(leftWinAgg.value(),
inputRecordTimestamp);
- } else {
- valueAndTime = ValueAndTimestamp.make(initializer.apply(),
inputRecordTimestamp);
+ final Set<Long> windowStartTimes = new HashSet<>();
+
+ // aggregate that will go in the current record’s left/right
window (if needed)
+ ValueAndTimestamp<Agg> leftWinAgg = null;
+ ValueAndTimestamp<Agg> rightWinAgg = null;
+
+ //if current record's left/right windows already exist
+ boolean leftWinAlreadyCreated = false;
+ boolean rightWinAlreadyCreated = false;
+
+ Long previousRecordTimestamp = null;
+
+ try (
+ final KeyValueIterator<Windowed<K>, ValueAndTimestamp<Agg>>
iterator = windowStore.backwardFetch(
+ key,
+ key,
+ Math.max(0, inputRecordTimestamp - 2 *
windows.timeDifferenceMs()),
+ // to catch the current record's right window, if it
exists, without more calls to the store
+ inputRecordTimestamp + 1)
+ ) {
+ while (iterator.hasNext()) {
+ final KeyValue<Windowed<K>, ValueAndTimestamp<Agg>>
windowBeingProcessed = iterator.next();
+ final long startTime =
windowBeingProcessed.key.window().start();
+ windowStartTimes.add(startTime);
+ final long endTime = startTime +
windows.timeDifferenceMs();
+ final long windowMaxRecordTimestamp =
windowBeingProcessed.value.timestamp();
+ if (startTime == inputRecordTimestamp + 1) {
+ //determine if current record's right window exists,
will only be true at most once, on the first pass
Review comment:
nit: maybe you can just remove this comment, the code seems pretty
explanatory. We can see what @vvcephei thinks
##########
File path:
streams/src/test/java/org/apache/kafka/streams/kstream/internals/KStreamSlidingWindowAggregateTest.java
##########
@@ -78,16 +100,28 @@
public void testAggregateSmallInput() {
final StreamsBuilder builder = new StreamsBuilder();
final String topic = "topic";
-
- final KTable<Windowed<String>, String> table = builder
- .stream(topic, Consumed.with(Serdes.String(), Serdes.String()))
- .groupByKey(Grouped.with(Serdes.String(), Serdes.String()))
-
.windowedBy(SlidingWindows.withTimeDifferenceAndGrace(ofMillis(10),
ofMillis(50)))
- .aggregate(
- MockInitializer.STRING_INIT,
- MockAggregator.TOSTRING_ADDER,
- Materialized.<String, String, WindowStore<Bytes,
byte[]>>as("topic-Canonized").withValueSerde(Serdes.String())
- );
+ final KTable<Windowed<String>, String> table;
+ if (inOrderIterator) {
+ table = builder
Review comment:
Instead of specifying the whole thing for both cases, you could just
create a
```
final WindowBytesStoreSupplier supplier = inOrderIterator ? new
InOrderMemoryWindowStoreSupplier(...) : Stores.InMemoryWindowStore(...)
```
and then pass that into the `Materialized` without having to list the whole
topology out twice.
##########
File path:
streams/src/test/java/org/apache/kafka/streams/kstream/internals/KStreamSlidingWindowAggregateTest.java
##########
@@ -78,16 +100,28 @@
public void testAggregateSmallInput() {
final StreamsBuilder builder = new StreamsBuilder();
final String topic = "topic";
-
- final KTable<Windowed<String>, String> table = builder
- .stream(topic, Consumed.with(Serdes.String(), Serdes.String()))
- .groupByKey(Grouped.with(Serdes.String(), Serdes.String()))
-
.windowedBy(SlidingWindows.withTimeDifferenceAndGrace(ofMillis(10),
ofMillis(50)))
- .aggregate(
- MockInitializer.STRING_INIT,
- MockAggregator.TOSTRING_ADDER,
- Materialized.<String, String, WindowStore<Bytes,
byte[]>>as("topic-Canonized").withValueSerde(Serdes.String())
- );
+ final KTable<Windowed<String>, String> table;
+ if (inOrderIterator) {
+ table = builder
+ .stream(topic, Consumed.with(Serdes.String(), Serdes.String()))
+ .groupByKey(Grouped.with(Serdes.String(), Serdes.String()))
+
.windowedBy(SlidingWindows.withTimeDifferenceAndGrace(ofMillis(10),
ofMillis(50)))
+ .aggregate(
+ MockInitializer.STRING_INIT,
+ MockAggregator.TOSTRING_ADDER,
+ Materialized.as(new
InOrderMemoryWindowStoreSupplier("InOrder", 50000L, 10L, false))
+ );
+ } else {
+ table = builder
+ .stream(topic, Consumed.with(Serdes.String(), Serdes.String()))
+ .groupByKey(Grouped.with(Serdes.String(), Serdes.String()))
+
.windowedBy(SlidingWindows.withTimeDifferenceAndGrace(ofMillis(10),
ofMillis(50)))
+ .aggregate(
+ MockInitializer.STRING_INIT,
+ MockAggregator.TOSTRING_ADDER,
+ Materialized.<String, String, WindowStore<Bytes,
byte[]>>as("topic-Canonized").withValueSerde(Serdes.String())
Review comment:
Where does `topic-Canonized` come from? Also, if we need to set the
ValueSerde to String here, then wouldn't we need to do so for the in-order case
as well? Does that mean we don't actually need the `withValueSerde` thing here?
##########
File path:
streams/src/main/java/org/apache/kafka/streams/state/internals/InMemoryWindowStore.java
##########
@@ -68,7 +68,7 @@
private volatile boolean open = false;
- InMemoryWindowStore(final String name,
+ public InMemoryWindowStore(final String name,
Review comment:
Need to make this public for the child class used to parametrize the
test so we continue to test both forward and reverse directions. The
alternative would be to just create a new standalone `ForwardOnlyWindowStore`
test utility class and stick it in the same package
(btw, need to fix parameter alignment)
----------------------------------------------------------------
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]