http://git-wip-us.apache.org/repos/asf/storm/blob/3ba3cabb/storm-client/src/jvm/org/apache/storm/windowing/CountTriggerPolicy.java ---------------------------------------------------------------------- diff --git a/storm-client/src/jvm/org/apache/storm/windowing/CountTriggerPolicy.java b/storm-client/src/jvm/org/apache/storm/windowing/CountTriggerPolicy.java index 17750b6..0612df0 100644 --- a/storm-client/src/jvm/org/apache/storm/windowing/CountTriggerPolicy.java +++ b/storm-client/src/jvm/org/apache/storm/windowing/CountTriggerPolicy.java @@ -25,14 +25,14 @@ import java.util.concurrent.atomic.AtomicInteger; * * @param <T> the type of event tracked by this policy. */ -public class CountTriggerPolicy<T> implements TriggerPolicy<T> { +public class CountTriggerPolicy<T> implements TriggerPolicy<T, Integer> { private final int count; private final AtomicInteger currentCount; private final TriggerHandler handler; - private final EvictionPolicy<T> evictionPolicy; + private final EvictionPolicy<T, ?> evictionPolicy; private boolean started; - public CountTriggerPolicy(int count, TriggerHandler handler, EvictionPolicy<T> evictionPolicy) { + public CountTriggerPolicy(int count, TriggerHandler handler, EvictionPolicy<T, ?> evictionPolicy) { this.count = count; this.currentCount = new AtomicInteger(); this.handler = handler; @@ -66,11 +66,21 @@ public class CountTriggerPolicy<T> implements TriggerPolicy<T> { } @Override + public Integer getState() { + return currentCount.get(); + } + + @Override + public void restoreState(Integer state) { + currentCount.set(state); + } + + @Override public String toString() { return "CountTriggerPolicy{" + - "count=" + count + - ", currentCount=" + currentCount + - ", started=" + started + - '}'; + "count=" + count + + ", currentCount=" + currentCount + + ", started=" + started + + '}'; } }
http://git-wip-us.apache.org/repos/asf/storm/blob/3ba3cabb/storm-client/src/jvm/org/apache/storm/windowing/Event.java ---------------------------------------------------------------------- diff --git a/storm-client/src/jvm/org/apache/storm/windowing/Event.java b/storm-client/src/jvm/org/apache/storm/windowing/Event.java index c967476..bb9a251 100644 --- a/storm-client/src/jvm/org/apache/storm/windowing/Event.java +++ b/storm-client/src/jvm/org/apache/storm/windowing/Event.java @@ -22,7 +22,7 @@ package org.apache.storm.windowing; * * @param <T> the type of the object thats wrapped. E.g Tuple */ -interface Event<T> { +public interface Event<T> { /** * The event timestamp in millis. This could be the time * when the source generated the tuple or the time http://git-wip-us.apache.org/repos/asf/storm/blob/3ba3cabb/storm-client/src/jvm/org/apache/storm/windowing/EvictionPolicy.java ---------------------------------------------------------------------- diff --git a/storm-client/src/jvm/org/apache/storm/windowing/EvictionPolicy.java b/storm-client/src/jvm/org/apache/storm/windowing/EvictionPolicy.java index fa44444..72dbb29 100644 --- a/storm-client/src/jvm/org/apache/storm/windowing/EvictionPolicy.java +++ b/storm-client/src/jvm/org/apache/storm/windowing/EvictionPolicy.java @@ -15,6 +15,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ + package org.apache.storm.windowing; /** @@ -23,30 +24,31 @@ package org.apache.storm.windowing; * * @param <T> the type of event that is tracked. */ -public interface EvictionPolicy<T> { +public interface EvictionPolicy<T, S> { /** * The action to be taken when {@link EvictionPolicy#evict(Event)} is invoked. */ public enum Action { /** - * expire the event and remove it from the queue + * expire the event and remove it from the queue. */ EXPIRE, /** - * process the event in the current window of events + * process the event in the current window of events. */ PROCESS, /** * don't include in the current window but keep the event - * in the queue for evaluating as a part of future windows + * in the queue for evaluating as a part of future windows. */ KEEP, /** * stop processing the queue, there cannot be anymore events - * satisfying the eviction policy + * satisfying the eviction policy. */ STOP } + /** * Decides if an event should be expired from the window, processed in the current * window or kept for later processing. @@ -68,15 +70,34 @@ public interface EvictionPolicy<T> { * Sets a context in the eviction policy that can be used while evicting the events. * E.g. For TimeEvictionPolicy, this could be used to set the reference timestamp. * - * @param context + * @param context the eviction context */ void setContext(EvictionContext context); /** - * Returns the current context that is part of this eviction policy + * Returns the current context that is part of this eviction policy. * * @return the eviction context */ EvictionContext getContext(); + /** + * Resets the eviction policy. + */ + void reset(); + + /** + * Return runtime state to be checkpointed by the framework for restoring the eviction policy + * in case of failures. + * + * @return the state + */ + S getState(); + + /** + * Restore the eviction policy from the state that was earlier checkpointed by the framework. + * + * @param state the state + */ + void restoreState(S state); } http://git-wip-us.apache.org/repos/asf/storm/blob/3ba3cabb/storm-client/src/jvm/org/apache/storm/windowing/StatefulWindowManager.java ---------------------------------------------------------------------- diff --git a/storm-client/src/jvm/org/apache/storm/windowing/StatefulWindowManager.java b/storm-client/src/jvm/org/apache/storm/windowing/StatefulWindowManager.java new file mode 100644 index 0000000..3d6da2b --- /dev/null +++ b/storm-client/src/jvm/org/apache/storm/windowing/StatefulWindowManager.java @@ -0,0 +1,164 @@ +/** + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.storm.windowing; + +import java.util.Collection; +import java.util.Iterator; +import java.util.NoSuchElementException; +import java.util.function.Supplier; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import static org.apache.storm.windowing.EvictionPolicy.Action.EXPIRE; +import static org.apache.storm.windowing.EvictionPolicy.Action.PROCESS; +import static org.apache.storm.windowing.EvictionPolicy.Action.STOP; + +/** + * Window manager that handles windows with state persistence. + */ +public class StatefulWindowManager<T> extends WindowManager<T> { + private static final Logger LOG = LoggerFactory.getLogger(StatefulWindowManager.class); + + public StatefulWindowManager(WindowLifecycleListener<T> lifecycleListener) { + super(lifecycleListener); + } + + /** + * Constructs a {@link StatefulWindowManager} + * @param lifecycleListener the {@link WindowLifecycleListener} + * @param queue a collection where the events in the window can be enqueued. + * <br/> + * <b>Note:</b> This collection has to be thread safe. + */ + public StatefulWindowManager(WindowLifecycleListener<T> lifecycleListener, Collection<Event<T>> queue) { + super(lifecycleListener, queue); + } + + @Override + protected void compactWindow() { + // NOOP + } + + @Override + public boolean onTrigger() { + Supplier<Iterator<T>> scanEventsStateful = this::scanEventsStateful; + Iterator<T> it = scanEventsStateful.get(); + boolean hasEvents = it.hasNext(); + if (hasEvents) { + final IteratorStatus status = new IteratorStatus(); + LOG.debug("invoking windowLifecycleListener onActivation with iterator"); + // reuse the retrieved iterator + Supplier<Iterator<T>> wrapper = new Supplier<Iterator<T>>() { + Iterator<T> initial = it; + @Override + public Iterator<T> get() { + if (status.isValid()) { + Iterator<T> res; + if (initial != null) { + res = initial; + initial = null; + } else { + res = scanEventsStateful.get(); + } + return expiringIterator(res, status); + } + throw new IllegalStateException("Stale window, the window is valid only within the corresponding execute"); + } + }; + windowLifecycleListener.onActivation(wrapper, null, null, evictionPolicy.getContext().getReferenceTime()); + // invalidate the iterator + status.invalidate(); + } else { + LOG.debug("No events in the window, skipping onActivation"); + } + triggerPolicy.reset(); + return hasEvents; + } + + private Iterator<T> scanEventsStateful() { + LOG.debug("Scan events, eviction policy {}", evictionPolicy); + evictionPolicy.reset(); + Iterator<T> it = new Iterator<T>() { + private Iterator<Event<T>> inner = queue.iterator(); + private T windowEvent; + private boolean stopped; + + @Override + public boolean hasNext() { + while (!stopped && windowEvent == null && inner.hasNext()) { + Event<T> cur = inner.next(); + EvictionPolicy.Action action = evictionPolicy.evict(cur); + if (action == EXPIRE) { + inner.remove(); + } else if (action == STOP) { + stopped = true; + } else if (action == PROCESS) { + windowEvent = cur.get(); + } + } + return windowEvent != null; + } + + @Override + public T next() { + if (!hasNext()) { + throw new NoSuchElementException(); + } + T res = windowEvent; + windowEvent = null; + return res; + } + }; + + return it; + + } + + private static <T> Iterator<T> expiringIterator(Iterator<T> inner, IteratorStatus status) { + return new Iterator<T>() { + @Override + public boolean hasNext() { + if (status.isValid()) { + return inner.hasNext(); + } + throw new IllegalStateException("Stale iterator, the iterator is valid only within the corresponding execute"); + } + + @Override + public T next() { + if (status.isValid()) { + return inner.next(); + } + throw new IllegalStateException("Stale iterator, the iterator is valid only within the corresponding execute"); + } + }; + } + + private static class IteratorStatus { + private boolean valid = true; + + void invalidate() { + valid = false; + } + + boolean isValid() { + return valid; + } + } +} http://git-wip-us.apache.org/repos/asf/storm/blob/3ba3cabb/storm-client/src/jvm/org/apache/storm/windowing/TimeEvictionPolicy.java ---------------------------------------------------------------------- diff --git a/storm-client/src/jvm/org/apache/storm/windowing/TimeEvictionPolicy.java b/storm-client/src/jvm/org/apache/storm/windowing/TimeEvictionPolicy.java index d484db3..ea4fb0d 100644 --- a/storm-client/src/jvm/org/apache/storm/windowing/TimeEvictionPolicy.java +++ b/storm-client/src/jvm/org/apache/storm/windowing/TimeEvictionPolicy.java @@ -23,11 +23,11 @@ import org.slf4j.LoggerFactory; /** * Eviction policy that evicts events based on time duration. */ -public class TimeEvictionPolicy<T> implements EvictionPolicy<T> { +public class TimeEvictionPolicy<T> implements EvictionPolicy<T, EvictionContext> { private static final Logger LOG = LoggerFactory.getLogger(TimeEvictionPolicy.class); private final int windowLength; - protected EvictionContext evictionContext; + protected volatile EvictionContext evictionContext; private long delta; /** @@ -86,10 +86,25 @@ public class TimeEvictionPolicy<T> implements EvictionPolicy<T> { } @Override + public void reset() { + // NOOP + } + + @Override + public EvictionContext getState() { + return evictionContext; + } + + @Override + public void restoreState(EvictionContext state) { + this.evictionContext = state; + } + + @Override public String toString() { return "TimeEvictionPolicy{" + - "windowLength=" + windowLength + - ", evictionContext=" + evictionContext + - '}'; + "windowLength=" + windowLength + + ", evictionContext=" + evictionContext + + '}'; } } http://git-wip-us.apache.org/repos/asf/storm/blob/3ba3cabb/storm-client/src/jvm/org/apache/storm/windowing/TimeTriggerPolicy.java ---------------------------------------------------------------------- diff --git a/storm-client/src/jvm/org/apache/storm/windowing/TimeTriggerPolicy.java b/storm-client/src/jvm/org/apache/storm/windowing/TimeTriggerPolicy.java index e23c2e2..2905bef 100644 --- a/storm-client/src/jvm/org/apache/storm/windowing/TimeTriggerPolicy.java +++ b/storm-client/src/jvm/org/apache/storm/windowing/TimeTriggerPolicy.java @@ -30,20 +30,20 @@ import java.util.concurrent.TimeUnit; /** * Invokes {@link TriggerHandler#onTrigger()} after the duration. */ -public class TimeTriggerPolicy<T> implements TriggerPolicy<T> { +public class TimeTriggerPolicy<T> implements TriggerPolicy<T, Void> { private static final Logger LOG = LoggerFactory.getLogger(TimeTriggerPolicy.class); private long duration; private final TriggerHandler handler; private final ScheduledExecutorService executor; - private final EvictionPolicy<T> evictionPolicy; + private final EvictionPolicy<T, ?> evictionPolicy; private ScheduledFuture<?> executorFuture; public TimeTriggerPolicy(long millis, TriggerHandler handler) { this(millis, handler, null); } - public TimeTriggerPolicy(long millis, TriggerHandler handler, EvictionPolicy<T> evictionPolicy) { + public TimeTriggerPolicy(long millis, TriggerHandler handler, EvictionPolicy<T, ?> evictionPolicy) { this.duration = millis; this.handler = handler; this.executor = Executors.newSingleThreadScheduledExecutor(); @@ -129,4 +129,14 @@ public class TimeTriggerPolicy<T> implements TriggerPolicy<T> { } }; } + + @Override + public Void getState() { + return null; + } + + @Override + public void restoreState(Void state) { + + } } http://git-wip-us.apache.org/repos/asf/storm/blob/3ba3cabb/storm-client/src/jvm/org/apache/storm/windowing/TriggerPolicy.java ---------------------------------------------------------------------- diff --git a/storm-client/src/jvm/org/apache/storm/windowing/TriggerPolicy.java b/storm-client/src/jvm/org/apache/storm/windowing/TriggerPolicy.java index 403b78d..2fc64eb 100644 --- a/storm-client/src/jvm/org/apache/storm/windowing/TriggerPolicy.java +++ b/storm-client/src/jvm/org/apache/storm/windowing/TriggerPolicy.java @@ -15,6 +15,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ + package org.apache.storm.windowing; /** @@ -22,7 +23,7 @@ package org.apache.storm.windowing; * * @param <T> the type of the event that is tracked */ -public interface TriggerPolicy<T> { +public interface TriggerPolicy<T, S> { /** * Tracks the event and could use this to invoke the trigger. * @@ -31,7 +32,7 @@ public interface TriggerPolicy<T> { void track(Event<T> event); /** - * resets the trigger policy + * resets the trigger policy. */ void reset(); @@ -46,4 +47,19 @@ public interface TriggerPolicy<T> { * Any clean up could be handled here. */ void shutdown(); + + /** + * Return runtime state to be checkpointed by the framework for restoring the trigger policy + * in case of failures. + * + * @return the state + */ + S getState(); + + /** + * Restore the trigger policy from the state that was earlier checkpointed by the framework. + * + * @param state the state + */ + void restoreState(S state); } http://git-wip-us.apache.org/repos/asf/storm/blob/3ba3cabb/storm-client/src/jvm/org/apache/storm/windowing/TupleWindowIterImpl.java ---------------------------------------------------------------------- diff --git a/storm-client/src/jvm/org/apache/storm/windowing/TupleWindowIterImpl.java b/storm-client/src/jvm/org/apache/storm/windowing/TupleWindowIterImpl.java new file mode 100644 index 0000000..8140723 --- /dev/null +++ b/storm-client/src/jvm/org/apache/storm/windowing/TupleWindowIterImpl.java @@ -0,0 +1,80 @@ +/** + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.storm.windowing; + +import com.google.common.collect.Iterators; +import java.util.ArrayList; +import java.util.Iterator; +import java.util.List; +import java.util.function.Supplier; +import org.apache.storm.tuple.Tuple; + +/** + * An iterator based implementation over the events in a window. + */ +public class TupleWindowIterImpl implements TupleWindow { + private final Supplier<Iterator<Tuple>> tuplesIt; + private final Supplier<Iterator<Tuple>> newTuplesIt; + private final Supplier<Iterator<Tuple>> expiredTuplesIt; + private final Long startTimestamp; + private final Long endTimestamp; + + public TupleWindowIterImpl(Supplier<Iterator<Tuple>> tuplesIt, + Supplier<Iterator<Tuple>> newTuplesIt, + Supplier<Iterator<Tuple>> expiredTuplesIt, + Long startTimestamp, Long endTimestamp) { + this.tuplesIt = tuplesIt; + this.newTuplesIt = newTuplesIt; + this.expiredTuplesIt = expiredTuplesIt; + this.startTimestamp = startTimestamp; + this.endTimestamp = endTimestamp; + } + + @Override + public List<Tuple> get() { + List<Tuple> tuples = new ArrayList<>(); + tuplesIt.get().forEachRemaining(t -> tuples.add(t)); + return tuples; + } + + @Override + public Iterator<Tuple> getIter() { + return Iterators.unmodifiableIterator(tuplesIt.get()); + } + + @Override + public List<Tuple> getNew() { + throw new UnsupportedOperationException("Not implemented"); + } + + @Override + public List<Tuple> getExpired() { + throw new UnsupportedOperationException("Not implemented"); + } + + @Override + public Long getEndTimestamp() { + return endTimestamp; + } + + @Override + public Long getStartTimestamp() { + return startTimestamp; + } +} http://git-wip-us.apache.org/repos/asf/storm/blob/3ba3cabb/storm-client/src/jvm/org/apache/storm/windowing/WatermarkCountEvictionPolicy.java ---------------------------------------------------------------------- diff --git a/storm-client/src/jvm/org/apache/storm/windowing/WatermarkCountEvictionPolicy.java b/storm-client/src/jvm/org/apache/storm/windowing/WatermarkCountEvictionPolicy.java index 0fe6f75..73fe325 100644 --- a/storm-client/src/jvm/org/apache/storm/windowing/WatermarkCountEvictionPolicy.java +++ b/storm-client/src/jvm/org/apache/storm/windowing/WatermarkCountEvictionPolicy.java @@ -17,20 +17,28 @@ */ package org.apache.storm.windowing; +import org.apache.storm.streams.Pair; + +import java.util.concurrent.atomic.AtomicLong; + /** * An eviction policy that tracks count based on watermark ts and * evicts events up to the watermark based on a threshold count. * * @param <T> the type of event tracked by this policy. */ -public class WatermarkCountEvictionPolicy<T> extends CountEvictionPolicy<T> { - private long processed = 0L; +public class WatermarkCountEvictionPolicy<T> implements EvictionPolicy<T, Pair<Long, Long>> { + protected final int threshold; + protected final AtomicLong currentCount; + private EvictionContext context; + + private volatile long processed; public WatermarkCountEvictionPolicy(int count) { - super(count); + threshold = count; + currentCount = new AtomicLong(); } - @Override public Action evict(Event<T> event) { if(getContext() == null) { //It is possible to get asked about eviction before we have a context, due to WindowManager.compactWindow. @@ -41,7 +49,7 @@ public class WatermarkCountEvictionPolicy<T> extends CountEvictionPolicy<T> { Action action; if (event.getTimestamp() <= getContext().getReferenceTime() && processed < currentCount.get()) { - action = super.evict(event); + action = doEvict(event); if (action == Action.PROCESS) { ++processed; } @@ -51,14 +59,37 @@ public class WatermarkCountEvictionPolicy<T> extends CountEvictionPolicy<T> { return action; } + private Action doEvict(Event<T> event) { + /* + * atomically decrement the count if its greater than threshold and + * return if the event should be evicted + */ + while (true) { + long curVal = currentCount.get(); + if (curVal > threshold) { + if (currentCount.compareAndSet(curVal, curVal - 1)) { + return Action.EXPIRE; + } + } else { + break; + } + } + return Action.PROCESS; + } + @Override public void track(Event<T> event) { // NOOP } @Override + public EvictionContext getContext() { + return context; + } + + @Override public void setContext(EvictionContext context) { - super.setContext(context); + this.context = context; if (context.getCurrentCount() != null) { currentCount.set(context.getCurrentCount()); } else { @@ -68,8 +99,24 @@ public class WatermarkCountEvictionPolicy<T> extends CountEvictionPolicy<T> { } @Override + public void reset() { + processed = 0; + } + + @Override + public Pair<Long, Long> getState() { + return Pair.of(currentCount.get(), processed); + } + + @Override + public void restoreState(Pair<Long, Long> state) { + currentCount.set(state.getFirst()); + processed = state.getSecond(); + } + + @Override public String toString() { return "WatermarkCountEvictionPolicy{" + - "} " + super.toString(); + "} " + super.toString(); } } http://git-wip-us.apache.org/repos/asf/storm/blob/3ba3cabb/storm-client/src/jvm/org/apache/storm/windowing/WatermarkCountTriggerPolicy.java ---------------------------------------------------------------------- diff --git a/storm-client/src/jvm/org/apache/storm/windowing/WatermarkCountTriggerPolicy.java b/storm-client/src/jvm/org/apache/storm/windowing/WatermarkCountTriggerPolicy.java index 3cfcaad..b6ab43b 100644 --- a/storm-client/src/jvm/org/apache/storm/windowing/WatermarkCountTriggerPolicy.java +++ b/storm-client/src/jvm/org/apache/storm/windowing/WatermarkCountTriggerPolicy.java @@ -25,16 +25,16 @@ import java.util.List; * * @param <T> the type of event tracked by this policy. */ -public class WatermarkCountTriggerPolicy<T> implements TriggerPolicy<T> { +public class WatermarkCountTriggerPolicy<T> implements TriggerPolicy<T, Long> { private final int count; private final TriggerHandler handler; - private final EvictionPolicy<T> evictionPolicy; + private final EvictionPolicy<T, ?> evictionPolicy; private final WindowManager<T> windowManager; - private long lastProcessedTs = 0; + private volatile long lastProcessedTs; private boolean started; public WatermarkCountTriggerPolicy(int count, TriggerHandler handler, - EvictionPolicy<T> evictionPolicy, WindowManager<T> windowManager) { + EvictionPolicy<T, ?> evictionPolicy, WindowManager<T> windowManager) { this.count = count; this.handler = handler; this.evictionPolicy = evictionPolicy; @@ -81,11 +81,21 @@ public class WatermarkCountTriggerPolicy<T> implements TriggerPolicy<T> { } @Override + public Long getState() { + return lastProcessedTs; + } + + @Override + public void restoreState(Long state) { + lastProcessedTs = state; + } + + @Override public String toString() { return "WatermarkCountTriggerPolicy{" + - "count=" + count + - ", lastProcessedTs=" + lastProcessedTs + - ", started=" + started + - '}'; + "count=" + count + + ", lastProcessedTs=" + lastProcessedTs + + ", started=" + started + + '}'; } } http://git-wip-us.apache.org/repos/asf/storm/blob/3ba3cabb/storm-client/src/jvm/org/apache/storm/windowing/WatermarkTimeEvictionPolicy.java ---------------------------------------------------------------------- diff --git a/storm-client/src/jvm/org/apache/storm/windowing/WatermarkTimeEvictionPolicy.java b/storm-client/src/jvm/org/apache/storm/windowing/WatermarkTimeEvictionPolicy.java index fdb3917..448abe9 100644 --- a/storm-client/src/jvm/org/apache/storm/windowing/WatermarkTimeEvictionPolicy.java +++ b/storm-client/src/jvm/org/apache/storm/windowing/WatermarkTimeEvictionPolicy.java @@ -81,4 +81,5 @@ public class WatermarkTimeEvictionPolicy<T> extends TimeEvictionPolicy<T> { "lag=" + lag + "} " + super.toString(); } + } http://git-wip-us.apache.org/repos/asf/storm/blob/3ba3cabb/storm-client/src/jvm/org/apache/storm/windowing/WatermarkTimeTriggerPolicy.java ---------------------------------------------------------------------- diff --git a/storm-client/src/jvm/org/apache/storm/windowing/WatermarkTimeTriggerPolicy.java b/storm-client/src/jvm/org/apache/storm/windowing/WatermarkTimeTriggerPolicy.java index 00620c4..ed5c9ff 100644 --- a/storm-client/src/jvm/org/apache/storm/windowing/WatermarkTimeTriggerPolicy.java +++ b/storm-client/src/jvm/org/apache/storm/windowing/WatermarkTimeTriggerPolicy.java @@ -24,16 +24,16 @@ import org.slf4j.LoggerFactory; * Handles watermark events and triggers {@link TriggerHandler#onTrigger()} for each window * interval that has events to be processed up to the watermark ts. */ -public class WatermarkTimeTriggerPolicy<T> implements TriggerPolicy<T> { +public class WatermarkTimeTriggerPolicy<T> implements TriggerPolicy<T, Long> { private static final Logger LOG = LoggerFactory.getLogger(WatermarkTimeTriggerPolicy.class); private final long slidingIntervalMs; private final TriggerHandler handler; - private final EvictionPolicy<T> evictionPolicy; + private final EvictionPolicy<T, ?> evictionPolicy; private final WindowManager<T> windowManager; - private long nextWindowEndTs = 0; + private volatile long nextWindowEndTs; private boolean started; - public WatermarkTimeTriggerPolicy(long slidingIntervalMs, TriggerHandler handler, EvictionPolicy<T> evictionPolicy, + public WatermarkTimeTriggerPolicy(long slidingIntervalMs, TriggerHandler handler, EvictionPolicy<T, ?> evictionPolicy, WindowManager<T> windowManager) { this.slidingIntervalMs = slidingIntervalMs; this.handler = handler; @@ -116,11 +116,21 @@ public class WatermarkTimeTriggerPolicy<T> implements TriggerPolicy<T> { } @Override + public Long getState() { + return nextWindowEndTs; + } + + @Override + public void restoreState(Long state) { + nextWindowEndTs = state; + } + + @Override public String toString() { return "WatermarkTimeTriggerPolicy{" + - "slidingIntervalMs=" + slidingIntervalMs + - ", nextWindowEndTs=" + nextWindowEndTs + - ", started=" + started + - '}'; + "slidingIntervalMs=" + slidingIntervalMs + + ", nextWindowEndTs=" + nextWindowEndTs + + ", started=" + started + + '}'; } } http://git-wip-us.apache.org/repos/asf/storm/blob/3ba3cabb/storm-client/src/jvm/org/apache/storm/windowing/Window.java ---------------------------------------------------------------------- diff --git a/storm-client/src/jvm/org/apache/storm/windowing/Window.java b/storm-client/src/jvm/org/apache/storm/windowing/Window.java index 2e2973a..43ce4c8 100644 --- a/storm-client/src/jvm/org/apache/storm/windowing/Window.java +++ b/storm-client/src/jvm/org/apache/storm/windowing/Window.java @@ -17,6 +17,9 @@ */ package org.apache.storm.windowing; +import org.apache.storm.topology.base.BaseStatefulWindowedBolt; + +import java.util.Iterator; import java.util.List; /** @@ -27,22 +30,44 @@ import java.util.List; public interface Window<T> { /** * Gets the list of events in the window. - * + * <p> + * <b>Note: </b> If the number of tuples in windows is huge, invoking {@code get} would + * load all the tuples into memory and may throw an OOM exception. Use windowing with persistence + * ({@link BaseStatefulWindowedBolt#withPersistence()}) and {@link Window#getIter} to retrieve an iterator over the events in the window. + * </p> * @return the list of events in the window. */ List<T> get(); /** + * Returns an iterator over the events in the window. + * <p> + * <b>Note: </b> This is only supported when using windowing with persistence {@link BaseStatefulWindowedBolt#withPersistence()}. + * </p> + * @return an {@link Iterator} over the events in the current window. + * @throws UnsupportedOperationException if not using {@link BaseStatefulWindowedBolt#withPersistence()} + */ + default Iterator<T> getIter() { + throw new UnsupportedOperationException("Not implemented"); + } + + /** * Get the list of newly added events in the window since the last time the window was generated. - * + * <p> + * <b>Note: </b> This is not supported when using windowing with persistence ({@link BaseStatefulWindowedBolt#withPersistence()}). + * </p> * @return the list of newly added events in the window. + * @throws UnsupportedOperationException if using {@link BaseStatefulWindowedBolt#withPersistence()} */ List<T> getNew(); /** * Get the list of events expired from the window since the last time the window was generated. - * + * <p> + * <b>Note: </b> This is not supported when using windowing with persistence ({@link BaseStatefulWindowedBolt#withPersistence()}). + * </p> * @return the list of events expired from the window. + * @throws UnsupportedOperationException if using {@link BaseStatefulWindowedBolt#withPersistence()} */ List<T> getExpired(); http://git-wip-us.apache.org/repos/asf/storm/blob/3ba3cabb/storm-client/src/jvm/org/apache/storm/windowing/WindowLifecycleListener.java ---------------------------------------------------------------------- diff --git a/storm-client/src/jvm/org/apache/storm/windowing/WindowLifecycleListener.java b/storm-client/src/jvm/org/apache/storm/windowing/WindowLifecycleListener.java index ea2c997..a3f9ee4 100644 --- a/storm-client/src/jvm/org/apache/storm/windowing/WindowLifecycleListener.java +++ b/storm-client/src/jvm/org/apache/storm/windowing/WindowLifecycleListener.java @@ -17,7 +17,9 @@ */ package org.apache.storm.windowing; +import java.util.Iterator; import java.util.List; +import java.util.function.Supplier; /** * A callback for expiry, activation of events tracked by the {@link WindowManager} @@ -39,5 +41,20 @@ public interface WindowLifecycleListener<T> { * @param expired the expired events since last activation. * @param referenceTime the reference (event or processing) time that resulted in activation */ - void onActivation(List<T> events, List<T> newEvents, List<T> expired, Long referenceTime); + default void onActivation(List<T> events, List<T> newEvents, List<T> expired, Long referenceTime) { + throw new UnsupportedOperationException("Not implemented"); + } + + /** + * Called on activation of the window due to the {@link TriggerPolicy}. This is typically invoked when + * the windows are persisted in state and is huge to be loaded entirely in memory. + * + * @param eventsIt a supplier of iterator over the list of current events in the window + * @param newEventsIt a supplier of iterator over the newly added events since the last ativation + * @param expiredIt a supplier of iterator over the expired events since the last activation + * @param referenceTime the reference (event or processing) time that resulted in activation + */ + default void onActivation(Supplier<Iterator<T>> eventsIt, Supplier<Iterator<T>> newEventsIt, Supplier<Iterator<T>> expiredIt, Long referenceTime) { + throw new UnsupportedOperationException("Not implemented"); + } } http://git-wip-us.apache.org/repos/asf/storm/blob/3ba3cabb/storm-client/src/jvm/org/apache/storm/windowing/WindowManager.java ---------------------------------------------------------------------- diff --git a/storm-client/src/jvm/org/apache/storm/windowing/WindowManager.java b/storm-client/src/jvm/org/apache/storm/windowing/WindowManager.java index f6cc521..d0d5fc3 100644 --- a/storm-client/src/jvm/org/apache/storm/windowing/WindowManager.java +++ b/storm-client/src/jvm/org/apache/storm/windowing/WindowManager.java @@ -15,16 +15,21 @@ * See the License for the specific language governing permissions and * limitations under the License. */ + package org.apache.storm.windowing; +import com.google.common.collect.ImmutableMap; import org.apache.storm.windowing.EvictionPolicy.Action; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import java.util.ArrayList; +import java.util.Collection; import java.util.HashSet; import java.util.Iterator; import java.util.List; +import java.util.Map; +import java.util.Optional; import java.util.Set; import java.util.concurrent.ConcurrentLinkedQueue; import java.util.concurrent.atomic.AtomicInteger; @@ -42,6 +47,8 @@ import static org.apache.storm.windowing.EvictionPolicy.Action.STOP; */ public class WindowManager<T> implements TriggerHandler { private static final Logger LOG = LoggerFactory.getLogger(WindowManager.class); + private static final String EVICTION_STATE_KEY = "es"; + private static final String TRIGGER_STATE_KEY = "ts"; /** * Expire old events every EXPIRE_EVENTS_THRESHOLD to @@ -52,29 +59,41 @@ public class WindowManager<T> implements TriggerHandler { */ public static final int EXPIRE_EVENTS_THRESHOLD = 100; - private final WindowLifecycleListener<T> windowLifecycleListener; - private final ConcurrentLinkedQueue<Event<T>> queue; + protected final Collection<Event<T>> queue; + protected EvictionPolicy<T, ?> evictionPolicy; + protected TriggerPolicy<T, ?> triggerPolicy; + protected final WindowLifecycleListener<T> windowLifecycleListener; private final List<T> expiredEvents; private final Set<Event<T>> prevWindowEvents; private final AtomicInteger eventsSinceLastExpiry; private final ReentrantLock lock; - private EvictionPolicy<T> evictionPolicy; - private TriggerPolicy<T> triggerPolicy; public WindowManager(WindowLifecycleListener<T> lifecycleListener) { + this(lifecycleListener, new ConcurrentLinkedQueue<>()); + } + + /** + * Constructs a {@link WindowManager} + * @param lifecycleListener the {@link WindowLifecycleListener} + * @param queue a collection where the events in the window can be enqueued. + * <br/> + * <b>Note:</b> This collection has to be thread safe. + */ + public WindowManager(WindowLifecycleListener<T> lifecycleListener, Collection<Event<T>> queue) { windowLifecycleListener = lifecycleListener; - queue = new ConcurrentLinkedQueue<>(); + this.queue = queue; expiredEvents = new ArrayList<>(); prevWindowEvents = new HashSet<>(); eventsSinceLastExpiry = new AtomicInteger(); lock = new ReentrantLock(true); + } - public void setEvictionPolicy(EvictionPolicy<T> evictionPolicy) { + public void setEvictionPolicy(EvictionPolicy<T, ?> evictionPolicy) { this.evictionPolicy = evictionPolicy; } - public void setTriggerPolicy(TriggerPolicy<T> triggerPolicy) { + public void setTriggerPolicy(TriggerPolicy<T, ?> triggerPolicy) { this.triggerPolicy = triggerPolicy; } @@ -165,7 +184,7 @@ public class WindowManager<T> implements TriggerHandler { * EXPIRE_EVENTS_THRESHOLD so that the window does not grow * too big. */ - private void compactWindow() { + protected void compactWindow() { if (eventsSinceLastExpiry.incrementAndGet() >= EXPIRE_EVENTS_THRESHOLD) { scanEvents(false); } @@ -289,4 +308,20 @@ public class WindowManager<T> implements TriggerHandler { ", triggerPolicy=" + triggerPolicy + '}'; } + + public void restoreState(Map<String, Optional<?>> state) { + Optional.ofNullable(state.get(EVICTION_STATE_KEY)) + .flatMap(x -> x) + .ifPresent(v -> ((EvictionPolicy) evictionPolicy).restoreState(v)); + Optional.ofNullable(state.get(TRIGGER_STATE_KEY)) + .flatMap(x -> x) + .ifPresent(v -> ((TriggerPolicy) triggerPolicy).restoreState(v)); + } + + public Map<String, Optional<?>> getState() { + return ImmutableMap.of( + EVICTION_STATE_KEY, Optional.ofNullable(evictionPolicy.getState()), + TRIGGER_STATE_KEY, Optional.ofNullable(triggerPolicy.getState()) + ); + } } http://git-wip-us.apache.org/repos/asf/storm/blob/3ba3cabb/storm-client/src/jvm/org/apache/storm/windowing/persistence/SimpleWindowPartitionCache.java ---------------------------------------------------------------------- diff --git a/storm-client/src/jvm/org/apache/storm/windowing/persistence/SimpleWindowPartitionCache.java b/storm-client/src/jvm/org/apache/storm/windowing/persistence/SimpleWindowPartitionCache.java new file mode 100644 index 0000000..3602882 --- /dev/null +++ b/storm-client/src/jvm/org/apache/storm/windowing/persistence/SimpleWindowPartitionCache.java @@ -0,0 +1,203 @@ +/** + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.storm.windowing.persistence; + +import java.util.HashMap; +import java.util.Iterator; +import java.util.Map; +import java.util.Objects; +import java.util.concurrent.ConcurrentMap; +import java.util.concurrent.ConcurrentSkipListMap; +import java.util.concurrent.locks.ReentrantLock; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * A simple implementation that evicts the largest un-pinned entry from the cache. This works well + * for caching window partitions since the access pattern is mostly sequential scans. + */ +public class SimpleWindowPartitionCache<K, V> implements WindowPartitionCache<K, V> { + private static final Logger LOG = LoggerFactory.getLogger(SimpleWindowPartitionCache.class); + + private final ConcurrentSkipListMap<K, V> map = new ConcurrentSkipListMap<>(); + private final Map<K, Long> pinned = new HashMap<>(); + private final long maximumSize; + private final RemovalListener<K, V> removalListener; + private final CacheLoader<K, V> cacheLoader; + private final ReentrantLock lock = new ReentrantLock(true); + private int size; + + @Override + public V get(K key) { + return getOrLoad(key, false); + } + + @Override + public V pinAndGet(K key) { + return getOrLoad(key, true); + } + + @Override + public boolean unpin(K key) { + LOG.debug("unpin '{}'", key); + boolean res = false; + try { + lock.lock(); + Long val = pinned.computeIfPresent(key, (k, v) -> v - 1); + if (val != null) { + if (val <= 0) { + pinned.remove(key); + } + res = true; + } + } finally { + lock.unlock(); + } + LOG.debug("pinned '{}'", pinned); + return res; + } + + @Override + public ConcurrentMap<K, V> asMap() { + return map; + } + + @Override + public void invalidate(K key) { + try { + lock.lock(); + if (isPinned(key)) { + LOG.debug("Entry '{}' is pinned, skipping invalidation", key); + } else { + LOG.debug("Invalidating entry '{}'", key); + V val = map.remove(key); + if (val != null) { + --size; + pinned.remove(key); + if (removalListener != null) { + removalListener.onRemoval(key, val, RemovalCause.EXPLICIT); + } + } + } + } finally { + lock.unlock(); + } + } + + // Get or load from the cache optionally pinning the entry + // so that it wont get evicted from the cache + private V getOrLoad(K key, boolean shouldPin) { + V val; + if (shouldPin) { + try { + lock.lock(); + val = load(key); + pin(key); + } finally { + lock.unlock(); + } + } else { + val = map.get(key); + if (val == null) { + try { + lock.lock(); + val = load(key); + } finally { + lock.unlock(); + } + } + } + + return val; + } + + private V load(K key) { + V val = map.get(key); + if (val == null) { + val = cacheLoader.load(key); + if (val == null) { + throw new NullPointerException("Null value for key " + key); + } + ensureCapacity(); + map.put(key, val); + ++size; + } + return val; + } + + private void ensureCapacity() { + if (size >= maximumSize) { + Iterator<Map.Entry<K, V>> it = map.descendingMap().entrySet().iterator(); + while (it.hasNext()) { + Map.Entry<K, V> next = it.next(); + if (!isPinned(next.getKey())) { + it.remove(); + if (removalListener != null) { + removalListener.onRemoval(next.getKey(), next.getValue(), RemovalCause.REPLACED); + } + --size; + break; + } + } + } + } + + private void pin(K key) { + LOG.debug("pin '{}'", key); + pinned.compute(key, (k, v) -> v == null ? 1L : v + 1); + LOG.debug("pinned '{}'", pinned); + } + + private boolean isPinned(K key) { + return pinned.getOrDefault(key, 0L) > 0; + } + + private SimpleWindowPartitionCache(long maximumSize, RemovalListener<K, V> removalListener, CacheLoader<K, V> cacheLoader) { + if (maximumSize <= 0) { + throw new IllegalArgumentException("maximumSize must be greater than 0"); + } + Objects.requireNonNull(cacheLoader); + this.maximumSize = maximumSize; + this.removalListener = removalListener; + this.cacheLoader = cacheLoader; + } + + public static <K, V> SimpleWindowPartitionCacheBuilder<K, V> newBuilder() { + return new SimpleWindowPartitionCacheBuilder<>(); + } + + public static class SimpleWindowPartitionCacheBuilder<K, V> implements WindowPartitionCache.Builder<K, V> { + private long maximumSize; + private RemovalListener<K, V> removalListener; + + public SimpleWindowPartitionCacheBuilder<K, V> maximumSize(long size) { + maximumSize = size; + return this; + } + + public SimpleWindowPartitionCacheBuilder<K, V> removalListener(RemovalListener<K, V> listener) { + removalListener = listener; + return this; + } + + public SimpleWindowPartitionCache<K, V> build(CacheLoader<K, V> loader) { + return new SimpleWindowPartitionCache<>(maximumSize, removalListener, loader); + } + } +} http://git-wip-us.apache.org/repos/asf/storm/blob/3ba3cabb/storm-client/src/jvm/org/apache/storm/windowing/persistence/WindowPartitionCache.java ---------------------------------------------------------------------- diff --git a/storm-client/src/jvm/org/apache/storm/windowing/persistence/WindowPartitionCache.java b/storm-client/src/jvm/org/apache/storm/windowing/persistence/WindowPartitionCache.java new file mode 100644 index 0000000..f1d37e7 --- /dev/null +++ b/storm-client/src/jvm/org/apache/storm/windowing/persistence/WindowPartitionCache.java @@ -0,0 +1,142 @@ +/** + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.storm.windowing.persistence; + +import java.util.concurrent.ConcurrentMap; + +/** + * A loading cache abstraction for caching {@link WindowState.WindowPartition}. + * + * @param <K> the key type + * @param <V> the value type + */ +public interface WindowPartitionCache<K, V> { + + /** + * Get value from the cache or load the value. + * + * @param key the key + * @return the value + */ + V get(K key); + + /** + * Get value from the cache or load the value pinning it + * so that the entry will never get evicted. + * + * @param key the key + * @return the value + */ + V pinAndGet(K key); + + /** + * Unpin an entry from the cache so that it can be a candidate for eviction. + * + * @param key the key + * @return true if the entry was unpinned, false otherwise + */ + boolean unpin(K key); + + /** + * Return a {@link ConcurrentMap} view of the current entries in the cache. + * + * @return the map of key-values currently cached. + */ + ConcurrentMap<K, V> asMap(); + + /** + * Invalidate an entry from the cache. + * + * @param key the key + */ + void invalidate(K key); + + /** + * The reason why an enrty got evicted from the cache. + */ + enum RemovalCause { + /** + * The entry was forcefully invalidated from the cache. + */ + EXPLICIT, + /** + * The entry was evicted from the cache due to overflow. + */ + REPLACED + } + + /** + * A callback interface for handling removal of events from the cache. + * + * @param <K> the key type + * @param <V> the value type + */ + interface RemovalListener<K, V> { + /** + * The method that is invoked when an entry is removed from the cache. + * + * @param key the key of the entry that was removed + * @param val the value of the entry that was removed + * @param removalCause the {@link RemovalCause} + */ + void onRemoval(K key, V val, RemovalCause removalCause); + } + + /** + * The interface for loading entires into the cache. + * + * @param <K> the key type + * @param <V> the value type + */ + interface CacheLoader<K, V> { + V load(K key); + } + + /** + * Builder interface for {@link WindowPartitionCache}. + * + * @param <K> the key type + * @param <V> the value type + */ + interface Builder<K, V> { + /** + * The maximum cache size. After this limit, entries are evicted from the cache. + * + * @param size the size + * @return the Builder + */ + Builder<K, V> maximumSize(long size); + + /** + * The {@link RemovalListener} to be invoked when entries are evicted. + * + * @param listener the listener + * @return the builder + */ + Builder<K, V> removalListener(RemovalListener<K, V> listener); + + /** + * Build the cache. + * + * @param loader the {@link CacheLoader} + * @return the {@link WindowPartitionCache} + */ + WindowPartitionCache<K, V> build(CacheLoader<K, V> loader); + } +} http://git-wip-us.apache.org/repos/asf/storm/blob/3ba3cabb/storm-client/src/jvm/org/apache/storm/windowing/persistence/WindowState.java ---------------------------------------------------------------------- diff --git a/storm-client/src/jvm/org/apache/storm/windowing/persistence/WindowState.java b/storm-client/src/jvm/org/apache/storm/windowing/persistence/WindowState.java new file mode 100644 index 0000000..d373636 --- /dev/null +++ b/storm-client/src/jvm/org/apache/storm/windowing/persistence/WindowState.java @@ -0,0 +1,424 @@ +/** + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * <p> + * http://www.apache.org/licenses/LICENSE-2.0 + * <p> + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.storm.windowing.persistence; + +import com.google.common.collect.ImmutableMap; +import java.util.AbstractCollection; +import java.util.ArrayList; +import java.util.Collection; +import java.util.Collections; +import java.util.Deque; +import java.util.HashSet; +import java.util.Iterator; +import java.util.LinkedList; +import java.util.Map; +import java.util.NoSuchElementException; +import java.util.Objects; +import java.util.Optional; +import java.util.Set; +import java.util.concurrent.ConcurrentLinkedQueue; +import java.util.concurrent.atomic.AtomicInteger; +import java.util.concurrent.locks.ReentrantLock; +import java.util.function.Supplier; +import org.apache.storm.state.KeyValueState; +import org.apache.storm.windowing.Event; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * A wrapper around the window related states that are checkpointed. + */ +public class WindowState<T> extends AbstractCollection<Event<T>> { + private static final Logger LOG = LoggerFactory.getLogger(WindowState.class); + + // number of events per window-partition + public static final int MAX_PARTITION_EVENTS = 1000; + public static final int MIN_PARTITIONS = 10; + private static final String PARTITION_IDS_KEY = "pk"; + private final KeyValueState<String, Deque<Long>> partitionIdsState; + private final KeyValueState<Long, WindowPartition<T>> windowPartitionsState; + private final KeyValueState<String, Optional<?>> windowSystemState; + // ordered partition keys + private volatile Deque<Long> partitionIds; + private volatile long latestPartitionId; + private volatile WindowPartition<T> latestPartition; + private volatile WindowPartitionCache<Long, WindowPartition<T>> cache; + private Supplier<Map<String, Optional<?>>> windowSystemStateSupplier; + private final ReentrantLock partitionIdsLock = new ReentrantLock(true); + private final WindowPartitionLock windowPartitionsLock = new WindowPartitionLock(); + private final long maxEventsInMemory; + private Set<Long> iteratorPins = new HashSet<>(); + + public WindowState(KeyValueState<Long, WindowPartition<T>> windowPartitionsState, + KeyValueState<String, Deque<Long>> partitionIdsState, + KeyValueState<String, Optional<?>> windowSystemState, + Supplier<Map<String, Optional<?>>> windowSystemStateSupplier, + long maxEventsInMemory) { + this.windowPartitionsState = windowPartitionsState; + this.partitionIdsState = partitionIdsState; + this.windowSystemState = windowSystemState; + this.windowSystemStateSupplier = windowSystemStateSupplier; + this.maxEventsInMemory = Math.max(MAX_PARTITION_EVENTS * MIN_PARTITIONS, maxEventsInMemory); + init(); + } + + @Override + public boolean add(Event<T> event) { + if (latestPartition.size() >= MAX_PARTITION_EVENTS) { + cache.unpin(latestPartition.getId()); + latestPartition = getPinnedPartition(getNextPartitionId()); + } + latestPartition.add(event); + return true; + } + + @Override + public Iterator<Event<T>> iterator() { + + return new Iterator<Event<T>>() { + private Iterator<Long> ids = getIds(); + private Iterator<Event<T>> current = Collections.emptyIterator(); + private Iterator<Event<T>> removeFrom; + private WindowPartition<T> curPartition; + + private Iterator<Long> getIds() { + try { + partitionIdsLock.lock(); + LOG.debug("Iterator partitionIds: {}", partitionIds); + return new ArrayList<>(partitionIds).iterator(); + } finally { + partitionIdsLock.unlock(); + } + } + + @Override + public void remove() { + if (removeFrom == null) { + throw new IllegalStateException("No calls to next() since last call to remove()"); + } + removeFrom.remove(); + removeFrom = null; + } + + @Override + public boolean hasNext() { + boolean curHasNext = current.hasNext(); + while (!curHasNext && ids.hasNext()) { + if (curPartition != null) { + unpin(curPartition.getId()); + } + curPartition = getPinnedPartition(ids.next()); + if (curPartition != null) { + iteratorPins.add(curPartition.getId()); + current = curPartition.iterator(); + curHasNext = current.hasNext(); + } + } + // un-pin the last partition + if (!curHasNext && curPartition != null) { + unpin(curPartition.getId()); + curPartition = null; + } + return curHasNext; + } + + @Override + public Event<T> next() { + if (!hasNext()) { + throw new NoSuchElementException(); + } + removeFrom = current; + return current.next(); + } + + private void unpin(long id) { + cache.unpin(id); + iteratorPins.remove(id); + } + }; + } + + public void clearIteratorPins() { + LOG.debug("clearIteratorPins '{}'", iteratorPins); + Iterator<Long> it = iteratorPins.iterator(); + while (it.hasNext()) { + cache.unpin(it.next()); + it.remove(); + } + } + + @Override + public int size() { + throw new UnsupportedOperationException(); + } + + /** + * Prepares the {@link WindowState} for commit. + * + * @param txid the transaction id + */ + public void prepareCommit(long txid) { + flush(); + partitionIdsState.prepareCommit(txid); + windowPartitionsState.prepareCommit(txid); + windowSystemState.prepareCommit(txid); + } + + /** + * Commits the {@link WindowState}. + * + * @param txid the transaction id + */ + public void commit(long txid) { + partitionIdsState.commit(txid); + windowPartitionsState.commit(txid); + windowSystemState.commit(txid); + } + + /** + * Rolls back the {@link WindowState}. + * + * @param reInit if the members should be synced with the values from the state. + */ + public void rollback(boolean reInit) { + partitionIdsState.rollback(); + windowPartitionsState.rollback(); + windowSystemState.rollback(); + // re-init cache and partitions + if (reInit) { + init(); + } + } + + private void init() { + initCache(); + initPartitions(); + } + + private void initPartitions() { + partitionIds = partitionIdsState.get(PARTITION_IDS_KEY, new LinkedList<>()); + if (partitionIds.isEmpty()) { + partitionIds.add(0L); + partitionIdsState.put(PARTITION_IDS_KEY, partitionIds); + } + latestPartitionId = partitionIds.peekLast(); + latestPartition = cache.pinAndGet(latestPartitionId); + } + + private void initCache() { + long size = maxEventsInMemory / MAX_PARTITION_EVENTS; + LOG.info("maxEventsInMemory: {}, partition size: {}, number of partitions: {}", + maxEventsInMemory, MAX_PARTITION_EVENTS, size); + cache = SimpleWindowPartitionCache.<Long, WindowPartition<T>>newBuilder() + .maximumSize(size) + .removalListener(new WindowPartitionCache.RemovalListener<Long, WindowPartition<T>>() { + @Override + public void onRemoval(Long pid, WindowPartition<T> p, WindowPartitionCache.RemovalCause removalCause) { + Objects.requireNonNull(pid, "Null partition id"); + Objects.requireNonNull(p, "Null window partition"); + LOG.debug("onRemoval for id '{}', WindowPartition '{}'", pid, p); + try { + windowPartitionsLock.lock(pid); + if (p.isEmpty() && pid != latestPartitionId) { + // if the empty partition was not invalidated by flush, but evicted from cache + if (removalCause != WindowPartitionCache.RemovalCause.EXPLICIT) { + deletePartition(pid); + windowPartitionsState.delete(pid); + } + } else if (p.isModified()) { + windowPartitionsState.put(pid, p); + } else { + LOG.debug("WindowPartition '{}' is not modified", pid); + } + } finally { + windowPartitionsLock.unlock(pid); + } + } + }).build(new WindowPartitionCache.CacheLoader<Long, WindowPartition<T>>() { + @Override + public WindowPartition<T> load(Long id) { + LOG.debug("Load partition: {}", id); + // load from state + try { + windowPartitionsLock.lock(id); + return windowPartitionsState.get(id, new WindowPartition<>(id)); + } finally { + windowPartitionsLock.unlock(id); + } + } + }); + } + + private void deletePartition(long pid) { + LOG.debug("Delete partition: {}", pid); + try { + partitionIdsLock.lock(); + partitionIds.remove(pid); + partitionIdsState.put(PARTITION_IDS_KEY, partitionIds); + } finally { + partitionIdsLock.unlock(); + } + } + + private long getNextPartitionId() { + try { + partitionIdsLock.lock(); + partitionIds.add(++latestPartitionId); + partitionIdsState.put(PARTITION_IDS_KEY, partitionIds); + } finally { + partitionIdsLock.unlock(); + } + return latestPartitionId; + } + + private WindowPartition<T> getPinnedPartition(long id) { + return cache.pinAndGet(id); + } + + private void flush() { + LOG.debug("Flushing modified partitions"); + cache.asMap().forEach((pid, p) -> { + Long pidToInvalidate = null; + try { + windowPartitionsLock.lock(pid); + if (p.isEmpty() && pid != latestPartitionId) { + LOG.debug("Invalidating empty partition {}", pid); + deletePartition(pid); + windowPartitionsState.delete(pid); + pidToInvalidate = pid; + } else if (p.isModified()) { + LOG.debug("Updating modified partition {}", pid); + p.clearModified(); + windowPartitionsState.put(pid, p); + } + } finally { + windowPartitionsLock.unlock(pid); + } + // invalidate after releasing the lock + // if the parition is pinned before we could invalidate, + // it will get invalidated in the next flush or when the entry gets evicted from the cache. + if (pidToInvalidate != null) { + cache.invalidate(pidToInvalidate); + } + }); + Map<String, Optional<?>> state = windowSystemStateSupplier.get(); + for (Map.Entry<String, Optional<?>> entry: state.entrySet()) { + windowSystemState.put(entry.getKey(), entry.getValue()); + } + } + + private static class WindowPartitionLock { + private final int numLocks = 8; + private final ImmutableMap<Long, ReentrantLock> locks; + + WindowPartitionLock() { + ImmutableMap.Builder<Long, ReentrantLock> builder = ImmutableMap.builder(); + for (long i = 0; i < numLocks; i++) { + builder.put(i, new ReentrantLock(true)); + } + locks = builder.build(); + } + + private void lock(long i) { + locks.get(i % numLocks).lock(); + } + + private void unlock(long i) { + locks.get(i % numLocks).unlock(); + } + } + + // the window partition that holds the events + public static class WindowPartition<T> implements Iterable<Event<T>> { + private final ConcurrentLinkedQueue<Event<T>> events = new ConcurrentLinkedQueue<>(); + private final AtomicInteger size = new AtomicInteger(); + private final long id; + private transient volatile boolean modified; + + public WindowPartition(long id) { + this.id = id; + } + + void add(Event<T> event) { + events.add(event); + size.incrementAndGet(); + setModified(); + } + + boolean isModified() { + return modified; + } + + void setModified() { + if (!modified) { + modified = true; + } + } + + void clearModified() { + modified = false; + } + + boolean isEmpty() { + return events.isEmpty(); + } + + @Override + public Iterator<Event<T>> iterator() { + return new Iterator<Event<T>>() { + Iterator<Event<T>> it = events.iterator(); + + @Override + public boolean hasNext() { + return it.hasNext(); + } + + @Override + public Event<T> next() { + return it.next(); + } + + @Override + public void remove() { + it.remove(); + size.decrementAndGet(); + setModified(); + } + }; + } + + public int size() { + return size.get(); + } + + public long getId() { + return id; + } + + // for unit tests + public Collection<Event<T>> getEvents() { + return Collections.unmodifiableCollection(events); + } + + @Override + public String toString() { + return "WindowPartition{id=" + id + ", size=" + size + '}'; + } + } +} http://git-wip-us.apache.org/repos/asf/storm/blob/3ba3cabb/storm-client/test/jvm/org/apache/storm/state/DefaultStateSerializerTest.java ---------------------------------------------------------------------- diff --git a/storm-client/test/jvm/org/apache/storm/state/DefaultStateSerializerTest.java b/storm-client/test/jvm/org/apache/storm/state/DefaultStateSerializerTest.java index c289cba..6d952e2 100644 --- a/storm-client/test/jvm/org/apache/storm/state/DefaultStateSerializerTest.java +++ b/storm-client/test/jvm/org/apache/storm/state/DefaultStateSerializerTest.java @@ -21,6 +21,7 @@ import org.apache.storm.spout.CheckPointState; import org.junit.Test; import java.util.ArrayList; +import java.util.Collections; import java.util.List; import static org.junit.Assert.*; @@ -45,7 +46,7 @@ public class DefaultStateSerializerTest { List<Class<?>> classesToRegister = new ArrayList<>(); classesToRegister.add(CheckPointState.class); - Serializer<CheckPointState> s3 = new DefaultStateSerializer<CheckPointState>(classesToRegister); + Serializer<CheckPointState> s3 = new DefaultStateSerializer<>(Collections.emptyMap(), null, classesToRegister); bytes = s2.serialize(cs); assertEquals(cs, (CheckPointState) s2.deserialize(bytes)); http://git-wip-us.apache.org/repos/asf/storm/blob/3ba3cabb/storm-client/test/jvm/org/apache/storm/topology/PersistentWindowedBoltExecutorTest.java ---------------------------------------------------------------------- diff --git a/storm-client/test/jvm/org/apache/storm/topology/PersistentWindowedBoltExecutorTest.java b/storm-client/test/jvm/org/apache/storm/topology/PersistentWindowedBoltExecutorTest.java new file mode 100644 index 0000000..d486487 --- /dev/null +++ b/storm-client/test/jvm/org/apache/storm/topology/PersistentWindowedBoltExecutorTest.java @@ -0,0 +1,299 @@ +/** + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.storm.topology; + +import com.google.common.collect.ImmutableMap; +import org.apache.storm.Config; +import org.apache.storm.generated.GlobalStreamId; +import org.apache.storm.state.KeyValueState; +import org.apache.storm.streams.Pair; +import org.apache.storm.task.OutputCollector; +import org.apache.storm.task.TopologyContext; +import org.apache.storm.tuple.Tuple; +import org.apache.storm.tuple.Values; +import org.apache.storm.windowing.Event; +import org.apache.storm.windowing.TimestampExtractor; +import org.apache.storm.windowing.TupleWindow; +import org.apache.storm.windowing.WaterMarkEvent; +import org.apache.storm.windowing.WaterMarkEventGenerator; +import org.apache.storm.windowing.persistence.WindowState; +import org.junit.Assert; +import org.junit.Before; +import org.junit.Test; +import org.junit.runner.RunWith; +import org.mockito.ArgumentCaptor; +import org.mockito.Captor; +import org.mockito.Mock; +import org.mockito.Mockito; +import org.mockito.MockitoAnnotations; +import org.mockito.invocation.InvocationOnMock; +import org.mockito.runners.MockitoJUnitRunner; +import org.mockito.stubbing.Answer; + +import java.util.ArrayList; +import java.util.Arrays; +import java.util.Collection; +import java.util.Collections; +import java.util.Deque; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.Optional; +import java.util.stream.Collectors; +import java.util.stream.LongStream; + +import static org.mockito.AdditionalAnswers.returnsArgAt; + +/** + * Unit tests for {@link PersistentWindowedBoltExecutor} + */ +@RunWith(MockitoJUnitRunner.class) +public class PersistentWindowedBoltExecutorTest { + private static final String LATE_STREAM = "late_stream"; + private static final String PARTITION_KEY = "pk"; + private static final String EVICTION_STATE_KEY = "es"; + private static final String TRIGGER_STATE_KEY = "ts"; + private static final int WINDOW_EVENT_COUNT = 5; + + private long tupleTs; + private PersistentWindowedBoltExecutor<KeyValueState<String, String>> executor; + private IStatefulWindowedBolt<KeyValueState<String, String>> mockBolt; + private Map<String, Object> testStormConf = new HashMap<>(); + private OutputCollector mockOutputCollector; + private TopologyContext mockTopologyContext; + private TimestampExtractor mockTimestampExtractor; + private WaterMarkEventGenerator mockWaterMarkEventGenerator; + + @Mock + private KeyValueState<String, Deque<Long>> mockPartitionState; + @Mock + private KeyValueState<Long, WindowState.WindowPartition<Tuple>> mockWindowState; + @Mock + private KeyValueState<String, Optional<?>> mockSystemState; + + @Captor + private ArgumentCaptor<Tuple> tupleCaptor; + @Captor + private ArgumentCaptor<Collection<Tuple>> anchorCaptor; + @Captor + private ArgumentCaptor<Long> longCaptor; + @Captor + private ArgumentCaptor<Values> valuesCaptor; + @Captor + private ArgumentCaptor<TupleWindow> tupleWindowCaptor; + @Captor + private ArgumentCaptor<Deque<Long>> partitionValuesCaptor; + @Captor + private ArgumentCaptor<WindowState.WindowPartition<Tuple>> windowValuesCaptor; + @Captor + private ArgumentCaptor<Optional<?>> systemValuesCaptor; + + @Before + public void setUp() throws Exception { + MockitoAnnotations.initMocks(this); + mockBolt = Mockito.mock(IStatefulWindowedBolt.class); + mockWaterMarkEventGenerator = Mockito.mock(WaterMarkEventGenerator.class); + mockTimestampExtractor = Mockito.mock(TimestampExtractor.class); + tupleTs = System.currentTimeMillis(); + Mockito.when(mockTimestampExtractor.extractTimestamp(Mockito.any())).thenReturn(tupleTs); + Mockito.when(mockBolt.getTimestampExtractor()).thenReturn(mockTimestampExtractor); + Mockito.when(mockBolt.isPersistent()).thenReturn(true); + mockTopologyContext = Mockito.mock(TopologyContext.class); + Mockito.when(mockTopologyContext.getThisStreams()).thenReturn(Collections.singleton(LATE_STREAM)); + mockOutputCollector = Mockito.mock(OutputCollector.class); + executor = new PersistentWindowedBoltExecutor<>(mockBolt); + testStormConf.put(Config.TOPOLOGY_BOLTS_WINDOW_LENGTH_COUNT, WINDOW_EVENT_COUNT); + testStormConf.put(Config.TOPOLOGY_BOLTS_SLIDING_INTERVAL_COUNT, WINDOW_EVENT_COUNT); + testStormConf.put(Config.TOPOLOGY_BOLTS_LATE_TUPLE_STREAM, LATE_STREAM); + testStormConf.put(Config.TOPOLOGY_BOLTS_WATERMARK_EVENT_INTERVAL_MS, 100_000); + testStormConf.put(Config.TOPOLOGY_MESSAGE_TIMEOUT_SECS, 30); + testStormConf.put(Config.TOPOLOGY_STATE_CHECKPOINT_INTERVAL, 1000); + Mockito.when(mockPartitionState.get(Mockito.any(), Mockito.any())).then(returnsArgAt(1)); + Mockito.when(mockWindowState.get(Mockito.any(), Mockito.any())).then(returnsArgAt(1)); + Mockito.when(mockSystemState.get(Mockito.any(), Mockito.any())).then(returnsArgAt(1)); + Mockito.when(mockSystemState.iterator()).thenReturn( + ImmutableMap.<String, Optional<?>>of("es", Optional.empty(), "ts", Optional.empty()).entrySet().iterator()); + executor.prepare(testStormConf, mockTopologyContext, mockOutputCollector, + mockWindowState, mockPartitionState, mockSystemState); + } + + @Test + public void testExecuteTuple() throws Exception { + Mockito.when(mockWaterMarkEventGenerator.track(Mockito.any(GlobalStreamId.class), Mockito.anyLong())).thenReturn(true); + Tuple mockTuple = Mockito.mock(Tuple.class); + executor.initState(null); + executor.waterMarkEventGenerator = mockWaterMarkEventGenerator; + executor.execute(mockTuple); + // should be ack-ed once + Mockito.verify(mockOutputCollector, Mockito.times(1)).ack(mockTuple); + } + + @Test + public void testExecuteLatetuple() throws Exception { + Mockito.when(mockWaterMarkEventGenerator.track(Mockito.any(GlobalStreamId.class), Mockito.anyLong())).thenReturn(false); + Tuple mockTuple = Mockito.mock(Tuple.class); + executor.initState(null); + executor.waterMarkEventGenerator = mockWaterMarkEventGenerator; + executor.execute(mockTuple); + // ack-ed once + Mockito.verify(mockOutputCollector, Mockito.times(1)).ack(mockTuple); + // late tuple emitted + ArgumentCaptor<String> stringCaptor = ArgumentCaptor.forClass(String.class); + Mockito.verify(mockOutputCollector, Mockito.times(1)) + .emit(stringCaptor.capture(), anchorCaptor.capture(), valuesCaptor.capture()); + Assert.assertEquals(LATE_STREAM, stringCaptor.getValue()); + Assert.assertEquals(Collections.singletonList(mockTuple), anchorCaptor.getValue()); + Assert.assertEquals(new Values(mockTuple), valuesCaptor.getValue()); + } + + @Test + public void testActivation() throws Exception { + Mockito.when(mockWaterMarkEventGenerator.track(Mockito.any(GlobalStreamId.class), Mockito.anyLong())).thenReturn(true); + executor.initState(null); + executor.waterMarkEventGenerator = mockWaterMarkEventGenerator; + + List<Tuple> mockTuples = getMockTuples(WINDOW_EVENT_COUNT); + mockTuples.forEach(t -> executor.execute(t)); + // all tuples acked + Mockito.verify(mockOutputCollector, Mockito.times(WINDOW_EVENT_COUNT)).ack(tupleCaptor.capture()); + Assert.assertArrayEquals(mockTuples.toArray(), tupleCaptor.getAllValues().toArray()); + + Mockito.doAnswer(new Answer<Void>() { + @Override + public Void answer(InvocationOnMock invocation) throws Throwable { + TupleWindow window = (TupleWindow) invocation.getArguments()[0]; + // iterate the tuples + Assert.assertEquals(WINDOW_EVENT_COUNT, window.get().size()); + // iterating multiple times should produce same events + Assert.assertEquals(WINDOW_EVENT_COUNT, window.get().size()); + Assert.assertEquals(WINDOW_EVENT_COUNT, window.get().size()); + return null; + } + }).when(mockBolt).execute(Mockito.any()); + // trigger the window + long activationTs = tupleTs + 1000; + executor.getWindowManager().add(new WaterMarkEvent<>(activationTs)); + executor.prePrepare(0); + + // partition ids + ArgumentCaptor<String> pkCatptor = ArgumentCaptor.forClass(String.class); + Mockito.verify(mockPartitionState, Mockito.times(1)).put(pkCatptor.capture(), partitionValuesCaptor.capture()); + Assert.assertEquals(PARTITION_KEY, pkCatptor.getValue()); + List<Long> expectedPartitionIds = Collections.singletonList(0L); + Assert.assertEquals(expectedPartitionIds, partitionValuesCaptor.getValue()); + + // window partitions + Mockito.verify(mockWindowState, Mockito.times(1)).put(longCaptor.capture(), windowValuesCaptor.capture()); + Assert.assertEquals((long) expectedPartitionIds.get(0), (long) longCaptor.getValue()); + Assert.assertEquals(WINDOW_EVENT_COUNT, windowValuesCaptor.getValue().size()); + List<Tuple> tuples = windowValuesCaptor.getValue() + .getEvents().stream().map(Event::get).collect(Collectors.toList()); + Assert.assertArrayEquals(mockTuples.toArray(), tuples.toArray()); + + // window system state + ArgumentCaptor<String> keyCaptor = ArgumentCaptor.forClass(String.class); + Mockito.verify(mockSystemState, Mockito.times(2)).put(keyCaptor.capture(), systemValuesCaptor.capture()); + Assert.assertEquals(EVICTION_STATE_KEY, keyCaptor.getAllValues().get(0)); + Assert.assertEquals(Optional.of(Pair.of((long)WINDOW_EVENT_COUNT, (long)WINDOW_EVENT_COUNT)), systemValuesCaptor.getAllValues().get(0)); + Assert.assertEquals(TRIGGER_STATE_KEY, keyCaptor.getAllValues().get(1)); + Assert.assertEquals(Optional.of(tupleTs), systemValuesCaptor.getAllValues().get(1)); + } + + @Test + public void testCacheEviction() { + Mockito.when(mockWaterMarkEventGenerator.track(Mockito.any(GlobalStreamId.class), Mockito.anyLong())).thenReturn(true); + executor.initState(null); + executor.waterMarkEventGenerator = mockWaterMarkEventGenerator; + int tupleCount = 20000; + List<Tuple> mockTuples = getMockTuples(tupleCount); + mockTuples.forEach(t -> executor.execute(t)); + + int numPartitions = tupleCount/WindowState.MAX_PARTITION_EVENTS; + int numEvictedPartitions = numPartitions - WindowState.MIN_PARTITIONS; + Mockito.verify(mockWindowState, Mockito.times(numEvictedPartitions)).put(longCaptor.capture(), windowValuesCaptor.capture()); + // number of evicted events + Assert.assertEquals(numEvictedPartitions*WindowState.MAX_PARTITION_EVENTS, windowValuesCaptor.getAllValues().stream() + .mapToInt(x -> x.size()).sum()); + + Map<Long, WindowState.WindowPartition<Tuple>> partitionMap = new HashMap<>(); + windowValuesCaptor.getAllValues().forEach(v -> partitionMap.put(v.getId(), v)); + + ArgumentCaptor<String> stringCaptor = ArgumentCaptor.forClass(String.class); + Mockito.verify(mockPartitionState, Mockito.times(numPartitions)).put(stringCaptor.capture(), partitionValuesCaptor.capture()); + // partition ids 0 .. 19 + Assert.assertEquals(LongStream.range(0, numPartitions).boxed().collect(Collectors.toList()), partitionValuesCaptor.getAllValues().get(numPartitions-1)); + + Mockito.when(mockWindowState.get(Mockito.any(), Mockito.any())).then(new Answer<Object>() { + @Override + public Object answer(InvocationOnMock invocation) throws Throwable { + Object[] args = invocation.getArguments(); + WindowState.WindowPartition<Tuple> evicted = partitionMap.get(args[0]); + return evicted != null ? evicted : args[1]; + } + }); + + Mockito.doAnswer(new Answer<Void>() { + @Override + public Void answer(InvocationOnMock invocation) throws Throwable { + Object[] args = invocation.getArguments(); + partitionMap.put((long)args[0], (WindowState.WindowPartition<Tuple>)args[1]); + return null; + } + }).when(mockWindowState).put(Mockito.any(), Mockito.any()); + + // trigger the window + long activationTs = tupleTs + 1000; + executor.getWindowManager().add(new WaterMarkEvent<>(activationTs)); + + Mockito.verify(mockBolt, Mockito.times(tupleCount/WINDOW_EVENT_COUNT)).execute(Mockito.any()); + } + + @Test + public void testRollbackBeforeInit() throws Exception { + executor.preRollback(); + Mockito.verify(mockBolt, Mockito.times(1)).preRollback(); + // partition ids + ArgumentCaptor<String> pkCatptor = ArgumentCaptor.forClass(String.class); + Mockito.verify(mockPartitionState, Mockito.times(1)).rollback(); + Mockito.verify(mockWindowState, Mockito.times(1)).rollback(); + Mockito.verify(mockSystemState, Mockito.times(1)).rollback(); + } + + @Test + public void testRollbackAfterInit() throws Exception { + executor.initState(null); + executor.prePrepare(0); + executor.preRollback(); + Mockito.verify(mockBolt, Mockito.times(1)).preRollback(); + Mockito.verify(mockPartitionState, Mockito.times(1)).rollback(); + ArgumentCaptor<String> stringArgumentCaptor = ArgumentCaptor.forClass(String.class); + Mockito.verify(mockPartitionState, Mockito.times(2)).put(stringArgumentCaptor.capture(), partitionValuesCaptor.capture()); + Mockito.verify(mockWindowState, Mockito.times(1)).rollback(); + Mockito.verify(mockSystemState, Mockito.times(1)).rollback(); + Mockito.verify(mockSystemState, Mockito.times(2)).iterator(); + } + + private List<Tuple> getMockTuples(long count) { + List<Tuple> tuples = new ArrayList<>(); + for (int i = 0; i < count; i++) { + tuples.add(Mockito.mock(Tuple.class)); + } + return tuples; + } +} \ No newline at end of file
