Skip to content
Merged
Show file tree
Hide file tree
Changes from 1 commit
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 @@ -20,6 +20,7 @@
import org.apache.kafka.common.serialization.Deserializer;
import org.apache.kafka.common.serialization.Serializer;
import org.apache.kafka.streaming.processor.Processor;
import org.apache.kafka.streaming.processor.ProcessorFactory;

/**
* KStream is an abstraction of a stream of key-value pairs.
Expand Down Expand Up @@ -133,7 +134,7 @@ public interface KStream<K, V> {
/**
* Processes all elements in this stream by applying a processor.
*
* @param processor the class of Processor
* @param processorFactory the class of ProcessorFactory
*/
<K1, V1> KStream<K1, V1> process(Processor<K, V> processor);
<K1, V1> KStream<K1, V1> process(ProcessorFactory processorFactory);
}
Original file line number Diff line number Diff line change
Expand Up @@ -36,20 +36,14 @@
import java.util.LinkedList;
import java.util.Map;

public class SlidingWindow<K, V> extends WindowSupport implements Window<K, V> {

private final Object lock = new Object();
public class SlidingWindow<K, V> implements Window<K, V> {
private String name;
private final long duration;
private final int maxCount;
private final Serializer<K> keySerializer;
private final Serializer<V> valueSerializer;
private final Deserializer<K> keyDeserializer;
private final Deserializer<V> valueDeserializer;
private ProcessorContext context;
private int slotNum;
private String name;
private final long duration;
private final int maxCount;
private LinkedList<K> list = new LinkedList<K>();
private HashMap<K, ValueList<V>> map = new HashMap<>();

public SlidingWindow(
String name,
Expand All @@ -69,182 +63,201 @@ public SlidingWindow(
}

@Override
public void init(ProcessorContext context) {
this.context = context;
RestoreFuncImpl restoreFunc = new RestoreFuncImpl();
context.register(this, restoreFunc);

for (ValueList<V> valueList : map.values()) {
valueList.clearDirtyValues();
}
this.slotNum = restoreFunc.slotNum;
public String name() {
return name;
}

@Override
public Iterator<V> findAfter(K key, final long timestamp) {
return find(key, timestamp, timestamp + duration);
public WindowInstance<K, V> build() {
return new SlidingWindowInstance();
}

@Override
public Iterator<V> findBefore(K key, final long timestamp) {
return find(key, timestamp - duration, timestamp);
}
public class SlidingWindowInstance extends WindowSupport implements WindowInstance<K, V> {
private final Object lock = new Object();
private ProcessorContext context;
private int slotNum; // used as a key for Kafka log compaction
private LinkedList<K> list = new LinkedList<K>();
private HashMap<K, ValueList<V>> map = new HashMap<>();

@Override
public Iterator<V> find(K key, final long timestamp) {
return find(key, timestamp - duration, timestamp + duration);
}
@Override
public void init(ProcessorContext context) {
this.context = context;
RestoreFuncImpl restoreFunc = new RestoreFuncImpl();
context.register(this, restoreFunc);

/*
* finds items in the window between startTime and endTime (both inclusive)
*/
private Iterator<V> find(K key, final long startTime, final long endTime) {
final ValueList<V> values = map.get(key);

if (values == null) {
return null;
} else {
return new FilteredIterator<V, Value<V>>(values.iterator()) {
@Override
protected V filter(Value<V> item) {
if (startTime <= item.timestamp && item.timestamp <= endTime)
return item.value;
else
return null;
}
};
for (ValueList<V> valueList : map.values()) {
valueList.clearDirtyValues();
}
this.slotNum = restoreFunc.slotNum;
}
}

@Override
public void put(K key, V value, long timestamp) {
synchronized (lock) {
slotNum++;
@Override
public Iterator<V> findAfter(K key, final long timestamp) {
return find(key, timestamp, timestamp + duration);
}

list.offerLast(key);
@Override
public Iterator<V> findBefore(K key, final long timestamp) {
return find(key, timestamp - duration, timestamp);
}

@Override
public Iterator<V> find(K key, final long timestamp) {
return find(key, timestamp - duration, timestamp + duration);
}

/*
* finds items in the window between startTime and endTime (both inclusive)
*/
private Iterator<V> find(K key, final long startTime, final long endTime) {
final ValueList<V> values = map.get(key);

ValueList<V> values = map.get(key);
if (values == null) {
values = new ValueList<>();
map.put(key, values);
return null;
} else {
return new FilteredIterator<V, Value<V>>(values.iterator()) {
@Override
protected V filter(Value<V> item) {
if (startTime <= item.timestamp && item.timestamp <= endTime)
return item.value;
else
return null;
}
};
}

values.add(slotNum, value, timestamp);
}
evictExcess();
evictExpired(timestamp - duration);
}

private void evictExcess() {
while (list.size() > maxCount) {
K oldestKey = list.pollFirst();
@Override
public void put(K key, V value, long timestamp) {
synchronized (lock) {
slotNum++;

ValueList<V> values = map.get(oldestKey);
values.removeFirst();
list.offerLast(key);

if (values.isEmpty()) map.remove(oldestKey);
}
}
ValueList<V> values = map.get(key);
if (values == null) {
values = new ValueList<>();
map.put(key, values);
}

private void evictExpired(long cutoffTime) {
while (true) {
K oldestKey = list.peekFirst();
values.add(slotNum, value, timestamp);
}
evictExcess();
evictExpired(timestamp - duration);
}

ValueList<V> values = map.get(oldestKey);
Stamped<V> oldestValue = values.first();
private void evictExcess() {
while (list.size() > maxCount) {
K oldestKey = list.pollFirst();

if (oldestValue.timestamp < cutoffTime) {
list.pollFirst();
ValueList<V> values = map.get(oldestKey);
values.removeFirst();

if (values.isEmpty()) map.remove(oldestKey);
} else {
break;
}
}
}

@Override
public String name() {
return name;
}
private void evictExpired(long cutoffTime) {
while (true) {
K oldestKey = list.peekFirst();

@Override
public void flush() {
IntegerSerializer intSerializer = new IntegerSerializer();
ByteArraySerializer byteArraySerializer = new ByteArraySerializer();
ValueList<V> values = map.get(oldestKey);
Stamped<V> oldestValue = values.first();

if (oldestValue.timestamp < cutoffTime) {
list.pollFirst();
values.removeFirst();

RecordCollector collector = ((ProcessorContextImpl)context).recordCollector();
if (values.isEmpty()) map.remove(oldestKey);
} else {
break;
}
}
}

for (Map.Entry<K, ValueList<V>> entry : map.entrySet()) {
ValueList<V> values = entry.getValue();
if (values.hasDirtyValues()) {
K key = entry.getKey();
@Override
public String name() {
return name;
}

byte[] keyBytes = keySerializer.serialize(name, key);
@Override
public void flush() {
IntegerSerializer intSerializer = new IntegerSerializer();
ByteArraySerializer byteArraySerializer = new ByteArraySerializer();

Iterator<Value<V>> iterator = values.dirtyValueIterator();
while (iterator.hasNext()) {
Value<V> dirtyValue = iterator.next();
byte[] slot = intSerializer.serialize("", dirtyValue.slotNum);
byte[] valBytes = valueSerializer.serialize(name, dirtyValue.value);
RecordCollector collector = ((ProcessorContextImpl) context).recordCollector();

byte[] combined = new byte[8 + 4 + keyBytes.length + 4 + valBytes.length];
for (Map.Entry<K, ValueList<V>> entry : map.entrySet()) {
ValueList<V> values = entry.getValue();
if (values.hasDirtyValues()) {
K key = entry.getKey();

int offset = 0;
offset += putLong(combined, offset, dirtyValue.timestamp);
offset += puts(combined, offset, keyBytes);
offset += puts(combined, offset, valBytes);
byte[] keyBytes = keySerializer.serialize(name, key);

if (offset != combined.length) throw new IllegalStateException("serialized length does not match");
Iterator<Value<V>> iterator = values.dirtyValueIterator();
while (iterator.hasNext()) {
Value<V> dirtyValue = iterator.next();
byte[] slot = intSerializer.serialize("", dirtyValue.slotNum);
byte[] valBytes = valueSerializer.serialize(name, dirtyValue.value);

collector.send(new ProducerRecord<>(name, context.id(), slot, combined), byteArraySerializer, byteArraySerializer);
byte[] combined = new byte[8 + 4 + keyBytes.length + 4 + valBytes.length];

int offset = 0;
offset += putLong(combined, offset, dirtyValue.timestamp);
offset += puts(combined, offset, keyBytes);
offset += puts(combined, offset, valBytes);

if (offset != combined.length)
throw new IllegalStateException("serialized length does not match");

collector.send(new ProducerRecord<>(name, context.id(), slot, combined), byteArraySerializer, byteArraySerializer);
}
values.clearDirtyValues();
}
values.clearDirtyValues();
}
}
}

@Override
public void close() {
// TODO
}
@Override
public void close() {
// TODO
}

@Override
public boolean persistent() {
// TODO: should not be persistent, right?
return false;
}
@Override
public boolean persistent() {
// TODO: should not be persistent, right?
return false;
}

private class RestoreFuncImpl implements RestoreFunc {
private class RestoreFuncImpl implements RestoreFunc {

final IntegerDeserializer intDeserializer;
int slotNum = 0;
final IntegerDeserializer intDeserializer;
int slotNum = 0;

RestoreFuncImpl() {
intDeserializer = new IntegerDeserializer();
}
RestoreFuncImpl() {
intDeserializer = new IntegerDeserializer();
}

@Override
public void apply(byte[] slot, byte[] bytes) {

slotNum = intDeserializer.deserialize("", slot);

int offset = 0;
// timestamp
long timestamp = getLong(bytes, offset);
offset += 8;
// key
int length = getInt(bytes, offset);
offset += 4;
K key = deserialize(bytes, offset, length, name, keyDeserializer);
offset += length;
// value
length = getInt(bytes, offset);
offset += 4;
V value = deserialize(bytes, offset, length, name, valueDeserializer);

put(key, value, timestamp);
@Override
public void apply(byte[] slot, byte[] bytes) {

slotNum = intDeserializer.deserialize("", slot);

int offset = 0;
// timestamp
long timestamp = getLong(bytes, offset);
offset += 8;
// key
int length = getInt(bytes, offset);
offset += 4;
K key = deserialize(bytes, offset, length, name, keyDeserializer);
offset += length;
// value
length = getInt(bytes, offset);
offset += 4;
V value = deserialize(bytes, offset, length, name, valueDeserializer);

put(key, value, timestamp);
}
}
}

Expand Down
Loading