Skip to content
Merged
Show file tree
Hide file tree
Changes from 4 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
Original file line number Diff line number Diff line change
Expand Up @@ -25,8 +25,8 @@
import org.apache.kafka.streams.processor.api.Record;
import org.apache.kafka.streams.processor.api.RecordMetadata;
import org.apache.kafka.streams.processor.internals.metrics.StreamsMetricsImpl;
import org.apache.kafka.streams.state.TimestampedKeyValueStore;
import org.apache.kafka.streams.state.ValueAndTimestamp;
import org.apache.kafka.streams.state.internals.KeyValueStoreWrapper;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

Expand Down Expand Up @@ -63,7 +63,7 @@ public void enableSendingOldValues() {


private class KStreamAggregateProcessor extends ContextualProcessor<KIn, VIn, KIn, Change<VAgg>> {
private TimestampedKeyValueStore<KIn, VAgg> store;
private KeyValueStoreWrapper<KIn, VAgg> store;
private Sensor droppedRecordsSensor;
private TimestampedTupleForwarder<KIn, VAgg> tupleForwarder;

Expand All @@ -74,9 +74,9 @@ public void init(final ProcessorContext<KIn, Change<VAgg>> context) {
Thread.currentThread().getName(),
context.taskId().toString(),
(StreamsMetricsImpl) context.metrics());
store = context.getStateStore(storeName);
store = new KeyValueStoreWrapper<>(context, storeName);
tupleForwarder = new TimestampedTupleForwarder<>(
store,
store.getStore(),
context,
new TimestampedCacheFlushListener<>(context),
sendOldValues);
Expand Down Expand Up @@ -118,7 +118,7 @@ public void process(final Record<KIn, VIn> record) {

newAgg = aggregator.apply(record.key(), record.value(), oldAgg);

store.put(record.key(), ValueAndTimestamp.make(newAgg, newTimestamp));
store.put(record.key(), newAgg, newTimestamp);
tupleForwarder.maybeForward(
record.withValue(new Change<>(newAgg, sendOldValues ? oldAgg : null))
.withTimestamp(newTimestamp));
Expand All @@ -141,11 +141,11 @@ public String[] storeNames() {
}

private class KStreamAggregateValueGetter implements KTableValueGetter<KIn, VAgg> {
private TimestampedKeyValueStore<KIn, VAgg> store;
private KeyValueStoreWrapper<KIn, VAgg> store;

@Override
public void init(final ProcessorContext<?, ?> context) {
store = context.getStateStore(storeName);
store = new KeyValueStoreWrapper<>(context, storeName);
}

@Override
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -24,8 +24,8 @@
import org.apache.kafka.streams.processor.api.Record;
import org.apache.kafka.streams.processor.api.RecordMetadata;
import org.apache.kafka.streams.processor.internals.metrics.StreamsMetricsImpl;
import org.apache.kafka.streams.state.TimestampedKeyValueStore;
import org.apache.kafka.streams.state.ValueAndTimestamp;
import org.apache.kafka.streams.state.internals.KeyValueStoreWrapper;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

Expand Down Expand Up @@ -58,7 +58,7 @@ public void enableSendingOldValues() {


private class KStreamReduceProcessor extends ContextualProcessor<K, V, K, Change<V>> {
private TimestampedKeyValueStore<K, V> store;
private KeyValueStoreWrapper<K, V> store;
private TimestampedTupleForwarder<K, V> tupleForwarder;
private Sensor droppedRecordsSensor;

Expand All @@ -70,9 +70,9 @@ public void init(final ProcessorContext<K, Change<V>> context) {
context.taskId().toString(),
(StreamsMetricsImpl) context.metrics()
);
store = context.getStateStore(storeName);
store = new KeyValueStoreWrapper<>(context, storeName);
tupleForwarder = new TimestampedTupleForwarder<>(
store,
store.getStore(),
context,
new TimestampedCacheFlushListener<>(context),
sendOldValues);
Expand Down Expand Up @@ -112,7 +112,7 @@ public void process(final Record<K, V> record) {
newTimestamp = Math.max(record.timestamp(), oldAggAndTimestamp.timestamp());
}

store.put(record.key(), ValueAndTimestamp.make(newAgg, newTimestamp));
store.put(record.key(), newAgg, newTimestamp);
tupleForwarder.maybeForward(
record.withValue(new Change<>(newAgg, sendOldValues ? oldAgg : null))
.withTimestamp(newTimestamp));
Expand All @@ -136,11 +136,11 @@ public String[] storeNames() {


private class KStreamReduceValueGetter implements KTableValueGetter<K, V> {
private TimestampedKeyValueStore<K, V> store;
private KeyValueStoreWrapper<K, V> store;

@Override
public void init(final ProcessorContext<?, ?> context) {
store = context.getStateStore(storeName);
store = new KeyValueStoreWrapper<>(context, storeName);
}

@Override
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -22,8 +22,8 @@
import org.apache.kafka.streams.processor.api.Processor;
import org.apache.kafka.streams.processor.api.ProcessorContext;
import org.apache.kafka.streams.processor.api.Record;
import org.apache.kafka.streams.state.TimestampedKeyValueStore;
import org.apache.kafka.streams.state.ValueAndTimestamp;
import org.apache.kafka.streams.state.internals.KeyValueStoreWrapper;

import static org.apache.kafka.streams.state.ValueAndTimestamp.getValueOrNull;

Expand Down Expand Up @@ -60,15 +60,15 @@ public Processor<KIn, Change<VIn>, KIn, Change<VAgg>> get() {
}

private class KTableAggregateProcessor implements Processor<KIn, Change<VIn>, KIn, Change<VAgg>> {
private TimestampedKeyValueStore<KIn, VAgg> store;
private KeyValueStoreWrapper<KIn, VAgg> store;
private TimestampedTupleForwarder<KIn, VAgg> tupleForwarder;

@SuppressWarnings("unchecked")
@Override
public void init(final ProcessorContext<KIn, Change<VAgg>> context) {
store = (TimestampedKeyValueStore<KIn, VAgg>) context.getStateStore(storeName);
store = new KeyValueStoreWrapper<>(context, storeName);
tupleForwarder = new TimestampedTupleForwarder<>(
store,
store.getStore(),
context,
new TimestampedCacheFlushListener<>(context),
sendOldValues);
Expand Down Expand Up @@ -116,7 +116,7 @@ public void process(final Record<KIn, Change<VIn>> record) {
}

// update the store with the new value
store.put(record.key(), ValueAndTimestamp.make(newAgg, newTimestamp));
store.put(record.key(), newAgg, newTimestamp);
tupleForwarder.maybeForward(
record.withValue(new Change<>(newAgg, sendOldValues ? oldAgg : null))
.withTimestamp(newTimestamp));
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -20,8 +20,8 @@
import org.apache.kafka.streams.processor.api.Processor;
import org.apache.kafka.streams.processor.api.ProcessorContext;
import org.apache.kafka.streams.processor.api.Record;
import org.apache.kafka.streams.state.TimestampedKeyValueStore;
import org.apache.kafka.streams.state.ValueAndTimestamp;
import org.apache.kafka.streams.state.internals.KeyValueStoreWrapper;

import static org.apache.kafka.streams.state.ValueAndTimestamp.getValueOrNull;

Expand Down Expand Up @@ -88,16 +88,16 @@ private ValueAndTimestamp<VIn> computeValue(final KIn key, final ValueAndTimesta

private class KTableFilterProcessor implements Processor<KIn, Change<VIn>, KIn, Change<VIn>> {
private ProcessorContext<KIn, Change<VIn>> context;
private TimestampedKeyValueStore<KIn, VIn> store;
private KeyValueStoreWrapper<KIn, VIn> store;
private TimestampedTupleForwarder<KIn, VIn> tupleForwarder;

@Override
public void init(final ProcessorContext<KIn, Change<VIn>> context) {
this.context = context;
if (queryableName != null) {
store = context.getStateStore(queryableName);
store = new KeyValueStoreWrapper<>(context, queryableName);
tupleForwarder = new TimestampedTupleForwarder<>(
store,
store.getStore(),
context,
new TimestampedCacheFlushListener<>(context),
sendOldValues);
Expand All @@ -117,7 +117,7 @@ public void process(final Record<KIn, Change<VIn>> record) {
}

if (queryableName != null) {
store.put(key, ValueAndTimestamp.make(newValue, record.timestamp()));
store.put(key, newValue, record.timestamp());
tupleForwarder.maybeForward(record.withValue(new Change<>(newValue, oldValue)));
} else {
context.forward(record.withValue(new Change<>(newValue, oldValue)));
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -175,7 +175,7 @@ private KTable<K, V> doFilter(final Predicate<? super K, ? super V> predicate,
final Serde<K> keySerde;
final Serde<V> valueSerde;
final String queryableStoreName;
final StoreBuilder<TimestampedKeyValueStore<K, V>> storeBuilder;
final StoreBuilder<?> storeBuilder;

if (materializedInternal != null) {
// we actually do not need to generate store names at all since if it is not specified, we will not
Expand Down Expand Up @@ -290,7 +290,7 @@ private <VR> KTable<K, VR> doMapValues(final ValueMapperWithKey<? super K, ? sup
final Serde<K> keySerde;
final Serde<VR> valueSerde;
final String queryableStoreName;
final StoreBuilder<TimestampedKeyValueStore<K, VR>> storeBuilder;
final StoreBuilder<?> storeBuilder;

if (materializedInternal != null) {
// we actually do not need to generate store names at all since if it is not specified, we will not
Expand Down Expand Up @@ -445,7 +445,7 @@ private <VR> KTable<K, VR> doTransformValues(final ValueTransformerWithKeySuppli
final Serde<K> keySerde;
final Serde<VR> valueSerde;
final String queryableStoreName;
final StoreBuilder<TimestampedKeyValueStore<K, VR>> storeBuilder;
final StoreBuilder<?> storeBuilder;

if (materializedInternal != null) {
// don't inherit parent value serde, since this operation may change the value type, more specifically:
Expand Down Expand Up @@ -750,7 +750,7 @@ private <VO, VR> KTable<K, VR> doJoin(final KTable<K, VO> other,
final Serde<K> keySerde;
final Serde<VR> valueSerde;
final String queryableStoreName;
final StoreBuilder<TimestampedKeyValueStore<K, VR>> storeBuilder;
final StoreBuilder<?> storeBuilder;

if (materializedInternal != null) {
if (materializedInternal.keySerde() == null) {
Expand Down Expand Up @@ -1270,7 +1270,7 @@ private <VR, KO, VO> KTable<K, VR> doJoinOnForeignKey(final KTable<KO, VO> forei
materializedInternal.queryableStoreName()
);

final StoreBuilder<TimestampedKeyValueStore<K, VR>> resultStore =
final StoreBuilder<?> resultStore =
new TimestampedKeyValueStoreMaterializer<>(materializedInternal).materialize();

final TableProcessorNode<K, VR> resultNode = new TableProcessorNode<>(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -20,12 +20,11 @@
import org.apache.kafka.streams.processor.api.Processor;
import org.apache.kafka.streams.processor.api.ProcessorContext;
import org.apache.kafka.streams.processor.api.Record;
import org.apache.kafka.streams.state.TimestampedKeyValueStore;
import org.apache.kafka.streams.state.ValueAndTimestamp;

import java.util.Collections;
import java.util.HashSet;
import java.util.Set;
import org.apache.kafka.streams.state.internals.KeyValueStoreWrapper;

public class KTableKTableJoinMerger<K, V> implements KTableProcessorSupplier<K, V, K, V> {

Expand Down Expand Up @@ -98,17 +97,17 @@ public static <K, V> KTableKTableJoinMerger<K, V> of(final KTableProcessorSuppli
}

private class KTableKTableJoinMergeProcessor extends ContextualProcessor<K, Change<V>, K, Change<V>> {
private TimestampedKeyValueStore<K, V> store;
private KeyValueStoreWrapper<K, V> store;
private TimestampedTupleForwarder<K, V> tupleForwarder;

@SuppressWarnings("unchecked")
@Override
public void init(final ProcessorContext<K, Change<V>> context) {
super.init(context);
if (queryableName != null) {
store = (TimestampedKeyValueStore<K, V>) context.getStateStore(queryableName);
store = new KeyValueStoreWrapper<>(context, queryableName);
tupleForwarder = new TimestampedTupleForwarder<>(
store,
store.getStore(),
context,
new TimestampedCacheFlushListener<>(context),
sendOldValues);
Expand All @@ -118,7 +117,7 @@ public void init(final ProcessorContext<K, Change<V>> context) {
@Override
public void process(final Record<K, Change<V>> record) {
if (queryableName != null) {
store.put(record.key(), ValueAndTimestamp.make(record.value().newValue, record.timestamp()));
store.put(record.key(), record.value().newValue, record.timestamp());
tupleForwarder.maybeForward(record);
} else {
if (sendOldValues) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -20,8 +20,8 @@
import org.apache.kafka.streams.processor.api.Processor;
import org.apache.kafka.streams.processor.api.ProcessorContext;
import org.apache.kafka.streams.processor.api.Record;
import org.apache.kafka.streams.state.TimestampedKeyValueStore;
import org.apache.kafka.streams.state.ValueAndTimestamp;
import org.apache.kafka.streams.state.internals.KeyValueStoreWrapper;

import static org.apache.kafka.streams.state.ValueAndTimestamp.getValueOrNull;

Expand Down Expand Up @@ -106,16 +106,16 @@ private ValueAndTimestamp<VOut> computeValueAndTimestamp(final KIn key, final Va

private class KTableMapValuesProcessor implements Processor<KIn, Change<VIn>, KIn, Change<VOut>> {
private ProcessorContext<KIn, Change<VOut>> context;
private TimestampedKeyValueStore<KIn, VOut> store;
private KeyValueStoreWrapper<KIn, VOut> store;
private TimestampedTupleForwarder<KIn, VOut> tupleForwarder;

@Override
public void init(final ProcessorContext<KIn, Change<VOut>> context) {
this.context = context;
if (queryableName != null) {
store = context.getStateStore(queryableName);
store = new KeyValueStoreWrapper<>(context, queryableName);
tupleForwarder = new TimestampedTupleForwarder<>(
store,
store.getStore(),
context,
new TimestampedCacheFlushListener<>(context),
sendOldValues);
Expand All @@ -128,7 +128,7 @@ public void process(final Record<KIn, Change<VIn>> record) {
final VOut oldValue = computeOldValue(record.key(), record.value());

if (queryableName != null) {
store.put(record.key(), ValueAndTimestamp.make(newValue, record.timestamp()));
store.put(record.key(), newValue, record.timestamp());
tupleForwarder.maybeForward(record.withValue(new Change<>(newValue, oldValue)));
} else {
context.forward(record.withValue(new Change<>(newValue, oldValue)));
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,8 +17,8 @@
package org.apache.kafka.streams.kstream.internals;

import org.apache.kafka.streams.processor.api.ProcessorContext;
import org.apache.kafka.streams.state.TimestampedKeyValueStore;
import org.apache.kafka.streams.state.ValueAndTimestamp;
import org.apache.kafka.streams.state.internals.KeyValueStoreWrapper;

public class KTableMaterializedValueGetterSupplier<K, V> implements KTableValueGetterSupplier<K, V> {
private final String storeName;
Expand All @@ -37,11 +37,11 @@ public String[] storeNames() {
}

private class KTableMaterializedValueGetter implements KTableValueGetter<K, V> {
private TimestampedKeyValueStore<K, V> store;
private KeyValueStoreWrapper<K, V> store;

@Override
public void init(final ProcessorContext<?, ?> context) {
store = context.getStateStore(storeName);
store = new KeyValueStoreWrapper<>(context, storeName);
}

@Override
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,10 +19,10 @@
import org.apache.kafka.streams.processor.api.Processor;
import org.apache.kafka.streams.processor.api.ProcessorContext;
import org.apache.kafka.streams.processor.api.Record;
import org.apache.kafka.streams.state.TimestampedKeyValueStore;
import org.apache.kafka.streams.state.ValueAndTimestamp;

import java.util.Collection;
import org.apache.kafka.streams.state.internals.KeyValueStoreWrapper;

public class KTablePassThrough<KIn, VIn> implements KTableProcessorSupplier<KIn, VIn, KIn, VIn> {
private final Collection<KStreamAggProcessorSupplier> parents;
Expand Down Expand Up @@ -79,11 +79,11 @@ public void process(final Record<KIn, Change<VIn>> record) {
}

private class KTablePassThroughValueGetter implements KTableValueGetter<KIn, VIn> {
private TimestampedKeyValueStore<KIn, VIn> store;
private KeyValueStoreWrapper<KIn, VIn> store;

@Override
public void init(final ProcessorContext<?, ?> context) {
store = context.getStateStore(storeName);
store = new KeyValueStoreWrapper<>(context, storeName);
}

@Override
Expand Down
Loading