Skip to content
Closed
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
Original file line number Diff line number Diff line change
Expand Up @@ -17,20 +17,20 @@

package org.apache.kafka.streams.examples;

import org.apache.kafka.common.serialization.IntegerSerializer;
import org.apache.kafka.common.serialization.StringSerializer;
import org.apache.kafka.common.serialization.IntegerDeserializer;
import org.apache.kafka.common.serialization.IntegerSerializer;
import org.apache.kafka.common.serialization.StringDeserializer;
import org.apache.kafka.common.serialization.StringSerializer;
import org.apache.kafka.streams.KafkaStreaming;
import org.apache.kafka.streams.StreamingConfig;
import org.apache.kafka.streams.processor.Processor;
import org.apache.kafka.streams.processor.ProcessorContext;
import org.apache.kafka.streams.processor.ProcessorSupplier;
import org.apache.kafka.streams.processor.TopologyBuilder;
import org.apache.kafka.streams.processor.ProcessorContext;
import org.apache.kafka.streams.StreamingConfig;
import org.apache.kafka.streams.state.Entry;
import org.apache.kafka.streams.state.InMemoryKeyValueStore;
import org.apache.kafka.streams.state.KeyValueIterator;
import org.apache.kafka.streams.state.KeyValueStore;
import org.apache.kafka.streams.state.Stores;

import java.util.Properties;

