diff --git a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/KStreamReduce.java b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/KStreamReduce.java index 09e4fab48e557..cd6728320c93c 100644 --- a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/KStreamReduce.java +++ b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/KStreamReduce.java @@ -59,7 +59,6 @@ private class KStreamReduceProcessor extends AbstractProcessor { public void init(final ProcessorContext context) { super.init(context); metrics = (StreamsMetricsImpl) context.metrics(); - store = (KeyValueStore) context.getStateStore(storeName); tupleForwarder = new TupleForwarder<>(store, context, new ForwardingCacheFlushListener(context), sendOldValues); } diff --git a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/KTableAggregate.java b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/KTableAggregate.java index c53224e845637..1c44486e3cfde 100644 --- a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/KTableAggregate.java +++ b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/KTableAggregate.java @@ -75,22 +75,29 @@ public void process(final K key, final Change value) { throw new StreamsException("Record key for KTable aggregate operator with state " + storeName + " should not be null."); } - T oldAgg = store.get(key); - - if (oldAgg == null) { - oldAgg = initializer.apply(); - } - - T newAgg = oldAgg; + final T oldAgg = store.get(key); + final T intermediateAgg; // first try to remove the old value - if (value.oldValue != null) { - newAgg = remove.apply(key, value.oldValue, newAgg); + if (value.oldValue != null && oldAgg != null) { + intermediateAgg = remove.apply(key, value.oldValue, oldAgg); + } else { + intermediateAgg = oldAgg; } // then try to add the new value + final T newAgg; if (value.newValue != null) { - newAgg = add.apply(key, value.newValue, newAgg); + final T initializedAgg; + if (intermediateAgg == null) { + initializedAgg = initializer.apply(); + } else { + initializedAgg = intermediateAgg; + } + + newAgg = add.apply(key, value.newValue, initializedAgg); + } else { + newAgg = intermediateAgg; } // update the store with the new value diff --git a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/KTableReduce.java b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/KTableReduce.java index 38c5a11f34db1..70db64432d868 100644 --- a/streams/src/main/java/org/apache/kafka/streams/kstream/internals/KTableReduce.java +++ b/streams/src/main/java/org/apache/kafka/streams/kstream/internals/KTableReduce.java @@ -71,20 +71,25 @@ public void process(final K key, final Change value) { } final V oldAgg = store.get(key); - V newAgg = oldAgg; + final V intermediateAgg; - // first try to add the new value + // first try to remove the old value + if (value.oldValue != null && oldAgg != null) { + intermediateAgg = removeReducer.apply(oldAgg, value.oldValue); + } else { + intermediateAgg = oldAgg; + } + + // then try to add the new value + final V newAgg; if (value.newValue != null) { - if (newAgg == null) { + if (intermediateAgg == null) { newAgg = value.newValue; } else { - newAgg = addReducer.apply(newAgg, value.newValue); + newAgg = addReducer.apply(intermediateAgg, value.newValue); } - } - - // then try to remove the old value - if (value.oldValue != null) { - newAgg = removeReducer.apply(newAgg, value.oldValue); + } else { + newAgg = intermediateAgg; } // update the store with the new value