Skip to content
Merged
Show file tree
Hide file tree
Changes from 13 commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1,373 changes: 43 additions & 1,330 deletions streams/src/main/java/org/apache/kafka/streams/kstream/KGroupedStream.java

Large diffs are not rendered by default.

Large diffs are not rendered by default.

1,151 changes: 80 additions & 1,071 deletions streams/src/main/java/org/apache/kafka/streams/kstream/KStream.java

Large diffs are not rendered by default.

832 changes: 20 additions & 812 deletions streams/src/main/java/org/apache/kafka/streams/kstream/KTable.java

Large diffs are not rendered by default.

Original file line number Diff line number Diff line change
Expand Up @@ -72,12 +72,6 @@ Set<String> ensureJoinableWith(final AbstractStream<K> other) {
return allSourceNodes;
}

String getOrCreateName(final String queryableStoreName, final String prefix) {
final String returnName = queryableStoreName != null ? queryableStoreName : builder.newStoreName(prefix);
Topic.validate(returnName);
return returnName;
}

static <T2, T1, R> ValueJoiner<T2, T1, R> reverseJoiner(final ValueJoiner<T1, T2, R> joiner) {
return new ValueJoiner<T2, T1, R>() {
@Override
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -39,13 +39,21 @@ public Long apply() {
return 0L;
}
};

final Aggregator<K, V, Long> countAggregator = new Aggregator<K, V, Long>() {
@Override
public Long apply(K aggKey, V value, Long aggregate) {
return aggregate + 1;
}
};

final Initializer<V> reduceInitializer = new Initializer<V>() {

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This is moved from TimeWindowedStreamImpl, just to be consistent with the other const functions.

@Override
public V apply() {
return null;
}
};

GroupedStreamAggregateBuilder(final InternalStreamsBuilder builder,
final Serde<K> keySerde,
final Serde<V> valueSerde,
Expand Down

Large diffs are not rendered by default.

Original file line number Diff line number Diff line change
Expand Up @@ -48,7 +48,6 @@ public class KGroupedTableImpl<K, V> extends AbstractStream<K> implements KGroup

protected final Serde<K> keySerde;
protected final Serde<V> valSerde;
private boolean isQueryable;
private final Initializer<Long> countInitializer = new Initializer<Long>() {
@Override
public Long apply() {
Expand Down Expand Up @@ -78,83 +77,6 @@ public Long apply(K aggKey, V value, Long aggregate) {
super(builder, name, Collections.singleton(sourceName));
this.keySerde = keySerde;
this.valSerde = valSerde;
this.isQueryable = true;
}

private void determineIsQueryable(final String queryableStoreName) {
if (queryableStoreName == null) {
isQueryable = false;
} // no need for else {} since isQueryable is true by default
}

@SuppressWarnings("deprecation")
@Override
public <T> KTable<K, T> aggregate(final Initializer<T> initializer,
final Aggregator<? super K, ? super V, T> adder,
final Aggregator<? super K, ? super V, T> subtractor,
final Serde<T> aggValueSerde,
final String queryableStoreName) {
determineIsQueryable(queryableStoreName);
return aggregate(initializer, adder, subtractor, keyValueStore(keySerde, aggValueSerde, getOrCreateName(queryableStoreName, AGGREGATE_NAME)));
}

@SuppressWarnings("deprecation")
@Override
public <T> KTable<K, T> aggregate(final Initializer<T> initializer,
final Aggregator<? super K, ? super V, T> adder,
final Aggregator<? super K, ? super V, T> subtractor,
final Serde<T> aggValueSerde) {
return aggregate(initializer, adder, subtractor, aggValueSerde, null);
}

@SuppressWarnings("deprecation")
@Override
public <T> KTable<K, T> aggregate(final Initializer<T> initializer,
final Aggregator<? super K, ? super V, T> adder,
final Aggregator<? super K, ? super V, T> subtractor,
final String queryableStoreName) {
determineIsQueryable(queryableStoreName);
return aggregate(initializer, adder, subtractor, null, getOrCreateName(queryableStoreName, AGGREGATE_NAME));
}

@Override
public <T> KTable<K, T> aggregate(final Initializer<T> initializer,
final Aggregator<? super K, ? super V, T> adder,
final Aggregator<? super K, ? super V, T> subtractor) {
return aggregate(initializer, adder, subtractor, (String) null);
}

@SuppressWarnings("deprecation")
@Override
public <T> KTable<K, T> aggregate(final Initializer<T> initializer,
final Aggregator<? super K, ? super V, T> adder,
final Aggregator<? super K, ? super V, T> subtractor,
final org.apache.kafka.streams.processor.StateStoreSupplier<KeyValueStore> storeSupplier) {
Objects.requireNonNull(initializer, "initializer can't be null");
Objects.requireNonNull(adder, "adder can't be null");
Objects.requireNonNull(subtractor, "subtractor can't be null");
Objects.requireNonNull(storeSupplier, "storeSupplier can't be null");
ProcessorSupplier<K, Change<V>> aggregateSupplier = new KTableAggregate<>(storeSupplier.name(), initializer, adder, subtractor);
return doAggregate(aggregateSupplier, AGGREGATE_NAME, storeSupplier);
}

@SuppressWarnings("deprecation")
private <T> KTable<K, T> doAggregate(final ProcessorSupplier<K, Change<V>> aggregateSupplier,
final String functionName,
final org.apache.kafka.streams.processor.StateStoreSupplier<KeyValueStore> storeSupplier) {
final String sinkName = builder.newProcessorName(KStreamImpl.SINK_NAME);
final String sourceName = builder.newProcessorName(KStreamImpl.SOURCE_NAME);
final String funcName = builder.newProcessorName(functionName);

buildAggregate(aggregateSupplier,
storeSupplier.name() + KStreamImpl.REPARTITION_TOPIC_SUFFIX,
funcName,
sourceName,
sinkName);
builder.internalTopologyBuilder.addStateStore(storeSupplier, funcName);

// return the KTable representation with the intermediate topic as the sources
return new KTableImpl<>(builder, funcName, aggregateSupplier, Collections.singleton(sourceName), storeSupplier.name(), isQueryable);
}

private void buildAggregate(final ProcessorSupplier<K, Change<V>> aggregateSupplier,
Expand Down Expand Up @@ -196,16 +118,7 @@ private <T> KTable<K, T> doAggregate(final ProcessorSupplier<K, Change<V>> aggre
.materialize(), funcName);

// return the KTable representation with the intermediate topic as the sources
return new KTableImpl<>(builder, funcName, aggregateSupplier, Collections.singleton(sourceName), materialized.storeName(), isQueryable);
}

@SuppressWarnings("deprecation")
@Override
public KTable<K, V> reduce(final Reducer<V> adder,
final Reducer<V> subtractor,
final String queryableStoreName) {
determineIsQueryable(queryableStoreName);
return reduce(adder, subtractor, keyValueStore(keySerde, valSerde, getOrCreateName(queryableStoreName, REDUCE_NAME)));
return new KTableImpl<>(builder, funcName, aggregateSupplier, Collections.singleton(sourceName), materialized.storeName(), materialized.isQueryable());
}

@Override
Expand All @@ -216,7 +129,13 @@ public KTable<K, V> reduce(final Reducer<V> adder,
Objects.requireNonNull(subtractor, "subtractor can't be null");
Objects.requireNonNull(materialized, "materialized can't be null");
final MaterializedInternal<K, V, KeyValueStore<Bytes, byte[]>> materializedInternal
= new MaterializedInternal<>(materialized, builder, REDUCE_NAME);
= new MaterializedInternal<>(materialized, builder, AGGREGATE_NAME);
if (materializedInternal.keySerde() == null) {
materializedInternal.withKeySerde(keySerde);
}
if (materializedInternal.valueSerde() == null) {
materializedInternal.withValueSerde(valSerde);
}
final ProcessorSupplier<K, Change<V>> aggregateSupplier = new KTableReduce<>(materializedInternal.storeName(),
adder,
subtractor);
Expand All @@ -226,37 +145,33 @@ public KTable<K, V> reduce(final Reducer<V> adder,
@Override
public KTable<K, V> reduce(final Reducer<V> adder,
final Reducer<V> subtractor) {
return reduce(adder, subtractor, (String) null);
return reduce(adder, subtractor, Materialized.<K, V, KeyValueStore<Bytes, byte[]>>with(keySerde, valSerde));
}

@SuppressWarnings("deprecation")
@Override
public KTable<K, V> reduce(final Reducer<V> adder,
final Reducer<V> subtractor,
final org.apache.kafka.streams.processor.StateStoreSupplier<KeyValueStore> storeSupplier) {
Objects.requireNonNull(adder, "adder can't be null");
Objects.requireNonNull(subtractor, "subtractor can't be null");
Objects.requireNonNull(storeSupplier, "storeSupplier can't be null");
ProcessorSupplier<K, Change<V>> aggregateSupplier = new KTableReduce<>(storeSupplier.name(), adder, subtractor);
return doAggregate(aggregateSupplier, REDUCE_NAME, storeSupplier);
}
public KTable<K, Long> count(final Materialized<K, Long, KeyValueStore<Bytes, byte[]>> materialized) {
final MaterializedInternal<K, Long, KeyValueStore<Bytes, byte[]>> materializedInternal
= new MaterializedInternal<>(materialized, builder, AGGREGATE_NAME);
if (materializedInternal.keySerde() == null) {
materializedInternal.withKeySerde(keySerde);
}
if (materializedInternal.valueSerde() == null) {
materializedInternal.withValueSerde(Serdes.Long());
}

@SuppressWarnings("deprecation")
@Override
public KTable<K, Long> count(final String queryableStoreName) {
determineIsQueryable(queryableStoreName);
return count(keyValueStore(keySerde, Serdes.Long(), getOrCreateName(queryableStoreName, AGGREGATE_NAME)));
final ProcessorSupplier<K, Change<V>> aggregateSupplier = new KTableAggregate<>(materializedInternal.storeName(),
countInitializer,
countAdder,
countSubtractor);

return doAggregate(aggregateSupplier, AGGREGATE_NAME, materializedInternal);
}

@Override
public KTable<K, Long> count(final Materialized<K, Long, KeyValueStore<Bytes, byte[]>> materialized) {
return aggregate(countInitializer,
countAdder,
countSubtractor,
materialized);
public KTable<K, Long> count() {
return count(Materialized.<K, Long, KeyValueStore<Bytes, byte[]>>with(keySerde, Serdes.Long()));
}

@SuppressWarnings("unchecked")
@Override
public <VR> KTable<K, VR> aggregate(final Initializer<VR> initializer,
final Aggregator<? super K, ? super V, VR> adder,
Expand All @@ -266,6 +181,7 @@ public <VR> KTable<K, VR> aggregate(final Initializer<VR> initializer,
Objects.requireNonNull(adder, "adder can't be null");
Objects.requireNonNull(subtractor, "subtractor can't be null");
Objects.requireNonNull(materialized, "materialized can't be null");

final MaterializedInternal<K, VR, KeyValueStore<Bytes, byte[]>> materializedInternal =
new MaterializedInternal<>(materialized, builder, AGGREGATE_NAME);
if (materializedInternal.keySerde() == null) {
Expand All @@ -279,18 +195,10 @@ public <VR> KTable<K, VR> aggregate(final Initializer<VR> initializer,
}

@Override
public KTable<K, Long> count() {
return count((String) null);
}

@SuppressWarnings("deprecation")
@Override
public KTable<K, Long> count(final org.apache.kafka.streams.processor.StateStoreSupplier<KeyValueStore> storeSupplier) {
return this.aggregate(
countInitializer,
countAdder,
countSubtractor,
storeSupplier);
public <T> KTable<K, T> aggregate(final Initializer<T> initializer,
final Aggregator<? super K, ? super V, T> adder,
final Aggregator<? super K, ? super V, T> subtractor) {
return aggregate(initializer, adder, subtractor, Materialized.<K, T, KeyValueStore<Bytes, byte[]>>with(keySerde, null));
}

}
Loading