Expand All @@ -48,7 +48,7 @@ public Processor<String, String> get() {
public void init(ProcessorContext context) {
this.context = context;
this.context.schedule(1000);
this.kvStore = new InMemoryKeyValueStore<>("local-state", context);
this.kvStore = Stores.create("local-state", context).withStringKeys().withIntegerValues().inMemory().build();
}

@Override
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -26,7 +26,6 @@
import org.apache.kafka.streams.kstream.internals.FilteredIterator;
import org.apache.kafka.streams.kstream.internals.WindowSupport;
import org.apache.kafka.streams.processor.StateRestoreCallback;
import org.apache.kafka.streams.processor.internals.ProcessorContextImpl;
import org.apache.kafka.streams.processor.internals.RecordCollector;
import org.apache.kafka.streams.processor.ProcessorContext;
import org.apache.kafka.streams.processor.internals.Stamped;
Expand Down Expand Up @@ -186,7 +185,7 @@ public void flush() {
IntegerSerializer intSerializer = new IntegerSerializer();
ByteArraySerializer byteArraySerializer = new ByteArraySerializer();

RecordCollector collector = ((ProcessorContextImpl) context).recordCollector();
RecordCollector collector = ((RecordCollector.Supplier) context).recordCollector();

for (Map.Entry<K, ValueList<V>> entry : map.entrySet()) {
ValueList<V> values = entry.getValue();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -37,7 +37,7 @@
import java.util.Map;
import java.util.Set;

public class ProcessorContextImpl implements ProcessorContext {
public class ProcessorContextImpl implements ProcessorContext, RecordCollector.Supplier {

private static final Logger log = LoggerFactory.getLogger(ProcessorContextImpl.class);

Expand Down Expand Up @@ -75,6 +75,7 @@ public ProcessorContextImpl(int id,
this.initialized = false;
}

@Override
public RecordCollector recordCollector() {
return this.collector;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,17 @@

public class RecordCollector {

/**
* A supplier of a {@link RecordCollector} instance.
*/
public static interface Supplier {
/**
* Get the record collector.
* @return the record collector
*/
public RecordCollector recordCollector();
}

private static final Logger log = LoggerFactory.getLogger(RecordCollector.class);

private final Producer<byte[], byte[]> producer;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -53,7 +53,7 @@ public void init(ProcessorContext context) {
@Override
public void process(K key, V value) {
// send to all the registered topics
RecordCollector collector = ((ProcessorContextImpl) context).recordCollector();
RecordCollector collector = ((RecordCollector.Supplier) context).recordCollector();
collector.send(new ProducerRecord<>(topic, key, value), keySerializer, valSerializer);
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,9 +17,9 @@

package org.apache.kafka.streams.state;

import org.apache.kafka.streams.processor.ProcessorContext;
import org.apache.kafka.common.utils.SystemTime;
import org.apache.kafka.common.utils.Time;
import org.apache.kafka.streams.processor.ProcessorContext;

import java.util.Iterator;
import java.util.List;
Expand All @@ -28,33 +28,28 @@
import java.util.TreeMap;

/**
* An in-memory key-value store based on a TreeMap
* An in-memory key-value store based on a TreeMap.
*
* @param <K> The key type
* @param <V> The value type
*
* @see Stores#create(String, ProcessorContext)
*/
public class InMemoryKeyValueStore<K, V> extends MeteredKeyValueStore<K, V> {

public InMemoryKeyValueStore(String name, ProcessorContext context) {
this(name, context, new SystemTime());
}

public InMemoryKeyValueStore(String name, ProcessorContext context, Time time) {
super(name, new MemoryStore<K, V>(name, context), context, "in-memory-state", time);
protected InMemoryKeyValueStore(String name, ProcessorContext context, Serdes<K, V> serdes, Time time) {
super(name, new MemoryStore<K, V>(name), context, serdes, "in-memory-state", time != null ? time : new SystemTime());
}

private static class MemoryStore<K, V> implements KeyValueStore<K, V> {

private final String name;
private final NavigableMap<K, V> map;
private final ProcessorContext context;

@SuppressWarnings("unchecked")
public MemoryStore(String name, ProcessorContext context) {
public MemoryStore(String name) {
super();
this.name = name;
this.map = new TreeMap<>();
this.context = context;
}

@Override
Expand Down Expand Up @@ -137,4 +132,4 @@ public void close() {

}
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -25,7 +25,6 @@
import org.apache.kafka.common.serialization.Deserializer;
import org.apache.kafka.common.serialization.Serializer;
import org.apache.kafka.common.utils.Time;
import org.apache.kafka.streams.processor.internals.ProcessorContextImpl;
import org.apache.kafka.streams.processor.internals.RecordCollector;

import java.util.HashSet;
Expand All @@ -35,6 +34,7 @@
public class MeteredKeyValueStore<K, V> implements KeyValueStore<K, V> {

protected final KeyValueStore<K, V> inner;
protected final Serdes<K, V> serialization;

private final Time time;
private final Sensor putTime;
Expand All @@ -54,8 +54,10 @@ public class MeteredKeyValueStore<K, V> implements KeyValueStore<K, V> {
private final ProcessorContext context;

// always wrap the logged store with the metered store
public MeteredKeyValueStore(final String name, final KeyValueStore<K, V> inner, ProcessorContext context, String metricGrp, Time time) {
public MeteredKeyValueStore(final String name, final KeyValueStore<K, V> inner, ProcessorContext context,
Serdes<K, V> serialization, String metricGrp, Time time) {
this.inner = inner;
this.serialization = serialization;

this.time = time;
this.metrics = context.metrics();
Expand All @@ -79,8 +81,8 @@ public MeteredKeyValueStore(final String name, final KeyValueStore<K, V> inner,
// register and possibly restore the state from the logs
long startNs = time.nanoseconds();
try {
final Deserializer<K> keyDeserializer = (Deserializer<K>) context.keyDeserializer();
final Deserializer<V> valDeserializer = (Deserializer<V>) context.valueDeserializer();
final Deserializer<K> keyDeserializer = serialization.keyDeserializer();
final Deserializer<V> valDeserializer = serialization.valueDeserializer();

context.register(this, new StateRestoreCallback() {
@Override
Expand Down Expand Up @@ -188,11 +190,11 @@ public void flush() {
}

private void logChange() {
RecordCollector collector = ((ProcessorContextImpl) context).recordCollector();
Serializer<K> keySerializer = (Serializer<K>) context.keySerializer();
Serializer<V> valueSerializer = (Serializer<V>) context.valueSerializer();

RecordCollector collector = ((RecordCollector.Supplier) context).recordCollector();
if (collector != null) {
Serializer<K> keySerializer = serialization.keySerializer();
Serializer<V> valueSerializer = serialization.valueSerializer();

for (K k : this.dirty) {
V v = this.inner.get(k);
collector.send(new ProducerRecord<>(this.topic, this.partition, k, v), keySerializer, valueSerializer);
Expand Down Expand Up @@ -239,4 +241,4 @@ public void close() {

}

}
}
Loading