Skip to content
Merged
Show file tree
Hide file tree
Changes from all 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
44 changes: 21 additions & 23 deletions streams/src/main/java/org/apache/kafka/streams/StreamsBuilder.java
Original file line number Diff line number Diff line change
Expand Up @@ -224,9 +224,9 @@ public synchronized <K, V> KTable<K, V> table(final String topic,
Objects.requireNonNull(materialized, "materialized can't be null");
final ConsumedInternal<K, V> consumedInternal = new ConsumedInternal<>(consumed);
materialized.withKeySerde(consumedInternal.keySerde()).withValueSerde(consumedInternal.valueSerde());
return internalStreamsBuilder.table(topic,
consumedInternal,
new MaterializedInternal<>(materialized, internalStreamsBuilder, topic + "-"));
final MaterializedInternal<K, V, KeyValueStore<Bytes, byte[]>> materializedInternal = new MaterializedInternal<>(materialized);
materializedInternal.generateStoreNameIfNeeded(internalStreamsBuilder, topic + "-");
return internalStreamsBuilder.table(topic, consumedInternal, materializedInternal);
}

/**
Expand Down Expand Up @@ -273,12 +273,10 @@ public synchronized <K, V> KTable<K, V> table(final String topic,
Objects.requireNonNull(topic, "topic can't be null");
Objects.requireNonNull(consumed, "consumed can't be null");
final ConsumedInternal<K, V> consumedInternal = new ConsumedInternal<>(consumed);
return internalStreamsBuilder.table(topic,
consumedInternal,
new MaterializedInternal<>(
Materialized.<K, V, KeyValueStore<Bytes, byte[]>>with(consumedInternal.keySerde(), consumedInternal.valueSerde()),
internalStreamsBuilder,
topic + "-"));
final MaterializedInternal<K, V, KeyValueStore<Bytes, byte[]>> materializedInternal =
new MaterializedInternal<>(Materialized.with(consumedInternal.keySerde(), consumedInternal.valueSerde()));
materializedInternal.generateStoreNameIfNeeded(internalStreamsBuilder, topic + "-");
return internalStreamsBuilder.table(topic, consumedInternal, materializedInternal);
}

/**
Expand All @@ -302,8 +300,9 @@ public synchronized <K, V> KTable<K, V> table(final String topic,
final Materialized<K, V, KeyValueStore<Bytes, byte[]>> materialized) {
Objects.requireNonNull(topic, "topic can't be null");
Objects.requireNonNull(materialized, "materialized can't be null");
final MaterializedInternal<K, V, KeyValueStore<Bytes, byte[]>> materializedInternal
= new MaterializedInternal<>(materialized, internalStreamsBuilder, topic + "-");
final MaterializedInternal<K, V, KeyValueStore<Bytes, byte[]>> materializedInternal = new MaterializedInternal<>(materialized);
materializedInternal.generateStoreNameIfNeeded(internalStreamsBuilder, topic + "-");

return internalStreamsBuilder.table(topic,
new ConsumedInternal<>(Consumed.with(materializedInternal.keySerde(),
materializedInternal.valueSerde())),
Expand Down Expand Up @@ -331,14 +330,11 @@ public synchronized <K, V> GlobalKTable<K, V> globalTable(final String topic,
Objects.requireNonNull(topic, "topic can't be null");
Objects.requireNonNull(consumed, "consumed can't be null");
final ConsumedInternal<K, V> consumedInternal = new ConsumedInternal<>(consumed);
final MaterializedInternal<K, V, KeyValueStore<Bytes, byte[]>> materialized =
new MaterializedInternal<>(
Materialized.<K, V, KeyValueStore<Bytes, byte[]>>with(consumedInternal.keySerde(), consumedInternal.valueSerde()),
internalStreamsBuilder,
topic + "-");

final MaterializedInternal<K, V, KeyValueStore<Bytes, byte[]>> materializedInternal =
new MaterializedInternal<>(Materialized.<K, V, KeyValueStore<Bytes, byte[]>>with(consumedInternal.keySerde(), consumedInternal.valueSerde()));
materializedInternal.generateStoreNameIfNeeded(internalStreamsBuilder, topic + "-");

return internalStreamsBuilder.globalTable(topic, consumedInternal, materialized);
return internalStreamsBuilder.globalTable(topic, consumedInternal, materializedInternal);
}

/**
Expand Down Expand Up @@ -402,9 +398,10 @@ public synchronized <K, V> GlobalKTable<K, V> globalTable(final String topic,
final ConsumedInternal<K, V> consumedInternal = new ConsumedInternal<>(consumed);
// always use the serdes from consumed
materialized.withKeySerde(consumedInternal.keySerde()).withValueSerde(consumedInternal.valueSerde());
return internalStreamsBuilder.globalTable(topic,
consumedInternal,
new MaterializedInternal<>(materialized, internalStreamsBuilder, topic + "-"));
final MaterializedInternal<K, V, KeyValueStore<Bytes, byte[]>> materializedInternal = new MaterializedInternal<>(materialized);
materializedInternal.generateStoreNameIfNeeded(internalStreamsBuilder, topic + "-");

return internalStreamsBuilder.globalTable(topic, consumedInternal, materializedInternal);
}

/**
Expand Down Expand Up @@ -436,8 +433,9 @@ public synchronized <K, V> GlobalKTable<K, V> globalTable(final String topic,
final Materialized<K, V, KeyValueStore<Bytes, byte[]>> materialized) {
Objects.requireNonNull(topic, "topic can't be null");
Objects.requireNonNull(materialized, "materialized can't be null");
final MaterializedInternal<K, V, KeyValueStore<Bytes, byte[]>> materializedInternal =
new MaterializedInternal<>(materialized, internalStreamsBuilder, topic + "-");
final MaterializedInternal<K, V, KeyValueStore<Bytes, byte[]>> materializedInternal = new MaterializedInternal<>(materialized);
materializedInternal.generateStoreNameIfNeeded(internalStreamsBuilder, topic + "-");

return internalStreamsBuilder.globalTable(topic,
new ConsumedInternal<>(Consumed.with(materializedInternal.keySerde(),
materializedInternal.valueSerde())),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -74,8 +74,10 @@ public KTable<K, V> reduce(final Reducer<V> reducer,
final Materialized<K, V, KeyValueStore<Bytes, byte[]>> materialized) {
Objects.requireNonNull(reducer, "reducer 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);

final MaterializedInternal<K, V, KeyValueStore<Bytes, byte[]>> materializedInternal = new MaterializedInternal<>(materialized);
materializedInternal.generateStoreNameIfNeeded(builder, REDUCE_NAME);

if (materializedInternal.keySerde() == null) {
materializedInternal.withKeySerde(keySerde);
}
Expand All @@ -97,8 +99,9 @@ public <VR> KTable<K, VR> aggregate(final Initializer<VR> initializer,
Objects.requireNonNull(aggregator, "aggregator 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);
final MaterializedInternal<K, VR, KeyValueStore<Bytes, byte[]>> materializedInternal = new MaterializedInternal<>(materialized);
materializedInternal.generateStoreNameIfNeeded(builder, AGGREGATE_NAME);

if (materializedInternal.keySerde() == null) {
materializedInternal.withKeySerde(keySerde);
}
Expand All @@ -117,14 +120,26 @@ public <VR> KTable<K, VR> aggregate(final Initializer<VR> initializer,

@Override
public KTable<K, Long> count() {
return count(Materialized.<K, Long, KeyValueStore<Bytes, byte[]>>with(keySerde, Serdes.Long()));
return doCount(Materialized.with(keySerde, Serdes.Long()));
}

@Override
public KTable<K, Long> count(final Materialized<K, Long, KeyValueStore<Bytes, byte[]>> materialized) {
Objects.requireNonNull(materialized, "materialized can't be null");
final MaterializedInternal<K, Long, KeyValueStore<Bytes, byte[]>> materializedInternal
= new MaterializedInternal<>(materialized, builder, AGGREGATE_NAME);

// TODO: remove this when we do a topology-incompatible release
// we used to burn a topology name here, so we have to keep doing it for compatibility
if (new MaterializedInternal<>(materialized).storeName() == null) {
builder.newStoreName(AGGREGATE_NAME);
}

return doCount(materialized);
}

private KTable<K, Long> doCount(final Materialized<K, Long, KeyValueStore<Bytes, byte[]>> materialized) {
final MaterializedInternal<K, Long, KeyValueStore<Bytes, byte[]>> materializedInternal = new MaterializedInternal<>(materialized);
materializedInternal.generateStoreNameIfNeeded(builder, AGGREGATE_NAME);

if (materializedInternal.keySerde() == null) {
materializedInternal.withKeySerde(keySerde);
}
Expand All @@ -133,9 +148,9 @@ public KTable<K, Long> count(final Materialized<K, Long, KeyValueStore<Bytes, by
}

return doAggregate(
new KStreamAggregate<>(materializedInternal.storeName(), aggregateBuilder.countInitializer, aggregateBuilder.countAggregator),
AGGREGATE_NAME,
materializedInternal);
new KStreamAggregate<>(materializedInternal.storeName(), aggregateBuilder.countInitializer, aggregateBuilder.countAggregator),
AGGREGATE_NAME,
materializedInternal);
}

@Override
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -128,8 +128,9 @@ public KTable<K, V> reduce(final Reducer<V> adder,
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, V, KeyValueStore<Bytes, byte[]>> materializedInternal
= new MaterializedInternal<>(materialized, builder, AGGREGATE_NAME);
final MaterializedInternal<K, V, KeyValueStore<Bytes, byte[]>> materializedInternal = new MaterializedInternal<>(materialized);
materializedInternal.generateStoreNameIfNeeded(builder, AGGREGATE_NAME);

if (materializedInternal.keySerde() == null) {
materializedInternal.withKeySerde(keySerde);
}
Expand All @@ -150,8 +151,9 @@ public KTable<K, V> reduce(final Reducer<V> adder,

@Override
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);
final MaterializedInternal<K, Long, KeyValueStore<Bytes, byte[]>> materializedInternal = new MaterializedInternal<>(materialized);
materializedInternal.generateStoreNameIfNeeded(builder, AGGREGATE_NAME);

if (materializedInternal.keySerde() == null) {
materializedInternal.withKeySerde(keySerde);
}
Expand Down Expand Up @@ -182,8 +184,9 @@ public <VR> KTable<K, VR> aggregate(final Initializer<VR> initializer,
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);
final MaterializedInternal<K, VR, KeyValueStore<Bytes, byte[]>> materializedInternal = new MaterializedInternal<>(materialized);
materializedInternal.generateStoreNameIfNeeded(builder, AGGREGATE_NAME);

if (materializedInternal.keySerde() == null) {
materializedInternal.withKeySerde(keySerde);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -154,7 +154,10 @@ public KTable<K, V> filter(final Predicate<? super K, ? super V> predicate,
final Materialized<K, V, KeyValueStore<Bytes, byte[]>> materialized) {
Objects.requireNonNull(predicate, "predicate can't be null");
Objects.requireNonNull(materialized, "materialized can't be null");
return doFilter(predicate, new MaterializedInternal<>(materialized, builder, FILTER_NAME), false);
final MaterializedInternal<K, V, KeyValueStore<Bytes, byte[]>> materializedInternal = new MaterializedInternal<>(materialized);
materializedInternal.generateStoreNameIfNeeded(builder, FILTER_NAME);

return doFilter(predicate, materializedInternal, false);
}

@Override
Expand All @@ -168,7 +171,10 @@ public KTable<K, V> filterNot(final Predicate<? super K, ? super V> predicate,
final Materialized<K, V, KeyValueStore<Bytes, byte[]>> materialized) {
Objects.requireNonNull(predicate, "predicate can't be null");
Objects.requireNonNull(materialized, "materialized can't be null");
return doFilter(predicate, new MaterializedInternal<>(materialized, builder, FILTER_NAME), true);
final MaterializedInternal<K, V, KeyValueStore<Bytes, byte[]>> materializedInternal = new MaterializedInternal<>(materialized);
materializedInternal.generateStoreNameIfNeeded(builder, FILTER_NAME);

return doFilter(predicate, materializedInternal, true);
}

private <VR> KTable<K, VR> doMapValues(final ValueMapperWithKey<? super K, ? super V, ? extends VR> mapper,
Expand Down Expand Up @@ -210,7 +216,10 @@ public <VR> KTable<K, VR> mapValues(final ValueMapper<? super V, ? extends VR> m
Objects.requireNonNull(mapper, "mapper can't be null");
Objects.requireNonNull(materialized, "materialized can't be null");

return doMapValues(withKey(mapper), new MaterializedInternal<>(materialized, builder, MAPVALUES_NAME));
final MaterializedInternal<K, VR, KeyValueStore<Bytes, byte[]>> materializedInternal = new MaterializedInternal<>(materialized);
materializedInternal.generateStoreNameIfNeeded(builder, MAPVALUES_NAME);

return doMapValues(withKey(mapper), materializedInternal);
}

@Override
Expand All @@ -219,7 +228,10 @@ public <VR> KTable<K, VR> mapValues(final ValueMapperWithKey<? super K, ? super
Objects.requireNonNull(mapper, "mapper can't be null");
Objects.requireNonNull(materialized, "materialized can't be null");

return doMapValues(mapper, new MaterializedInternal<>(materialized, builder, MAPVALUES_NAME));
final MaterializedInternal<K, VR, KeyValueStore<Bytes, byte[]>> materializedInternal = new MaterializedInternal<>(materialized);
materializedInternal.generateStoreNameIfNeeded(builder, MAPVALUES_NAME);

return doMapValues(mapper, materializedInternal);
}

@Override
Expand All @@ -233,8 +245,10 @@ public <VR> KTable<K, VR> transformValues(final ValueTransformerWithKeySupplier<
final Materialized<K, VR, KeyValueStore<Bytes, byte[]>> materialized,
final String... stateStoreNames) {
Objects.requireNonNull(materialized, "materialized can't be null");
return doTransformValues(transformerSupplier,
new MaterializedInternal<>(materialized, builder, TRANSFORMVALUES_NAME), stateStoreNames);
final MaterializedInternal<K, VR, KeyValueStore<Bytes, byte[]>> materializedInternal = new MaterializedInternal<>(materialized);
materializedInternal.generateStoreNameIfNeeded(builder, TRANSFORMVALUES_NAME);

return doTransformValues(transformerSupplier, materializedInternal, stateStoreNames);
}

private <VR> KTable<K, VR> doTransformValues(final ValueTransformerWithKeySupplier<? super K, ? super V, ? extends VR> transformerSupplier,
Expand Down Expand Up @@ -304,7 +318,10 @@ public <VO, VR> KTable<K, VR> join(final KTable<K, VO> other,
Objects.requireNonNull(other, "other can't be null");
Objects.requireNonNull(joiner, "joiner can't be null");
Objects.requireNonNull(materialized, "materialized can't be null");
return doJoin(other, joiner, new MaterializedInternal<>(materialized, builder, MERGE_NAME), false, false);
final MaterializedInternal<K, VR, KeyValueStore<Bytes, byte[]>> materializedInternal = new MaterializedInternal<>(materialized);
materializedInternal.generateStoreNameIfNeeded(builder, MERGE_NAME);

return doJoin(other, joiner, materializedInternal, false, false);
}

@Override
Expand All @@ -317,7 +334,10 @@ public <V1, R> KTable<K, R> outerJoin(final KTable<K, V1> other,
public <VO, VR> KTable<K, VR> outerJoin(final KTable<K, VO> other,
final ValueJoiner<? super V, ? super VO, ? extends VR> joiner,
final Materialized<K, VR, KeyValueStore<Bytes, byte[]>> materialized) {
return doJoin(other, joiner, new MaterializedInternal<>(materialized, builder, MERGE_NAME), true, true);
final MaterializedInternal<K, VR, KeyValueStore<Bytes, byte[]>> materializedInternal = new MaterializedInternal<>(materialized);
materializedInternal.generateStoreNameIfNeeded(builder, MERGE_NAME);

return doJoin(other, joiner, materializedInternal, true, true);
}

@Override
Expand All @@ -330,11 +350,9 @@ public <V1, R> KTable<K, R> leftJoin(final KTable<K, V1> other,
public <VO, VR> KTable<K, VR> leftJoin(final KTable<K, VO> other,
final ValueJoiner<? super V, ? super VO, ? extends VR> joiner,
final Materialized<K, VR, KeyValueStore<Bytes, byte[]>> materialized) {
return doJoin(other,
joiner,
new MaterializedInternal<>(materialized, builder, MERGE_NAME),
true,
false);
final MaterializedInternal<K, VR, KeyValueStore<Bytes, byte[]>> materializedInternal = new MaterializedInternal<>(materialized);
materializedInternal.generateStoreNameIfNeeded(builder, MERGE_NAME);
return doJoin(other, joiner, materializedInternal, true, false);
}

@SuppressWarnings("unchecked")
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -25,18 +25,17 @@

public class MaterializedInternal<K, V, S extends StateStore> extends Materialized<K, V, S> {

private final boolean queryable;
private final boolean queriable;


public MaterializedInternal(final Materialized<K, V, S> materialized,
final InternalNameProvider nameProvider,
final String generatedStorePrefix) {
public MaterializedInternal(final Materialized<K, V, S> materialized) {
super(materialized);
queriable = storeName() != null;
}

public void generateStoreNameIfNeeded(final InternalNameProvider nameProvider,
final String generatedStorePrefix) {
if (storeName() == null) {
queryable = false;
storeName = nameProvider.newStoreName(generatedStorePrefix);
} else {
queryable = true;
}
}

Expand All @@ -63,7 +62,7 @@ public boolean loggingEnabled() {
return loggingEnabled;
}

public Map<String, String> logConfig() {
Map<String, String> logConfig() {
return topicConfig;
}

Expand All @@ -72,6 +71,6 @@ boolean cachingEnabled() {
}

boolean isQueryable() {
return queryable;
return queriable;
}
}
Loading