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

Reply via email to