STORM-2614: Enhance stateful windowing to persist the window state

Right now the tuples in window are stored in memory. This limits the usage to 
windows
that fit in memory. Also the source tuples cannot be acked until the window 
expiry.
By persisting the window transparently in the state backend and 
caching/iterating them as needed,
we could support larger windows and also support windowed bolts with 
user/application state.


Project: http://git-wip-us.apache.org/repos/asf/storm/repo
Commit: http://git-wip-us.apache.org/repos/asf/storm/commit/3ba3cabb
Tree: http://git-wip-us.apache.org/repos/asf/storm/tree/3ba3cabb
Diff: http://git-wip-us.apache.org/repos/asf/storm/diff/3ba3cabb

Branch: refs/heads/master
Commit: 3ba3cabb01743f6a7181378b89957fff8d14e850
Parents: 1c2ac2e
Author: Arun Mahadevan <[email protected]>
Authored: Wed Jun 28 21:43:26 2017 +0530
Committer: Arun Mahadevan <[email protected]>
Committed: Tue Aug 8 07:55:46 2017 +0530

----------------------------------------------------------------------
 docs/State-checkpointing.md                     |   2 +
 docs/Windowing.md                               | 108 +++++
 .../starter/PersistentWindowingTopology.java    | 176 ++++++++
 .../storm/starter/SlidingWindowTopology.java    |   1 +
 .../apache/storm/starter/StatefulTopology.java  |   2 +
 .../starter/StatefulWindowingTopology.java      |   1 +
 external/storm-hbase/pom.xml                    |   2 -
 .../hbase/state/HBaseKeyValueStateProvider.java |  31 +-
 .../redis/state/RedisKeyValueStateProvider.java |  47 +-
 pom.xml                                         |   6 +
 .../src/jvm/org/apache/storm/Config.java        |   6 +
 .../serialization/SerializationFactory.java     |  63 +--
 .../apache/storm/state/DefaultStateEncoder.java |   5 +-
 .../storm/state/DefaultStateSerializer.java     |  64 ++-
 .../storm/topology/IStatefulWindowedBolt.java   |  22 +-
 .../PersistentWindowedBoltExecutor.java         | 256 +++++++++++
 .../storm/topology/StatefulBoltExecutor.java    |   7 +-
 .../apache/storm/topology/TopologyBuilder.java  |  21 +-
 .../storm/topology/WindowedBoltExecutor.java    |  86 +++-
 .../topology/base/BaseStatefulWindowedBolt.java |  39 ++
 .../windowing/AbstractTridentWindowManager.java |   4 +-
 .../strategy/SlidingCountWindowStrategy.java    |   4 +-
 .../strategy/SlidingDurationWindowStrategy.java |   4 +-
 .../strategy/TumblingCountWindowStrategy.java   |   4 +-
 .../TumblingDurationWindowStrategy.java         |   4 +-
 .../windowing/strategy/WindowStrategy.java      |   4 +-
 .../storm/windowing/CountEvictionPolicy.java    |  17 +-
 .../storm/windowing/CountTriggerPolicy.java     |  24 +-
 .../jvm/org/apache/storm/windowing/Event.java   |   2 +-
 .../apache/storm/windowing/EvictionPolicy.java  |  35 +-
 .../storm/windowing/StatefulWindowManager.java  | 164 +++++++
 .../storm/windowing/TimeEvictionPolicy.java     |  25 +-
 .../storm/windowing/TimeTriggerPolicy.java      |  16 +-
 .../apache/storm/windowing/TriggerPolicy.java   |  20 +-
 .../storm/windowing/TupleWindowIterImpl.java    |  80 ++++
 .../windowing/WatermarkCountEvictionPolicy.java |  61 ++-
 .../windowing/WatermarkCountTriggerPolicy.java  |  26 +-
 .../windowing/WatermarkTimeEvictionPolicy.java  |   1 +
 .../windowing/WatermarkTimeTriggerPolicy.java   |  26 +-
 .../jvm/org/apache/storm/windowing/Window.java  |  31 +-
 .../windowing/WindowLifecycleListener.java      |  19 +-
 .../apache/storm/windowing/WindowManager.java   |  51 ++-
 .../persistence/SimpleWindowPartitionCache.java | 203 +++++++++
 .../persistence/WindowPartitionCache.java       | 142 +++++++
 .../windowing/persistence/WindowState.java      | 424 +++++++++++++++++++
 .../storm/state/DefaultStateSerializerTest.java |   3 +-
 .../PersistentWindowedBoltExecutorTest.java     | 299 +++++++++++++
 .../SimpleWindowPartitionCacheTest.java         | 233 ++++++++++
 .../storm/windowing/WindowManagerTest.java      |  46 +-
 .../windowing/persistence/WindowStateTest.java  | 246 +++++++++++
 50 files changed, 2966 insertions(+), 197 deletions(-)
----------------------------------------------------------------------


http://git-wip-us.apache.org/repos/asf/storm/blob/3ba3cabb/docs/State-checkpointing.md
----------------------------------------------------------------------
diff --git a/docs/State-checkpointing.md b/docs/State-checkpointing.md
index 4be6e23..42af3a7 100644
--- a/docs/State-checkpointing.md
+++ b/docs/State-checkpointing.md
@@ -22,6 +22,7 @@ For example a word count bolt could use the key value state 
abstraction for the
 last committed by the framework during the previous run.
 3. In the execute method, update the word count.
 
+```java
     public class WordCountBolt extends BaseStatefulBolt<KeyValueState<String, 
Long>> {
     private KeyValueState<String, Long> wordCounts;
     private OutputCollector collector;
@@ -45,6 +46,7 @@ last committed by the framework during the previous run.
         }
     ...
     }
+```
 
 4. The framework periodically checkpoints the state of the bolt (default every 
second). The frequency
 can be changed by setting the storm config 
`topology.state.checkpoint.interval.ms`

http://git-wip-us.apache.org/repos/asf/storm/blob/3ba3cabb/docs/Windowing.md
----------------------------------------------------------------------
diff --git a/docs/Windowing.md b/docs/Windowing.md
index 91a6509..01d8857 100644
--- a/docs/Windowing.md
+++ b/docs/Windowing.md
@@ -266,3 +266,111 @@ tuples can be received within the timeout period.
 An example toplogy `SlidingWindowTopology` shows how to use the apis to 
compute a sliding window sum and a tumbling window 
 average.
 
+## Stateful windowing
+The default windowing implementation in storm stores the tuples in memory 
until they are processed and expired from the 
+window. This limits the use cases to windows that
+fit entirely in memory. Also the source tuples cannot be ack-ed until the 
window expiry requiring large message timeouts
+(topology.message.timeout.secs should be larger than the window length + 
sliding interval). This also puts extra loads 
+due to the complex acking and anchoring requirements.
+ 
+To address the above limitations and to support larger window sizes, storm 
provides stateful windowing support via `IStatefulWindowedBolt`. 
+User bolts should typically extend `BaseStatefulWindowedBolt` for the 
windowing operations with the framework automatically 
+managing the state of the window in the background.
+
+If the sources provide a monotonically increasing identifier as a part of the 
message, the framework can use this
+to periodically checkpoint the last expired and evaluated message ids, to 
avoid duplicate window evaluations in case of 
+failures or restarts. During recovery, the tuples with message ids lower than 
last expired id are discarded and tuples with 
+message id between the last expired and last evaluated message ids are fed 
into the system without activating any previously
+activated windows.
+The tuples beyond the last evaluated message ids are processed as usual. This 
can be enabled by setting
+the `messageIdField` as shown below,
+
+```java
+topologyBuilder.setBolt("mybolt",
+                   new MyStatefulWindowedBolt()
+                   .withWindow(...) // windowing configuarations
+                   .withMessageIdField("msgid"), // a monotonically increasing 
'long' field in the tuple
+                   parallelism)
+               .shuffleGrouping("spout");
+```
+
+However, this option is feasible only if the sources can provide a 
monotonically increasing identifier in the tuple and the same is maintained
+while re-emitting the messages in case of failures. With this option the 
tuples are still buffered in memory until processed
+and expired from the window.
+
+For more details take a look at the sample topology in storm-starter 
[StatefulWindowingTopology](../examples/storm-starter/src/jvm/org/apache/storm/starter/StatefulWindowingTopology.java)
 which will help you get started.
+
+### Window checkpointing
+
+With window checkpointing, the monotonically increasing id is no longer 
required since the framework transparently saves the state of the window 
periodically into the configured state backend.
+The state that is saved includes the tuples in the window, any system state 
that is required to recover the state of processing
+and also the user state.
+
+```java
+topologyBuilder.setBolt("mybolt",
+                   new MyStatefulPersistentWindowedBolt()
+                   .withWindow(...) // windowing configuarations
+                   .withPersistence() // persist the window state
+                   .withMaxEventsInMemory(25000), // max number of events to 
be cached in memory
+                    parallelism)
+               .shuffleGrouping("spout");
+
+```
+
+The `withPersistence` instructs the framework to transparently save the tuples 
in window along with
+any associated system and user state to the state backend. The 
`withMaxEventsInMemory` is an optional 
+configuration that specifies the maximum number of tuples that may be kept in 
memory. The tuples are transparently loaded from 
+the state backend as required and the ones that are most likely to be used 
again are retained in memory.
+
+The state backend can be configured by setting the topology state provider 
config,
+
+```java
+// use redis for state persistence
+conf.put(Config.TOPOLOGY_STATE_PROVIDER, 
"org.apache.storm.redis.state.RedisKeyValueStateProvider");
+
+```
+Currently storm supports Redis and HBase as state backends and uses the 
underlying state-checkpointing
+framework for saving the window state. For more details on state checkpointing 
see [State-checkpointing.md](State-checkpointing.md)
+
+Here is an example of a persistent windowed bolt that uses the window 
checkpointing to save its state. The `initState`
+is invoked with the last saved state (user state) at initialization time. The 
execute method is invoked based on the configured
+windowing parameters and the tuples in the active window can be accessed via 
an `iterator` as shown below.
+
+```java
+public class MyStatefulPersistentWindowedBolt extends 
BaseStatefulWindowedBolt<K, V> {
+  private KeyValueState<K, V> state;
+  
+  @Override
+  public void initState(KeyValueState<K, V> state) {
+    this.state = state;
+   // ...
+   // restore the state from the last saved state.
+   // ...
+  }
+  
+  @Override
+  public void execute(TupleWindow window) {      
+    // iterate over tuples in the current window
+    Iterator<Tuple> it = window.getIter();
+    while (it.hasNext()) {
+        // compute some result based on the tuples in window
+    }
+    
+    // possibly update any state to be maintained across windows
+    state.put(STATE_KEY, updatedValue);
+    
+    // emit the results downstream
+    collector.emit(new Values(result));
+  }
+}
+```
+
+**Note:** In case of persistent windowed bolts, use `TupleWindow.getIter` to 
retrieve an iterator over the
+events in the window. If the number of tuples in windows is huge, invoking 
`TupleWindow.get` would
+try to load all the tuples into memory and may throw an OOM exception.
+
+**Note:** In case of persistent windowed bolts the `TupleWindow.getNew` and 
`TupleWindow.getExpired` are currently not supported
+and will throw an `UnsupportedOperationException`.
+
+For more details take a look at the sample topology in storm-starter 
[PersistentWindowingTopology](../examples/storm-starter/src/jvm/org/apache/storm/starter/PersistentWindowingTopology.java)
+which will help you get started.

http://git-wip-us.apache.org/repos/asf/storm/blob/3ba3cabb/examples/storm-starter/src/jvm/org/apache/storm/starter/PersistentWindowingTopology.java
----------------------------------------------------------------------
diff --git 
a/examples/storm-starter/src/jvm/org/apache/storm/starter/PersistentWindowingTopology.java
 
b/examples/storm-starter/src/jvm/org/apache/storm/starter/PersistentWindowingTopology.java
new file mode 100644
index 0000000..566ab69
--- /dev/null
+++ 
b/examples/storm-starter/src/jvm/org/apache/storm/starter/PersistentWindowingTopology.java
@@ -0,0 +1,176 @@
+/*
+ * 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.starter;
+
+import static org.apache.storm.topology.base.BaseWindowedBolt.Duration;
+
+import java.util.Iterator;
+import java.util.Map;
+import java.util.concurrent.TimeUnit;
+import org.apache.storm.Config;
+import org.apache.storm.StormSubmitter;
+import org.apache.storm.starter.spout.RandomIntegerSpout;
+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.topology.OutputFieldsDeclarer;
+import org.apache.storm.topology.TopologyBuilder;
+import org.apache.storm.topology.base.BaseStatefulWindowedBolt;
+import org.apache.storm.tuple.Fields;
+import org.apache.storm.tuple.Tuple;
+import org.apache.storm.tuple.Values;
+import org.apache.storm.windowing.TupleWindow;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+/**
+ * An example that demonstrates the usage of {@link 
org.apache.storm.topology.IStatefulWindowedBolt} with window
+ * persistence.
+ * <p>
+ * The framework automatically checkpoints the tuples in the window along with 
the bolt's state and restores the same
+ * during restarts.
+ * </p>
+ *
+ * <p>
+ * This topology uses 'redis' for state persistence, so you should also start 
a redis instance before deploying.
+ * If you are running in local mode you can just start a redis server locally 
which will be used for storing the state. The default
+ * RedisKeyValueStateProvider parameters can be overridden by setting {@link 
Config#TOPOLOGY_STATE_PROVIDER_CONFIG}, for e.g.
+ * <pre>
+ * {
+ *   "jedisPoolConfig": {
+ *     "host": "redis-server-host",
+ *     "port": 6379,
+ *     "timeout": 2000,
+ *     "database": 0,
+ *     "password": "xyz"
+ *   }
+ * }
+ * </pre>
+ * </p>
+ */
+public class PersistentWindowingTopology {
+    private static final Logger LOG = 
LoggerFactory.getLogger(PersistentWindowingTopology.class);
+
+    // wrapper to hold global and window averages
+    private static class Averages {
+        private final double global;
+        private final double window;
+
+        Averages(double global, double window) {
+            this.global = global;
+            this.window = window;
+        }
+
+        @Override
+        public String toString() {
+            return "Averages{" + "global=" + String.format("%.2f", global) + 
", window=" + String.format("%.2f", window) + '}';
+        }
+    }
+
+    /**
+     * A bolt that uses stateful persistence to store the windows along with 
the state (global avg).
+     */
+    private static class AvgBolt extends 
BaseStatefulWindowedBolt<KeyValueState<String, Pair<Long, Long>>> {
+        private static final String STATE_KEY = "avg";
+
+        private OutputCollector collector;
+        private KeyValueState<String, Pair<Long, Long>> state;
+        private Pair<Long, Long> globalAvg;
+
+        @Override
+        public void prepare(Map<String, Object> topoConf, TopologyContext 
context, OutputCollector collector) {
+            this.collector = collector;
+        }
+
+        @Override
+        public void initState(KeyValueState<String, Pair<Long, Long>> state) {
+            this.state = state;
+            globalAvg = state.get(STATE_KEY, Pair.of(0L, 0L));
+            LOG.info("initState with global avg [" + (double) 
globalAvg.getFirst() / globalAvg.getSecond() + "]");
+        }
+
+        @Override
+        public void execute(TupleWindow window) {
+            int sum = 0;
+            int count = 0;
+            // iterate over tuples in the current window
+            Iterator<Tuple> it = window.getIter();
+            while (it.hasNext()) {
+                Tuple tuple = it.next();
+                sum += tuple.getInteger(0);
+                ++count;
+            }
+            LOG.debug("Count : {}", count);
+            globalAvg = Pair.of(globalAvg.getFirst() + sum, 
globalAvg.getSecond() + count);
+            // update the value in state
+            state.put(STATE_KEY, globalAvg);
+            // emit the averages downstream
+            collector.emit(new Values(new Averages((double) 
globalAvg.getFirst() / globalAvg.getSecond(), (double) sum / count)));
+        }
+
+        @Override
+        public void declareOutputFields(OutputFieldsDeclarer declarer) {
+            declarer.declare(new Fields("avg"));
+        }
+    }
+
+
+    /**
+     * Create and deploy the topology.
+     *
+     * @param args args
+     * @throws Exception exception
+     */
+    public static void main(String[] args) throws Exception {
+        TopologyBuilder builder = new TopologyBuilder();
+
+        // generate random numbers
+        builder.setSpout("spout", new RandomIntegerSpout());
+
+        // emits sliding window and global averages
+        builder.setBolt("avgbolt", new AvgBolt()
+            .withWindow(new Duration(10, TimeUnit.SECONDS), new Duration(2, 
TimeUnit.SECONDS))
+            // persist the window in state
+            .withPersistence()
+            // max number of events to be cached in memory
+            .withMaxEventsInMemory(25000), 1)
+            .shuffleGrouping("spout");
+
+        // print the values to stdout
+        builder.setBolt("printer", (x, y) -> 
System.out.println(x.getValue(0)), 1).shuffleGrouping("avgbolt");
+
+        Config conf = new Config();
+        conf.setDebug(false);
+
+        // checkpoint the state every 5 seconds
+        conf.put(Config.TOPOLOGY_STATE_CHECKPOINT_INTERVAL, 5000);
+
+        // use redis for state persistence
+        conf.put(Config.TOPOLOGY_STATE_PROVIDER, 
"org.apache.storm.redis.state.RedisKeyValueStateProvider");
+
+        String topoName = "test";
+        if (args != null && args.length > 0) {
+            topoName = args[0];
+        }
+        conf.setNumWorkers(1);
+        StormSubmitter.submitTopologyWithProgressBar(topoName, conf, 
builder.createTopology());
+    }
+
+}

http://git-wip-us.apache.org/repos/asf/storm/blob/3ba3cabb/examples/storm-starter/src/jvm/org/apache/storm/starter/SlidingWindowTopology.java
----------------------------------------------------------------------
diff --git 
a/examples/storm-starter/src/jvm/org/apache/storm/starter/SlidingWindowTopology.java
 
b/examples/storm-starter/src/jvm/org/apache/storm/starter/SlidingWindowTopology.java
index 1aa086b..fc7cf4a 100644
--- 
a/examples/storm-starter/src/jvm/org/apache/storm/starter/SlidingWindowTopology.java
+++ 
b/examples/storm-starter/src/jvm/org/apache/storm/starter/SlidingWindowTopology.java
@@ -15,6 +15,7 @@
  * See the License for the specific language governing permissions and
  * limitations under the License.
  */
+
 package org.apache.storm.starter;
 
 import java.util.List;

http://git-wip-us.apache.org/repos/asf/storm/blob/3ba3cabb/examples/storm-starter/src/jvm/org/apache/storm/starter/StatefulTopology.java
----------------------------------------------------------------------
diff --git 
a/examples/storm-starter/src/jvm/org/apache/storm/starter/StatefulTopology.java 
b/examples/storm-starter/src/jvm/org/apache/storm/starter/StatefulTopology.java
index 3372f91..e407ce8 100644
--- 
a/examples/storm-starter/src/jvm/org/apache/storm/starter/StatefulTopology.java
+++ 
b/examples/storm-starter/src/jvm/org/apache/storm/starter/StatefulTopology.java
@@ -15,6 +15,7 @@
  * See the License for the specific language governing permissions and
  * limitations under the License.
  */
+
 package org.apache.storm.starter;
 
 import java.util.Map;
@@ -63,6 +64,7 @@ import org.slf4j.LoggerFactory;
  */
 public class StatefulTopology {
     private static final Logger LOG = 
LoggerFactory.getLogger(StatefulTopology.class);
+
     /**
      * A bolt that uses {@link KeyValueState} to save its state.
      */

http://git-wip-us.apache.org/repos/asf/storm/blob/3ba3cabb/examples/storm-starter/src/jvm/org/apache/storm/starter/StatefulWindowingTopology.java
----------------------------------------------------------------------
diff --git 
a/examples/storm-starter/src/jvm/org/apache/storm/starter/StatefulWindowingTopology.java
 
b/examples/storm-starter/src/jvm/org/apache/storm/starter/StatefulWindowingTopology.java
index 3cf8ffe..30c1ab2 100644
--- 
a/examples/storm-starter/src/jvm/org/apache/storm/starter/StatefulWindowingTopology.java
+++ 
b/examples/storm-starter/src/jvm/org/apache/storm/starter/StatefulWindowingTopology.java
@@ -15,6 +15,7 @@
  * See the License for the specific language governing permissions and
  * limitations under the License.
  */
+
 package org.apache.storm.starter;
 
 import java.util.Map;

http://git-wip-us.apache.org/repos/asf/storm/blob/3ba3cabb/external/storm-hbase/pom.xml
----------------------------------------------------------------------
diff --git a/external/storm-hbase/pom.xml b/external/storm-hbase/pom.xml
index 1175770..0921e6d 100644
--- a/external/storm-hbase/pom.xml
+++ b/external/storm-hbase/pom.xml
@@ -37,7 +37,6 @@
 
     <properties>
         <hdfs.version>${hadoop.version}</hdfs.version>
-        <caffeine.version>2.3.5</caffeine.version>
     </properties>
 
     <dependencies>
@@ -91,7 +90,6 @@
         <dependency>
             <groupId>com.github.ben-manes.caffeine</groupId>
             <artifactId>caffeine</artifactId>
-            <version>${caffeine.version}</version>
         </dependency>
         <dependency>
             <groupId>org.apache.storm</groupId>

http://git-wip-us.apache.org/repos/asf/storm/blob/3ba3cabb/external/storm-hbase/src/main/java/org/apache/storm/hbase/state/HBaseKeyValueStateProvider.java
----------------------------------------------------------------------
diff --git 
a/external/storm-hbase/src/main/java/org/apache/storm/hbase/state/HBaseKeyValueStateProvider.java
 
b/external/storm-hbase/src/main/java/org/apache/storm/hbase/state/HBaseKeyValueStateProvider.java
index 001400a..ce91086 100644
--- 
a/external/storm-hbase/src/main/java/org/apache/storm/hbase/state/HBaseKeyValueStateProvider.java
+++ 
b/external/storm-hbase/src/main/java/org/apache/storm/hbase/state/HBaseKeyValueStateProvider.java
@@ -47,7 +47,7 @@ public class HBaseKeyValueStateProvider implements 
StateProvider {
     @Override
     public State newState(String namespace, Map stormConf, TopologyContext 
context) {
         try {
-            return getHBaseKeyValueState(namespace, stormConf, 
getStateConfig(stormConf));
+            return getHBaseKeyValueState(namespace, stormConf, context, 
getStateConfig(stormConf));
         } catch (Exception ex) {
             LOG.error("Error loading config from storm conf {}", stormConf);
             throw new RuntimeException(ex);
@@ -55,12 +55,12 @@ public class HBaseKeyValueStateProvider implements 
StateProvider {
     }
 
     StateConfig getStateConfig(Map stormConf) throws Exception {
-        StateConfig stateConfig = null;
-        String providerConfig = null;
+        StateConfig stateConfig;
+        String providerConfig;
         ObjectMapper mapper = new ObjectMapper();
         mapper.setVisibility(PropertyAccessor.FIELD, 
JsonAutoDetect.Visibility.ANY);
-        if 
(stormConf.containsKey(org.apache.storm.Config.TOPOLOGY_STATE_PROVIDER_CONFIG)) 
{
-            providerConfig = (String) 
stormConf.get(org.apache.storm.Config.TOPOLOGY_STATE_PROVIDER_CONFIG);
+        if (stormConf.containsKey(Config.TOPOLOGY_STATE_PROVIDER_CONFIG)) {
+            providerConfig = (String) 
stormConf.get(Config.TOPOLOGY_STATE_PROVIDER_CONFIG);
             stateConfig = mapper.readValue(providerConfig, StateConfig.class);
         } else {
             stateConfig = new StateConfig();
@@ -74,7 +74,8 @@ public class HBaseKeyValueStateProvider implements 
StateProvider {
         return stateConfig;
     }
 
-    private HBaseKeyValueState getHBaseKeyValueState(String namespace, Map 
stormConf, StateConfig config) throws Exception {
+    private HBaseKeyValueState getHBaseKeyValueState(String namespace, 
Map<String, Object> stormConf, TopologyContext context,
+                                                     StateConfig config) 
throws Exception {
         Map<String, Object> conf = getHBaseConfigMap(stormConf, 
config.hbaseConfigKey);
         final Configuration hbConfig = getHBaseConfigurationInstance(conf);
 
@@ -85,7 +86,7 @@ public class HBaseKeyValueStateProvider implements 
StateProvider {
         HBaseClient hbaseClient = new HBaseClient(hbaseConfMap, hbConfig, 
config.tableName);
 
         return new HBaseKeyValueState(hbaseClient, config.columnFamily, 
namespace,
-                getKeySerializer(config), getValueSerializer(config));
+                getKeySerializer(stormConf, context, config), 
getValueSerializer(stormConf, context, config));
     }
 
     private Configuration getHBaseConfigurationInstance(Map<String, Object> 
conf) {
@@ -114,28 +115,28 @@ public class HBaseKeyValueStateProvider implements 
StateProvider {
         }
     }
 
-    private Serializer getKeySerializer(StateConfig config) throws Exception {
-        Serializer serializer = null;
+    private Serializer getKeySerializer(Map<String, Object> topoConf, 
TopologyContext context, StateConfig config) throws Exception {
+        Serializer serializer;
         if (config.keySerializerClass != null) {
-            Class<?> klass = (Class<?>) 
Class.forName(config.keySerializerClass);
+            Class<?> klass = Class.forName(config.keySerializerClass);
             serializer = (Serializer) klass.newInstance();
         } else if (config.keyClass != null) {
-            serializer = new 
DefaultStateSerializer(Collections.singletonList(Class.forName(config.keyClass)));
+            serializer = new DefaultStateSerializer(topoConf, context, 
Collections.singletonList(Class.forName(config.keyClass)));
         } else {
-            serializer = new DefaultStateSerializer();
+            serializer = new DefaultStateSerializer(topoConf, context);
         }
         return serializer;
     }
 
-    private Serializer getValueSerializer(StateConfig config) throws Exception 
{
+    private Serializer getValueSerializer(Map<String, Object> topoConf, 
TopologyContext context, StateConfig config) throws Exception {
         Serializer serializer = null;
         if (config.valueSerializerClass != null) {
             Class<?> klass = (Class<?>) 
Class.forName(config.valueSerializerClass);
             serializer = (Serializer) klass.newInstance();
         } else if (config.valueClass != null) {
-            serializer = new 
DefaultStateSerializer(Collections.singletonList(Class.forName(config.valueClass)));
+            serializer = new DefaultStateSerializer(topoConf, context, 
Collections.singletonList(Class.forName(config.valueClass)));
         } else {
-            serializer = new DefaultStateSerializer();
+            serializer = new DefaultStateSerializer(topoConf, context);
         }
         return serializer;
     }

http://git-wip-us.apache.org/repos/asf/storm/blob/3ba3cabb/external/storm-redis/src/main/java/org/apache/storm/redis/state/RedisKeyValueStateProvider.java
----------------------------------------------------------------------
diff --git 
a/external/storm-redis/src/main/java/org/apache/storm/redis/state/RedisKeyValueStateProvider.java
 
b/external/storm-redis/src/main/java/org/apache/storm/redis/state/RedisKeyValueStateProvider.java
index 6281074..01434ab 100644
--- 
a/external/storm-redis/src/main/java/org/apache/storm/redis/state/RedisKeyValueStateProvider.java
+++ 
b/external/storm-redis/src/main/java/org/apache/storm/redis/state/RedisKeyValueStateProvider.java
@@ -17,23 +17,21 @@
  */
 package org.apache.storm.redis.state;
 
+import com.fasterxml.jackson.annotation.JsonAutoDetect;
+import com.fasterxml.jackson.annotation.PropertyAccessor;
+import com.fasterxml.jackson.databind.ObjectMapper;
+import org.apache.storm.Config;
 import org.apache.storm.redis.common.config.JedisClusterConfig;
+import org.apache.storm.redis.common.config.JedisPoolConfig;
 import org.apache.storm.state.DefaultStateSerializer;
 import org.apache.storm.state.Serializer;
 import org.apache.storm.state.State;
 import org.apache.storm.state.StateProvider;
 import org.apache.storm.task.TopologyContext;
-import com.fasterxml.jackson.annotation.JsonAutoDetect;
-import com.fasterxml.jackson.annotation.PropertyAccessor;
-import com.fasterxml.jackson.core.type.TypeReference;
-import com.fasterxml.jackson.databind.ObjectMapper;
-import org.apache.storm.redis.common.config.JedisPoolConfig;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
-import redis.clients.jedis.JedisCluster;
 
 import java.util.Collections;
-import java.util.HashMap;
 import java.util.Map;
 
 /**
@@ -45,7 +43,7 @@ public class RedisKeyValueStateProvider implements 
StateProvider {
     @Override
     public State newState(String namespace, Map<String, Object> topoConf, 
TopologyContext context) {
         try {
-            return getRedisKeyValueState(namespace, getStateConfig(topoConf));
+            return getRedisKeyValueState(namespace, topoConf, context, 
getStateConfig(topoConf));
         } catch (Exception ex) {
             LOG.error("Error loading config from storm conf {}", topoConf);
             throw new RuntimeException(ex);
@@ -53,12 +51,12 @@ public class RedisKeyValueStateProvider implements 
StateProvider {
     }
 
     StateConfig getStateConfig(Map<String, Object> topoConf) throws Exception {
-        StateConfig stateConfig = null;
-        String providerConfig = null;
+        StateConfig stateConfig;
+        String providerConfig;
         ObjectMapper mapper = new ObjectMapper();
         mapper.setVisibility(PropertyAccessor.FIELD, 
JsonAutoDetect.Visibility.ANY);
-        if 
(topoConf.containsKey(org.apache.storm.Config.TOPOLOGY_STATE_PROVIDER_CONFIG)) {
-            providerConfig = (String) 
topoConf.get(org.apache.storm.Config.TOPOLOGY_STATE_PROVIDER_CONFIG);
+        if (topoConf.containsKey(Config.TOPOLOGY_STATE_PROVIDER_CONFIG)) {
+            providerConfig = (String) 
topoConf.get(Config.TOPOLOGY_STATE_PROVIDER_CONFIG);
             stateConfig = mapper.readValue(providerConfig, StateConfig.class);
         } else {
             stateConfig = new StateConfig();
@@ -66,7 +64,8 @@ public class RedisKeyValueStateProvider implements 
StateProvider {
         return stateConfig;
     }
 
-    private RedisKeyValueState getRedisKeyValueState(String namespace, 
StateConfig config) throws Exception {
+    private RedisKeyValueState getRedisKeyValueState(String namespace, 
Map<String, Object> topoConf, TopologyContext context,
+                                                     StateConfig config) 
throws Exception {
         JedisPoolConfig jedisPoolConfig = getJedisPoolConfig(config);
         JedisClusterConfig jedisClusterConfig = getJedisClusterConfig(config);
 
@@ -75,34 +74,36 @@ public class RedisKeyValueStateProvider implements 
StateProvider {
         }
 
         if (jedisPoolConfig != null) {
-            return new RedisKeyValueState(namespace, jedisPoolConfig, 
getKeySerializer(config), getValueSerializer(config));
+            return new RedisKeyValueState(namespace, jedisPoolConfig,
+                getKeySerializer(topoConf, context, config), 
getValueSerializer(topoConf, context, config));
         } else {
-            return new RedisKeyValueState(namespace, jedisClusterConfig, 
getKeySerializer(config), getValueSerializer(config));
+            return new RedisKeyValueState(namespace, jedisClusterConfig,
+                getKeySerializer(topoConf, context, config), 
getValueSerializer(topoConf, context, config));
         }
     }
 
-    private Serializer getKeySerializer(StateConfig config) throws Exception {
-        Serializer serializer = null;
+    private Serializer getKeySerializer(Map<String, Object> topoConf, 
TopologyContext context, StateConfig config) throws Exception {
+        Serializer serializer;
         if (config.keySerializerClass != null) {
             Class<?> klass = (Class<?>) 
Class.forName(config.keySerializerClass);
             serializer = (Serializer) klass.newInstance();
         } else if (config.keyClass != null) {
-            serializer = new 
DefaultStateSerializer(Collections.singletonList(Class.forName(config.keyClass)));
+            serializer = new DefaultStateSerializer(topoConf, context, 
Collections.singletonList(Class.forName(config.keyClass)));
         } else {
-            serializer = new DefaultStateSerializer();
+            serializer = new DefaultStateSerializer(topoConf, context);
         }
         return serializer;
     }
 
-    private Serializer getValueSerializer(StateConfig config) throws Exception 
{
-        Serializer serializer = null;
+    private Serializer getValueSerializer(Map<String, Object> topoConf, 
TopologyContext context, StateConfig config) throws Exception {
+        Serializer serializer;
         if (config.valueSerializerClass != null) {
             Class<?> klass = (Class<?>) 
Class.forName(config.valueSerializerClass);
             serializer = (Serializer) klass.newInstance();
         } else if (config.valueClass != null) {
-            serializer = new 
DefaultStateSerializer(Collections.singletonList(Class.forName(config.valueClass)));
+            serializer = new DefaultStateSerializer(topoConf, context, 
Collections.singletonList(Class.forName(config.valueClass)));
         } else {
-            serializer = new DefaultStateSerializer();
+            serializer = new DefaultStateSerializer(topoConf, context);
         }
         return serializer;
     }

http://git-wip-us.apache.org/repos/asf/storm/blob/3ba3cabb/pom.xml
----------------------------------------------------------------------
diff --git a/pom.xml b/pom.xml
index 319854b..de79ef2 100644
--- a/pom.xml
+++ b/pom.xml
@@ -313,6 +313,7 @@
         <dropwizard.version>1.1.2</dropwizard.version>
         <j2html.version>1.0.0</j2html.version>
         <jool.version>0.9.12</jool.version>
+        <caffeine.version>2.3.5</caffeine.version>
 
         <!-- see intellij profile below... This fixes an annoyance with 
intellij -->
         <provided.scope>provided</provided.scope>
@@ -1051,6 +1052,11 @@
                 <artifactId>sysout-over-slf4j</artifactId>
                 <version>${sysout-over-slf4j.version}</version>
             </dependency>
+            <dependency>
+                <groupId>com.github.ben-manes.caffeine</groupId>
+                <artifactId>caffeine</artifactId>
+                <version>${caffeine.version}</version>
+            </dependency>
         </dependencies>
     </dependencyManagement>
 

http://git-wip-us.apache.org/repos/asf/storm/blob/3ba3cabb/storm-client/src/jvm/org/apache/storm/Config.java
----------------------------------------------------------------------
diff --git a/storm-client/src/jvm/org/apache/storm/Config.java 
b/storm-client/src/jvm/org/apache/storm/Config.java
index 4a654fd..56e33c2 100644
--- a/storm-client/src/jvm/org/apache/storm/Config.java
+++ b/storm-client/src/jvm/org/apache/storm/Config.java
@@ -342,6 +342,12 @@ public class Config extends HashMap<String, Object> {
     public static final String TOPOLOGY_SKIP_MISSING_KRYO_REGISTRATIONS= 
"topology.skip.missing.kryo.registrations";
 
     /**
+     * List of classes to register during state serialization
+     */
+    @isStringList
+    public static final String TOPOLOGY_STATE_KRYO_REGISTER = 
"topology.state.kryo.register";
+
+    /**
      * A list of classes implementing IMetricsConsumer (See storm.yaml.example 
for exact config format).
      * Each listed class will be routed all the metrics data generated by the 
storm metrics API.
      * Each listed class maps 1:1 to a system bolt named 
__metrics_ClassName#N, and it's parallelism is configurable.

http://git-wip-us.apache.org/repos/asf/storm/blob/3ba3cabb/storm-client/src/jvm/org/apache/storm/serialization/SerializationFactory.java
----------------------------------------------------------------------
diff --git 
a/storm-client/src/jvm/org/apache/storm/serialization/SerializationFactory.java 
b/storm-client/src/jvm/org/apache/storm/serialization/SerializationFactory.java
index 1faa926..d6c54a3 100644
--- 
a/storm-client/src/jvm/org/apache/storm/serialization/SerializationFactory.java
+++ 
b/storm-client/src/jvm/org/apache/storm/serialization/SerializationFactory.java
@@ -83,31 +83,11 @@ public class SerializationFactory {
             }
         }
 
-        Map<String, String> registrations = normalizeKryoRegister(conf);
-
         kryoFactory.preRegister(k, conf);
 
         boolean skipMissing = (Boolean) 
conf.get(Config.TOPOLOGY_SKIP_MISSING_KRYO_REGISTRATIONS);
-        for(Map.Entry<String, String> entry: registrations.entrySet()) {
-            String serializerClassName = entry.getValue();
-            try {
-                Class klass = Class.forName(entry.getKey());
-                Class serializerClass = null;
-                if(serializerClassName!=null)
-                    serializerClass = Class.forName(serializerClassName);
-                if(serializerClass == null) {
-                    k.register(klass);
-                } else {
-                    k.register(klass, resolveSerializerInstance(k, klass, 
serializerClass, conf));
-                }
-            } catch (ClassNotFoundException e) {
-                if(skipMissing) {
-                    LOG.info("Could not find serialization or class for " + 
serializerClassName + ". Skipping registration...");
-                } else {
-                    throw new RuntimeException(e);
-                }
-            }
-        }
+
+        register(k, conf.get(Config.TOPOLOGY_KRYO_REGISTER), conf, 
skipMissing);
 
         kryoFactory.postRegister(k, conf);
 
@@ -136,6 +116,34 @@ public class SerializationFactory {
         return k;
     }
 
+    public static void register(Kryo k, List<String> classesToRegister) {
+        register(k, classesToRegister, Collections.emptyMap(), true);
+    }
+
+    public static void register(Kryo k, Object kryoRegistrations, Map<String, 
Object> conf, boolean skipMissing) {
+        Map<String, String> registrations = 
normalizeKryoRegister(kryoRegistrations);
+        for(Map.Entry<String, String> entry: registrations.entrySet()) {
+            String serializerClassName = entry.getValue();
+            try {
+                Class klass = Class.forName(entry.getKey());
+                Class serializerClass = null;
+                if(serializerClassName!=null)
+                    serializerClass = Class.forName(serializerClassName);
+                if(serializerClass == null) {
+                    k.register(klass);
+                } else {
+                    k.register(klass, resolveSerializerInstance(k, klass, 
serializerClass, conf));
+                }
+            } catch (ClassNotFoundException e) {
+                if(skipMissing) {
+                    LOG.info("Could not find serialization or class for " + 
serializerClassName + ". Skipping registration...");
+                } else {
+                    throw new RuntimeException(e);
+                }
+            }
+        }
+    }
+
     public static class IdDictionary {
         Map<String, Map<String, Integer>> streamNametoId = new HashMap<>();
         Map<String, Map<Integer, String>> streamIdToName = new HashMap<>();
@@ -223,15 +231,14 @@ public class SerializationFactory {
         }
     }
 
-    private static Map<String, String> normalizeKryoRegister(Map<String, 
Object> conf) {
+    private static Map<String, String> normalizeKryoRegister(Object 
kryoRegistrations) {
         // TODO: de-duplicate this logic with the code in nimbus
-        Object res = conf.get(Config.TOPOLOGY_KRYO_REGISTER);
-        if(res==null) return new TreeMap<>();
+        if(kryoRegistrations==null) return new TreeMap<>();
         Map<String, String> ret = new HashMap<>();
-        if(res instanceof Map) {
-            ret = (Map<String, String>) res;
+        if(kryoRegistrations instanceof Map) {
+            ret = (Map<String, String>) kryoRegistrations;
         } else {
-            for(Object o: (List) res) {
+            for(Object o: (List) kryoRegistrations) {
                 if(o instanceof Map) {
                     ret.putAll((Map) o);
                 } else {

http://git-wip-us.apache.org/repos/asf/storm/blob/3ba3cabb/storm-client/src/jvm/org/apache/storm/state/DefaultStateEncoder.java
----------------------------------------------------------------------
diff --git 
a/storm-client/src/jvm/org/apache/storm/state/DefaultStateEncoder.java 
b/storm-client/src/jvm/org/apache/storm/state/DefaultStateEncoder.java
index 5530a6e..b1c9f14 100644
--- a/storm-client/src/jvm/org/apache/storm/state/DefaultStateEncoder.java
+++ b/storm-client/src/jvm/org/apache/storm/state/DefaultStateEncoder.java
@@ -18,7 +18,8 @@
 
 package org.apache.storm.state;
 
-import com.google.common.base.Optional;
+
+import java.util.Optional;
 
 /**
  * Default state encoder class for encoding/decoding key values. This class 
assumes encoded types of key and value are
@@ -28,7 +29,7 @@ public class DefaultStateEncoder<K, V> implements 
StateEncoder<K, V, byte[], byt
 
     public static final Serializer<Optional<byte[]>> internalValueSerializer = 
new DefaultStateSerializer<>();
 
-    public static final byte[] TOMBSTONE = 
internalValueSerializer.serialize(Optional.<byte[]>absent());
+    public static final byte[] TOMBSTONE = 
internalValueSerializer.serialize(Optional.<byte[]>empty());
 
     private final Serializer<K> keySerializer;
     private final Serializer<V> valueSerializer;

http://git-wip-us.apache.org/repos/asf/storm/blob/3ba3cabb/storm-client/src/jvm/org/apache/storm/state/DefaultStateSerializer.java
----------------------------------------------------------------------
diff --git 
a/storm-client/src/jvm/org/apache/storm/state/DefaultStateSerializer.java 
b/storm-client/src/jvm/org/apache/storm/state/DefaultStateSerializer.java
index 45e2bfc..bb61921 100644
--- a/storm-client/src/jvm/org/apache/storm/state/DefaultStateSerializer.java
+++ b/storm-client/src/jvm/org/apache/storm/state/DefaultStateSerializer.java
@@ -20,20 +20,42 @@ package org.apache.storm.state;
 import com.esotericsoftware.kryo.Kryo;
 import com.esotericsoftware.kryo.io.Input;
 import com.esotericsoftware.kryo.io.Output;
+import org.apache.storm.Config;
+import org.apache.storm.serialization.KryoTupleDeserializer;
+import org.apache.storm.serialization.KryoTupleSerializer;
+import org.apache.storm.serialization.SerializationFactory;
+import org.apache.storm.task.TopologyContext;
+import org.apache.storm.tuple.TupleImpl;
 import org.objenesis.strategy.StdInstantiatorStrategy;
 
+import java.util.ArrayList;
 import java.util.Collections;
 import java.util.List;
+import java.util.Map;
+import java.util.Optional;
+import java.util.stream.Collectors;
 
 /**
  * A default implementation that uses Kryo to serialize and de-serialize
  * the state.
  */
 public class DefaultStateSerializer<T> implements Serializer<T> {
+    private final TopologyContext context;
+    private final Map<String, Object> topoConf;
+    private final List<String> registrations = new ArrayList<>();
+
     private final ThreadLocal<Kryo> kryo = new ThreadLocal<Kryo>() {
         @Override
         protected Kryo initialValue() {
             Kryo obj = new Kryo();
+            if (context != null && topoConf != null) {
+                KryoTupleSerializer ser = new KryoTupleSerializer(topoConf, 
context);
+                KryoTupleDeserializer deser = new 
KryoTupleDeserializer(topoConf, context);
+                obj.register(TupleImpl.class, new TupleSerializer(ser, deser));
+            }
+            if (!registrations.isEmpty()) {
+                SerializationFactory.register(obj, registrations);
+            }
             obj.setInstantiatorStrategy(new 
Kryo.DefaultInstantiatorStrategy(new StdInstantiatorStrategy()));
             return obj;
         }
@@ -52,14 +74,22 @@ public class DefaultStateSerializer<T> implements 
Serializer<T> {
      *
      * @param classesToRegister the classes to register.
      */
-    public DefaultStateSerializer(List<Class<?>> classesToRegister) {
-        for (Class<?> klazz : classesToRegister) {
-            kryo.get().register(klazz);
-        }
+    public DefaultStateSerializer(Map<String, Object> topoConf, 
TopologyContext context, List<Class<?>> classesToRegister) {
+        this.context = context;
+        this.topoConf = topoConf;
+        
registrations.addAll(classesToRegister.stream().map(Class::getName).collect(Collectors.toSet()));
+        // other classes from config
+        registrations.addAll((List<String>) 
topoConf.getOrDefault(Config.TOPOLOGY_STATE_KRYO_REGISTER, 
Collections.emptyList()));
+        // defaults
+        registrations.add(Optional.class.getName());
+    }
+
+    public DefaultStateSerializer(Map<String, Object> topoConf, 
TopologyContext context) {
+        this(topoConf, context, Collections.emptyList());
     }
 
     public DefaultStateSerializer() {
-        this(Collections.<Class<?>>emptyList());
+        this(Collections.emptyMap(), null);
     }
 
     @Override
@@ -74,4 +104,28 @@ public class DefaultStateSerializer<T> implements 
Serializer<T> {
         Input input = new Input(b);
         return (T) kryo.get().readClassAndObject(input);
     }
+
+    private static class TupleSerializer extends 
com.esotericsoftware.kryo.Serializer<TupleImpl> {
+        private final KryoTupleSerializer tupleSerializer;
+        private final KryoTupleDeserializer tupleDeserializer;
+
+        TupleSerializer(KryoTupleSerializer tupleSerializer, 
KryoTupleDeserializer tupleDeserializer) {
+            this.tupleSerializer = tupleSerializer;
+            this.tupleDeserializer = tupleDeserializer;
+        }
+
+        @Override
+        public void write(Kryo kryo, Output output, TupleImpl tuple) {
+            byte[] bytes = tupleSerializer.serialize(tuple);
+            output.writeInt(bytes.length);
+            output.write(bytes);
+        }
+
+        @Override
+        public TupleImpl read(Kryo kryo, Input input, Class<TupleImpl> type) {
+            int length = input.readInt();
+            byte[] bytes = input.readBytes(length);
+            return (TupleImpl) tupleDeserializer.deserialize(bytes);
+        }
+    }
 }

http://git-wip-us.apache.org/repos/asf/storm/blob/3ba3cabb/storm-client/src/jvm/org/apache/storm/topology/IStatefulWindowedBolt.java
----------------------------------------------------------------------
diff --git 
a/storm-client/src/jvm/org/apache/storm/topology/IStatefulWindowedBolt.java 
b/storm-client/src/jvm/org/apache/storm/topology/IStatefulWindowedBolt.java
index cd7baf8..d8fb0e7 100644
--- a/storm-client/src/jvm/org/apache/storm/topology/IStatefulWindowedBolt.java
+++ b/storm-client/src/jvm/org/apache/storm/topology/IStatefulWindowedBolt.java
@@ -15,12 +15,32 @@
  * See the License for the specific language governing permissions and
  * limitations under the License.
  */
+
 package org.apache.storm.topology;
 
 import org.apache.storm.state.State;
 
 /**
- * A windowed bolt abstraction for supporting windowing operation with state
+ * A windowed bolt abstraction for supporting windowing operation with state.
  */
 public interface IStatefulWindowedBolt<T extends State> extends 
IStatefulComponent<T>, IWindowedBolt {
+    /**
+     * If the stateful windowed bolt should have its windows persisted in 
state and maintain a subset
+     * of events in memory.
+     * <p>
+     * The default is to keep all the window events in memory.
+     * </p>
+     *
+     * @return true if the windows should be persisted
+     */
+    default boolean isPersistent() {
+        return false;
+    }
+
+    /**
+     * The maximum number of window events to keep in memory.
+     */
+    default long maxEventsInMemory() {
+        return 1_000_000L; // default
+    }
 }

http://git-wip-us.apache.org/repos/asf/storm/blob/3ba3cabb/storm-client/src/jvm/org/apache/storm/topology/PersistentWindowedBoltExecutor.java
----------------------------------------------------------------------
diff --git 
a/storm-client/src/jvm/org/apache/storm/topology/PersistentWindowedBoltExecutor.java
 
b/storm-client/src/jvm/org/apache/storm/topology/PersistentWindowedBoltExecutor.java
new file mode 100644
index 0000000..d2b6516
--- /dev/null
+++ 
b/storm-client/src/jvm/org/apache/storm/topology/PersistentWindowedBoltExecutor.java
@@ -0,0 +1,256 @@
+/**
+ * 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 java.util.ArrayList;
+import java.util.Deque;
+import java.util.HashMap;
+import java.util.Iterator;
+import java.util.LinkedList;
+import java.util.List;
+import java.util.Map;
+import java.util.Optional;
+import java.util.concurrent.ConcurrentLinkedQueue;
+import java.util.concurrent.atomic.AtomicInteger;
+import java.util.function.Supplier;
+
+import org.apache.storm.Config;
+import org.apache.storm.state.KeyValueState;
+import org.apache.storm.state.State;
+import org.apache.storm.state.StateFactory;
+import org.apache.storm.task.OutputCollector;
+import org.apache.storm.task.TopologyContext;
+import org.apache.storm.topology.base.BaseWindowedBolt;
+import org.apache.storm.tuple.Tuple;
+import org.apache.storm.windowing.DefaultEvictionContext;
+import org.apache.storm.windowing.EventImpl;
+import org.apache.storm.windowing.WindowLifecycleListener;
+import org.apache.storm.windowing.persistence.WindowState;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import static 
org.apache.storm.windowing.persistence.WindowState.WindowPartition;
+
+/**
+ * Wraps a {@link IStatefulWindowedBolt} and handles the execution. Uses state 
and the underlying
+ * checkpointing mechanisms to save the tuples in window to state. The tuples 
are also kept in-memory
+ * by transparently caching the window partitions and checkpointing them as 
needed.
+ */
+public class PersistentWindowedBoltExecutor<T extends State> extends 
WindowedBoltExecutor implements IStatefulBolt<T> {
+    private static final Logger LOG = 
LoggerFactory.getLogger(PersistentWindowedBoltExecutor.class);
+    private final IStatefulWindowedBolt<T> statefulWindowedBolt;
+    private transient OutputCollector outputCollector;
+    private transient WindowState<Tuple> state;
+    private transient boolean stateInitialized;
+    private transient boolean prePrepared;
+    private transient KeyValueState<String, Optional<?>> windowSystemState;
+
+    public PersistentWindowedBoltExecutor(IStatefulWindowedBolt<T> bolt) {
+        super(bolt);
+        statefulWindowedBolt = bolt;
+    }
+
+    @Override
+    public void prepare(Map<String, Object> topoConf, TopologyContext context, 
OutputCollector collector) {
+        List<String> registrations = (List<String>) 
topoConf.getOrDefault(Config.TOPOLOGY_STATE_KRYO_REGISTER, new ArrayList<>());
+        registrations.add(ConcurrentLinkedQueue.class.getName());
+        registrations.add(LinkedList.class.getName());
+        registrations.add(AtomicInteger.class.getName());
+        registrations.add(EventImpl.class.getName());
+        registrations.add(WindowPartition.class.getName());
+        registrations.add(DefaultEvictionContext.class.getName());
+        topoConf.put(Config.TOPOLOGY_STATE_KRYO_REGISTER, registrations);
+        prepare(topoConf, context, collector, getWindowState(topoConf, 
context), getPartitionState(topoConf, context),
+            getWindowSystemState(topoConf, context));
+    }
+
+    // package access for unit tests
+    void prepare(Map<String, Object> topoConf, TopologyContext context, 
OutputCollector collector,
+                 KeyValueState<Long, WindowPartition<Tuple>> windowState,
+                 KeyValueState<String, Deque<Long>> partitionState,
+                 KeyValueState<String, Optional<?>> windowSystemState) {
+        outputCollector = collector;
+        this.windowSystemState = windowSystemState;
+        state = new WindowState<>(windowState, partitionState, 
windowSystemState, this::getState,
+            statefulWindowedBolt.maxEventsInMemory());
+        doPrepare(topoConf, context, new NoAckOutputCollector(collector), 
state, true);
+        restoreWindowSystemState();
+    }
+
+    private void restoreWindowSystemState() {
+        Map<String, Optional<?>> map = new HashMap<>();
+        for (Map.Entry<String, Optional<?>> entry : windowSystemState) {
+            map.put(entry.getKey(), entry.getValue());
+        }
+        restoreState(map);
+    }
+
+    @Override
+    protected void validate(Map<String, Object> topoConf,
+                            BaseWindowedBolt.Count windowLengthCount,
+                            BaseWindowedBolt.Duration windowLengthDuration,
+                            BaseWindowedBolt.Count slidingIntervalCount,
+                            BaseWindowedBolt.Duration slidingIntervalDuration) 
{
+        if (windowLengthCount == null && windowLengthDuration == null) {
+            throw new IllegalArgumentException("Window length is not 
specified");
+        }
+        int interval = getCheckpointIntervalMillis(topoConf);
+        int timeout = getTopologyTimeoutMillis(topoConf);
+        if (interval > timeout) {
+            throw new 
IllegalArgumentException(Config.TOPOLOGY_STATE_CHECKPOINT_INTERVAL + interval
+                + " is more than " + Config.TOPOLOGY_MESSAGE_TIMEOUT_SECS
+                + " value " + timeout);
+        }
+    }
+
+    private int getCheckpointIntervalMillis(Map<String, Object> topoConf) {
+        int checkpointInterval = Integer.MAX_VALUE;
+        if (topoConf.get(Config.TOPOLOGY_STATE_CHECKPOINT_INTERVAL) != null) {
+            checkpointInterval = ((Number) 
topoConf.get(Config.TOPOLOGY_STATE_CHECKPOINT_INTERVAL)).intValue();
+        }
+        return checkpointInterval;
+    }
+
+    @Override
+    protected void start() {
+        if (stateInitialized) {
+            super.start();
+        } else {
+            LOG.debug("Will invoke start after state is initialized");
+        }
+    }
+
+    @Override
+    public void execute(Tuple input) {
+        if (!stateInitialized) {
+            throw new IllegalStateException("execute invoked before initState 
with input tuple " + input);
+        }
+        super.execute(input);
+        // StatefulBoltExecutor does the actual ack when the state is saved.
+        outputCollector.ack(input);
+    }
+
+    @Override
+    public void initState(T state) {
+        if (stateInitialized) {
+            String msg = "initState invoked when the state is already 
initialized";
+            LOG.warn(msg);
+            throw new IllegalStateException(msg);
+        } else {
+            statefulWindowedBolt.initState(state);
+            stateInitialized = true;
+            start();
+        }
+    }
+
+    @Override
+    public void prePrepare(long txid) {
+        if (stateInitialized) {
+            LOG.debug("Prepare streamState, txid {}", txid);
+            statefulWindowedBolt.prePrepare(txid);
+            state.prepareCommit(txid);
+            prePrepared = true;
+        } else {
+            String msg = "Cannot prepare before initState";
+            LOG.warn(msg);
+            throw new IllegalStateException(msg);
+        }
+    }
+
+    @Override
+    public void preCommit(long txid) {
+        // preCommit can be invoked during recovery before the state is 
initialized
+        if (prePrepared || !stateInitialized) {
+            LOG.debug("Commit streamState, txid {}", txid);
+            statefulWindowedBolt.preCommit(txid);
+            state.commit(txid);
+        } else {
+            String msg = "preCommit before prePrepare in initialized state";
+            LOG.warn(msg);
+            throw new IllegalStateException(msg);
+        }
+    }
+
+    @Override
+    public void preRollback() {
+        LOG.debug("Rollback streamState, stateInitialized {}", 
stateInitialized);
+        statefulWindowedBolt.preRollback();
+        state.rollback(stateInitialized);
+        if (stateInitialized) {
+            restoreWindowSystemState();
+        }
+    }
+
+    @Override
+    protected WindowLifecycleListener<Tuple> newWindowLifecycleListener() {
+        return new WindowLifecycleListener<Tuple>() {
+            @Override
+            public void onExpiry(List<Tuple> events) {
+                /*
+                 * NO-OP: the events are ack-ed in execute
+                 */
+            }
+
+            @Override
+            public void onActivation(Supplier<Iterator<Tuple>> eventsIt,
+                                     Supplier<Iterator<Tuple>> newEventsIt,
+                                     Supplier<Iterator<Tuple>> expiredIt,
+                                     Long timestamp) {
+                /*
+                 * Here we don't set the tuples in windowedOutputCollector's 
context and emit un-anchored.
+                 * The checkpoint tuple will trigger a checkpoint in the 
receiver with the emitted tuples.
+                 */
+                boltExecute(eventsIt, newEventsIt, expiredIt, timestamp);
+                state.clearIteratorPins();
+            }
+        };
+    }
+
+    private KeyValueState<Long, WindowPartition<Tuple>> 
getWindowState(Map<String, Object> topoConf, TopologyContext context) {
+        String namespace = context.getThisComponentId() + "-" + 
context.getThisTaskId() + "-window";
+        return (KeyValueState<Long, WindowPartition<Tuple>>) 
StateFactory.getState(namespace, topoConf, context);
+    }
+
+    private KeyValueState<String, Deque<Long>> getPartitionState(Map<String, 
Object> topoConf, TopologyContext context) {
+        String namespace = context.getThisComponentId() + "-" + 
context.getThisTaskId() + "-window-partitions";
+        return (KeyValueState<String, Deque<Long>>) 
StateFactory.getState(namespace, topoConf, context);
+    }
+
+    private KeyValueState<String, Optional<?>> 
getWindowSystemState(Map<String, Object> topoConf, TopologyContext context) {
+        String namespace = context.getThisComponentId() + "-" + 
context.getThisTaskId() + "-window-systemstate";
+        return (KeyValueState<String, Optional<?>>) 
StateFactory.getState(namespace, topoConf, context);
+    }
+
+    /**
+     * Creates an {@link OutputCollector} wrapper that ignores acks.
+     * The {@link PersistentWindowedBoltExecutor} acks the tuples in execute 
and
+     * this is to prevent double ack-ing
+     */
+    private static class NoAckOutputCollector extends OutputCollector {
+
+        public NoAckOutputCollector(OutputCollector delegate) {
+            super(delegate);
+        }
+
+        @Override
+        public void ack(Tuple input) {
+            // NOOP
+        }
+    }
+}

http://git-wip-us.apache.org/repos/asf/storm/blob/3ba3cabb/storm-client/src/jvm/org/apache/storm/topology/StatefulBoltExecutor.java
----------------------------------------------------------------------
diff --git 
a/storm-client/src/jvm/org/apache/storm/topology/StatefulBoltExecutor.java 
b/storm-client/src/jvm/org/apache/storm/topology/StatefulBoltExecutor.java
index 9efb6c9..9394c36 100644
--- a/storm-client/src/jvm/org/apache/storm/topology/StatefulBoltExecutor.java
+++ b/storm-client/src/jvm/org/apache/storm/topology/StatefulBoltExecutor.java
@@ -36,9 +36,9 @@ import java.util.concurrent.ConcurrentLinkedQueue;
 
 import static org.apache.storm.spout.CheckPointState.Action;
 import static org.apache.storm.spout.CheckPointState.Action.COMMIT;
+import static org.apache.storm.spout.CheckPointState.Action.INITSTATE;
 import static org.apache.storm.spout.CheckPointState.Action.PREPARE;
 import static org.apache.storm.spout.CheckPointState.Action.ROLLBACK;
-import static org.apache.storm.spout.CheckPointState.Action.INITSTATE;
 /**
  * Wraps a {@link IStatefulBolt} and manages the state of the bolt.
  */
@@ -123,6 +123,11 @@ public class StatefulBoltExecutor<T extends State> extends 
BaseStatefulBoltExecu
                 }
                 pendingTuples.clear();
             } else {
+                /*
+                 * If a worker crashes, the states of all workers are rolled 
back and an initState message is sent across
+                 * the topology so that crashed workers can initialize their 
state.
+                 * The bolts that have their state already initialized need 
not be re-initialized.
+                 */
                 LOG.debug("Bolt state is already initialized, ignoring tuple 
{}, action {}, txid {}",
                           checkpointTuple, action, txid);
             }

http://git-wip-us.apache.org/repos/asf/storm/blob/3ba3cabb/storm-client/src/jvm/org/apache/storm/topology/TopologyBuilder.java
----------------------------------------------------------------------
diff --git 
a/storm-client/src/jvm/org/apache/storm/topology/TopologyBuilder.java 
b/storm-client/src/jvm/org/apache/storm/topology/TopologyBuilder.java
index 4cdc08b..99ef9cd 100644
--- a/storm-client/src/jvm/org/apache/storm/topology/TopologyBuilder.java
+++ b/storm-client/src/jvm/org/apache/storm/topology/TopologyBuilder.java
@@ -18,7 +18,16 @@
 package org.apache.storm.topology;
 
 import org.apache.storm.Config;
-import org.apache.storm.generated.*;
+import org.apache.storm.generated.Bolt;
+import org.apache.storm.generated.ComponentCommon;
+import org.apache.storm.generated.ComponentObject;
+import org.apache.storm.generated.GlobalStreamId;
+import org.apache.storm.generated.Grouping;
+import org.apache.storm.generated.NullStruct;
+import org.apache.storm.generated.SpoutSpec;
+import org.apache.storm.generated.StateSpoutSpec;
+import org.apache.storm.generated.StormTopology;
+import org.apache.storm.generated.SharedMemory;
 import org.apache.storm.grouping.CustomStreamGrouping;
 import org.apache.storm.grouping.PartialKeyGrouping;
 import org.apache.storm.hooks.IWorkerHook;
@@ -35,6 +44,7 @@ import org.apache.storm.task.TopologyContext;
 import org.apache.storm.tuple.Fields;
 import org.apache.storm.tuple.Tuple;
 import org.apache.storm.utils.Utils;
+import org.apache.storm.windowing.TupleWindow;
 import org.json.simple.JSONValue;
 import org.json.simple.parser.ParseException;
 
@@ -47,7 +57,6 @@ import java.util.List;
 import java.util.Map;
 import java.util.Set;
 
-import org.apache.storm.windowing.TupleWindow;
 import static org.apache.storm.spout.CheckpointSpout.CHECKPOINT_COMPONENT_ID;
 import static org.apache.storm.spout.CheckpointSpout.CHECKPOINT_STREAM_ID;
 
@@ -323,7 +332,13 @@ public class TopologyBuilder {
      */
     public <T extends State> BoltDeclarer setBolt(String id, 
IStatefulWindowedBolt<T> bolt, Number parallelism_hint) throws 
IllegalArgumentException {
         hasStatefulBolt = true;
-        return setBolt(id, new StatefulBoltExecutor<T>(new 
StatefulWindowedBoltExecutor<T>(bolt)), parallelism_hint);
+        IStatefulBolt<T> executor;
+        if (bolt.isPersistent()) {
+            executor = new PersistentWindowedBoltExecutor<>(bolt);
+        } else {
+            executor = new StatefulWindowedBoltExecutor<T>(bolt);
+        }
+        return setBolt(id, new StatefulBoltExecutor<T>(executor), 
parallelism_hint);
     }
 
     /**

http://git-wip-us.apache.org/repos/asf/storm/blob/3ba3cabb/storm-client/src/jvm/org/apache/storm/topology/WindowedBoltExecutor.java
----------------------------------------------------------------------
diff --git 
a/storm-client/src/jvm/org/apache/storm/topology/WindowedBoltExecutor.java 
b/storm-client/src/jvm/org/apache/storm/topology/WindowedBoltExecutor.java
index 0efe557..c54ec3c 100644
--- a/storm-client/src/jvm/org/apache/storm/topology/WindowedBoltExecutor.java
+++ b/storm-client/src/jvm/org/apache/storm/topology/WindowedBoltExecutor.java
@@ -23,18 +23,20 @@ import org.apache.storm.spout.CheckpointSpout;
 import org.apache.storm.task.IOutputCollector;
 import org.apache.storm.task.OutputCollector;
 import org.apache.storm.task.TopologyContext;
-import org.apache.storm.topology.base.BaseWindowedBolt;
 import org.apache.storm.tuple.Fields;
 import org.apache.storm.tuple.Tuple;
 import org.apache.storm.tuple.Values;
 import org.apache.storm.windowing.CountEvictionPolicy;
 import org.apache.storm.windowing.CountTriggerPolicy;
+import org.apache.storm.windowing.Event;
 import org.apache.storm.windowing.EvictionPolicy;
+import org.apache.storm.windowing.StatefulWindowManager;
 import org.apache.storm.windowing.TimeEvictionPolicy;
 import org.apache.storm.windowing.TimeTriggerPolicy;
 import org.apache.storm.windowing.TimestampExtractor;
 import org.apache.storm.windowing.TriggerPolicy;
 import org.apache.storm.windowing.TupleWindowImpl;
+import org.apache.storm.windowing.TupleWindowIterImpl;
 import org.apache.storm.windowing.WaterMarkEventGenerator;
 import org.apache.storm.windowing.WatermarkCountEvictionPolicy;
 import org.apache.storm.windowing.WatermarkCountTriggerPolicy;
@@ -45,11 +47,17 @@ import org.apache.storm.windowing.WindowManager;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
+import java.util.Collection;
 import java.util.HashSet;
+import java.util.Iterator;
 import java.util.List;
 import java.util.Map;
+import java.util.Objects;
+import java.util.Optional;
 import java.util.Set;
+import java.util.concurrent.ConcurrentLinkedQueue;
 import java.util.concurrent.TimeUnit;
+import java.util.function.Supplier;
 
 import static org.apache.storm.topology.base.BaseWindowedBolt.Count;
 import static org.apache.storm.topology.base.BaseWindowedBolt.Duration;
@@ -69,8 +77,8 @@ public class WindowedBoltExecutor implements IRichBolt {
     private transient int maxLagMs;
     private TimestampExtractor timestampExtractor;
     private transient String lateTupleStream;
-    private transient TriggerPolicy<Tuple> triggerPolicy;
-    private transient EvictionPolicy<Tuple> evictionPolicy;
+    private transient TriggerPolicy<Tuple, ?> triggerPolicy;
+    private transient EvictionPolicy<Tuple,?> evictionPolicy;
     private transient Duration windowLengthDuration;
     // package level for unit tests
     transient WaterMarkEventGenerator<Tuple> waterMarkEventGenerator;
@@ -80,7 +88,7 @@ public class WindowedBoltExecutor implements IRichBolt {
         timestampExtractor = bolt.getTimestampExtractor();
     }
 
-    private int getTopologyTimeoutMillis(Map<String, Object> topoConf) {
+    protected int getTopologyTimeoutMillis(Map<String, Object> topoConf) {
         if (topoConf.get(Config.TOPOLOGY_ENABLE_MESSAGE_TIMEOUTS) != null) {
             boolean timeOutsEnabled = (boolean) 
topoConf.get(Config.TOPOLOGY_ENABLE_MESSAGE_TIMEOUTS);
             if (!timeOutsEnabled) {
@@ -118,8 +126,8 @@ public class WindowedBoltExecutor implements IRichBolt {
         }
     }
 
-    private void validate(Map<String, Object> topoConf, Count 
windowLengthCount, Duration windowLengthDuration,
-                          Count slidingIntervalCount, Duration 
slidingIntervalDuration) {
+    protected void validate(Map<String, Object> topoConf, Count 
windowLengthCount, Duration windowLengthDuration,
+                            Count slidingIntervalCount, Duration 
slidingIntervalDuration) {
 
         int topologyTimeout = getTopologyTimeoutMillis(topoConf);
         int maxSpoutPending = getMaxSpoutPending(topoConf);
@@ -145,8 +153,12 @@ public class WindowedBoltExecutor implements IRichBolt {
     }
 
     private WindowManager<Tuple> 
initWindowManager(WindowLifecycleListener<Tuple> lifecycleListener, Map<String, 
Object> topoConf,
-                                                   TopologyContext context) {
-        WindowManager<Tuple> manager = new WindowManager<>(lifecycleListener);
+                                                   TopologyContext context, 
Collection<Event<Tuple>> queue, boolean stateful) {
+
+        WindowManager<Tuple> manager = stateful ?
+            new StatefulWindowManager<>(lifecycleListener, queue)
+            : new WindowManager<>(lifecycleListener, queue);
+
         Count windowLengthCount = null;
         Duration slidingIntervalDuration = null;
         Count slidingIntervalCount = null;
@@ -207,6 +219,14 @@ public class WindowedBoltExecutor implements IRichBolt {
         return manager;
     }
 
+    protected void restoreState(Map<String, Optional<?>> state) {
+        windowManager.restoreState(state);
+    }
+
+    protected Map<String, Optional<?>> getState() {
+        return windowManager.getState();
+    }
+
     private Set<GlobalStreamId> getComponentStreams(TopologyContext context) {
         Set<GlobalStreamId> streams = new HashSet<>();
         for (GlobalStreamId streamId : context.getThisSources().keySet()) {
@@ -233,8 +253,8 @@ public class WindowedBoltExecutor implements IRichBolt {
         return timestampExtractor != null;
     }
 
-    private TriggerPolicy<Tuple> getTriggerPolicy(Count slidingIntervalCount, 
Duration slidingIntervalDuration,
-                                                  WindowManager<Tuple> 
manager, EvictionPolicy<Tuple> evictionPolicy) {
+    private TriggerPolicy<Tuple, ?> getTriggerPolicy(Count 
slidingIntervalCount, Duration slidingIntervalDuration,
+                                                  WindowManager<Tuple> 
manager, EvictionPolicy<Tuple, ?> evictionPolicy) {
         if (slidingIntervalCount != null) {
             if (isTupleTs()) {
                 return new 
WatermarkCountTriggerPolicy<>(slidingIntervalCount.value, manager, 
evictionPolicy, manager);
@@ -250,7 +270,7 @@ public class WindowedBoltExecutor implements IRichBolt {
         }
     }
 
-    private EvictionPolicy<Tuple> getEvictionPolicy(Count windowLengthCount, 
Duration windowLengthDuration) {
+    private EvictionPolicy<Tuple, ?> getEvictionPolicy(Count 
windowLengthCount, Duration windowLengthDuration) {
         if (windowLengthCount != null) {
             if (isTupleTs()) {
                 return new 
WatermarkCountEvictionPolicy<>(windowLengthCount.value);
@@ -268,10 +288,20 @@ public class WindowedBoltExecutor implements IRichBolt {
 
     @Override
     public void prepare(Map<String, Object> topoConf, TopologyContext context, 
OutputCollector collector) {
+        doPrepare(topoConf, context, collector, new ConcurrentLinkedQueue<>(), 
false);
+    }
+
+    // NOTE: the queue has to be thread safe.
+    protected void doPrepare(Map<String, Object> topoConf, TopologyContext 
context, OutputCollector collector,
+                             Collection<Event<Tuple>> queue, boolean stateful) 
{
+        Objects.requireNonNull(topoConf);
+        Objects.requireNonNull(context);
+        Objects.requireNonNull(collector);
+        Objects.requireNonNull(queue);
         this.windowedOutputCollector = new WindowedOutputCollector(collector);
         bolt.prepare(topoConf, context, windowedOutputCollector);
         this.listener = newWindowLifecycleListener();
-        this.windowManager = initWindowManager(listener, topoConf, context);
+        this.windowManager = initWindowManager(listener, topoConf, context, 
queue, stateful);
         start();
         LOG.info("Initialized window manager {} ", windowManager);
     }
@@ -301,6 +331,11 @@ public class WindowedBoltExecutor implements IRichBolt {
         bolt.cleanup();
     }
 
+    // for unit tests
+    WindowManager<Tuple> getWindowManager() {
+        return windowManager;
+    }
+
     @Override
     public void declareOutputFields(OutputFieldsDeclarer declarer) {
         String lateTupleStream = (String) 
getComponentConfiguration().get(Config.TOPOLOGY_BOLTS_LATE_TUPLE_STREAM);
@@ -327,19 +362,30 @@ public class WindowedBoltExecutor implements IRichBolt {
             @Override
             public void onActivation(List<Tuple> tuples, List<Tuple> 
newTuples, List<Tuple> expiredTuples, Long timestamp) {
                 windowedOutputCollector.setContext(tuples);
-                bolt.execute(new TupleWindowImpl(tuples, newTuples, 
expiredTuples, getWindowStartTs(timestamp), timestamp));
+                boltExecute(tuples, newTuples, expiredTuples, timestamp);
             }
 
-            private Long getWindowStartTs(Long endTs) {
-                Long res = null;
-                if (endTs != null && windowLengthDuration != null) {
-                    res = endTs - windowLengthDuration.value;
-                }
-                return res;
-            }
         };
     }
 
+    protected void boltExecute(List<Tuple> tuples, List<Tuple> newTuples, 
List<Tuple> expiredTuples, Long timestamp) {
+        bolt.execute(new TupleWindowImpl(tuples, newTuples, expiredTuples, 
getWindowStartTs(timestamp), timestamp));
+    }
+
+    protected void boltExecute(Supplier<Iterator<Tuple>> tuples,
+                               Supplier<Iterator<Tuple>> newTuples,
+                               Supplier<Iterator<Tuple>> expiredTuples, Long 
timestamp) {
+        bolt.execute(new TupleWindowIterImpl(tuples, newTuples, expiredTuples, 
getWindowStartTs(timestamp), timestamp));
+    }
+
+    private Long getWindowStartTs(Long endTs) {
+        Long res = null;
+        if (endTs != null && windowLengthDuration != null) {
+            res = endTs - windowLengthDuration.value;
+        }
+        return res;
+    }
+
     /**
      * Creates an {@link OutputCollector} wrapper that automatically
      * anchors the tuples to inputTuples while emitting.

http://git-wip-us.apache.org/repos/asf/storm/blob/3ba3cabb/storm-client/src/jvm/org/apache/storm/topology/base/BaseStatefulWindowedBolt.java
----------------------------------------------------------------------
diff --git 
a/storm-client/src/jvm/org/apache/storm/topology/base/BaseStatefulWindowedBolt.java
 
b/storm-client/src/jvm/org/apache/storm/topology/base/BaseStatefulWindowedBolt.java
index 7e80830..d98ebf6 100644
--- 
a/storm-client/src/jvm/org/apache/storm/topology/base/BaseStatefulWindowedBolt.java
+++ 
b/storm-client/src/jvm/org/apache/storm/topology/base/BaseStatefulWindowedBolt.java
@@ -15,6 +15,7 @@
  * See the License for the specific language governing permissions and
  * limitations under the License.
  */
+
 package org.apache.storm.topology.base;
 
 import org.apache.storm.Config;
@@ -23,6 +24,12 @@ import org.apache.storm.topology.IStatefulWindowedBolt;
 import org.apache.storm.windowing.TimestampExtractor;
 
 public abstract class BaseStatefulWindowedBolt<T extends State> extends 
BaseWindowedBolt implements IStatefulWindowedBolt<T> {
+    // if the windows should be persisted in state
+    private boolean persistent;
+
+    // max number of window events in memory
+    private long maxEventsInMemory;
+
     /**
      * {@inheritDoc}
      */
@@ -151,6 +158,38 @@ public abstract class BaseStatefulWindowedBolt<T extends 
State> extends BaseWind
         return this;
     }
 
+    /**
+     * If set, the stateful windowed bolt would use the backend state for 
window persistence and
+     * only keep a sub-set of events in memory as specified by {@link 
#withMaxEventsInMemory(long)}.
+     */
+    public BaseStatefulWindowedBolt<T> withPersistence() {
+        persistent = true;
+        return this;
+    }
+
+    /**
+     * The maximum number of window events to keep in memory. This is 
meaningful only if
+     * {@link #withPersistence()} is also set. As the number of events in 
memory grows close
+     * to the maximum, the events that are less likely to be used again are 
evicted and persisted.
+     * The default value for this is {@code 1,000,000}.
+     *
+     * @param maxEventsInMemory the maximum number of window events to keep in 
memory
+     */
+    public BaseStatefulWindowedBolt<T> withMaxEventsInMemory(long 
maxEventsInMemory) {
+        this.maxEventsInMemory = maxEventsInMemory;
+        return this;
+    }
+
+    @Override
+    public boolean isPersistent() {
+        return persistent;
+    }
+
+    @Override
+    public long maxEventsInMemory() {
+        return maxEventsInMemory > 0 ? maxEventsInMemory : 
IStatefulWindowedBolt.super.maxEventsInMemory();
+    }
+
     @Override
     public void preCommit(long txid) {
         // NOOP

http://git-wip-us.apache.org/repos/asf/storm/blob/3ba3cabb/storm-client/src/jvm/org/apache/storm/trident/windowing/AbstractTridentWindowManager.java
----------------------------------------------------------------------
diff --git 
a/storm-client/src/jvm/org/apache/storm/trident/windowing/AbstractTridentWindowManager.java
 
b/storm-client/src/jvm/org/apache/storm/trident/windowing/AbstractTridentWindowManager.java
index a8fbb41..e412eb1 100644
--- 
a/storm-client/src/jvm/org/apache/storm/trident/windowing/AbstractTridentWindowManager.java
+++ 
b/storm-client/src/jvm/org/apache/storm/trident/windowing/AbstractTridentWindowManager.java
@@ -56,7 +56,7 @@ public abstract class AbstractTridentWindowManager<T> 
implements ITridentWindowM
     protected final Queue<TriggerResult> pendingTriggers = new 
ConcurrentLinkedQueue<>();
     protected final AtomicInteger triggerId = new AtomicInteger();
     private final String windowTriggerCountId;
-    private final TriggerPolicy<T> triggerPolicy;
+    private final TriggerPolicy<T, ?> triggerPolicy;
 
     public AbstractTridentWindowManager(WindowConfig windowConfig, String 
windowTaskId, WindowsStore windowStore,
                                         Aggregator aggregator, 
BatchOutputCollector delegateCollector) {
@@ -70,7 +70,7 @@ public abstract class AbstractTridentWindowManager<T> 
implements ITridentWindowM
         windowManager = new WindowManager<>(new 
TridentWindowLifeCycleListener());
 
         WindowStrategy<T> windowStrategy = windowConfig.getWindowStrategy();
-        EvictionPolicy<T> evictionPolicy = windowStrategy.getEvictionPolicy();
+        EvictionPolicy<T, ?> evictionPolicy = 
windowStrategy.getEvictionPolicy();
         windowManager.setEvictionPolicy(evictionPolicy);
         triggerPolicy = windowStrategy.getTriggerPolicy(windowManager, 
evictionPolicy);
         windowManager.setTriggerPolicy(triggerPolicy);

http://git-wip-us.apache.org/repos/asf/storm/blob/3ba3cabb/storm-client/src/jvm/org/apache/storm/trident/windowing/strategy/SlidingCountWindowStrategy.java
----------------------------------------------------------------------
diff --git 
a/storm-client/src/jvm/org/apache/storm/trident/windowing/strategy/SlidingCountWindowStrategy.java
 
b/storm-client/src/jvm/org/apache/storm/trident/windowing/strategy/SlidingCountWindowStrategy.java
index c26b795..bd72a77 100644
--- 
a/storm-client/src/jvm/org/apache/storm/trident/windowing/strategy/SlidingCountWindowStrategy.java
+++ 
b/storm-client/src/jvm/org/apache/storm/trident/windowing/strategy/SlidingCountWindowStrategy.java
@@ -43,7 +43,7 @@ public class SlidingCountWindowStrategy<T> extends 
BaseWindowStrategy<T> {
      * @return
      */
     @Override
-    public TriggerPolicy<T> getTriggerPolicy(TriggerHandler triggerHandler, 
EvictionPolicy<T> evictionPolicy) {
+    public TriggerPolicy<T, ?> getTriggerPolicy(TriggerHandler triggerHandler, 
EvictionPolicy<T, ?> evictionPolicy) {
         return new CountTriggerPolicy<>(windowConfig.getSlidingLength(), 
triggerHandler, evictionPolicy);
     }
 
@@ -53,7 +53,7 @@ public class SlidingCountWindowStrategy<T> extends 
BaseWindowStrategy<T> {
      * @return
      */
     @Override
-    public EvictionPolicy<T> getEvictionPolicy() {
+    public EvictionPolicy<T, ?> getEvictionPolicy() {
         return new CountEvictionPolicy<>(windowConfig.getWindowLength());
     }
 }

http://git-wip-us.apache.org/repos/asf/storm/blob/3ba3cabb/storm-client/src/jvm/org/apache/storm/trident/windowing/strategy/SlidingDurationWindowStrategy.java
----------------------------------------------------------------------
diff --git 
a/storm-client/src/jvm/org/apache/storm/trident/windowing/strategy/SlidingDurationWindowStrategy.java
 
b/storm-client/src/jvm/org/apache/storm/trident/windowing/strategy/SlidingDurationWindowStrategy.java
index 9e71220..301001c 100644
--- 
a/storm-client/src/jvm/org/apache/storm/trident/windowing/strategy/SlidingDurationWindowStrategy.java
+++ 
b/storm-client/src/jvm/org/apache/storm/trident/windowing/strategy/SlidingDurationWindowStrategy.java
@@ -44,7 +44,7 @@ public final class SlidingDurationWindowStrategy<T> extends 
BaseWindowStrategy<T
      * @return
      */
     @Override
-    public TriggerPolicy<T> getTriggerPolicy(TriggerHandler triggerHandler, 
EvictionPolicy<T> evictionPolicy) {
+    public TriggerPolicy<T, ?> getTriggerPolicy(TriggerHandler triggerHandler, 
EvictionPolicy<T, ?> evictionPolicy) {
         return new TimeTriggerPolicy<>(windowConfig.getSlidingLength(), 
triggerHandler, evictionPolicy);
     }
 
@@ -54,7 +54,7 @@ public final class SlidingDurationWindowStrategy<T> extends 
BaseWindowStrategy<T
      * @return
      */
     @Override
-    public EvictionPolicy<T> getEvictionPolicy() {
+    public EvictionPolicy<T, ?> getEvictionPolicy() {
         return new TimeEvictionPolicy<>(windowConfig.getWindowLength());
     }
 }

http://git-wip-us.apache.org/repos/asf/storm/blob/3ba3cabb/storm-client/src/jvm/org/apache/storm/trident/windowing/strategy/TumblingCountWindowStrategy.java
----------------------------------------------------------------------
diff --git 
a/storm-client/src/jvm/org/apache/storm/trident/windowing/strategy/TumblingCountWindowStrategy.java
 
b/storm-client/src/jvm/org/apache/storm/trident/windowing/strategy/TumblingCountWindowStrategy.java
index 5e4d6fe..f7404a7 100644
--- 
a/storm-client/src/jvm/org/apache/storm/trident/windowing/strategy/TumblingCountWindowStrategy.java
+++ 
b/storm-client/src/jvm/org/apache/storm/trident/windowing/strategy/TumblingCountWindowStrategy.java
@@ -44,7 +44,7 @@ public final class TumblingCountWindowStrategy<T> extends 
BaseWindowStrategy<T>
      * @return
      */
     @Override
-    public TriggerPolicy<T> getTriggerPolicy(TriggerHandler triggerHandler, 
EvictionPolicy<T> evictionPolicy) {
+    public TriggerPolicy<T, ?> getTriggerPolicy(TriggerHandler triggerHandler, 
EvictionPolicy<T, ?> evictionPolicy) {
         return new CountTriggerPolicy<>(windowConfig.getSlidingLength(), 
triggerHandler, evictionPolicy);
     }
 
@@ -54,7 +54,7 @@ public final class TumblingCountWindowStrategy<T> extends 
BaseWindowStrategy<T>
      * @return
      */
     @Override
-    public EvictionPolicy<T> getEvictionPolicy() {
+    public EvictionPolicy<T, ?> getEvictionPolicy() {
         return new CountEvictionPolicy<>(windowConfig.getWindowLength());
     }
 }

http://git-wip-us.apache.org/repos/asf/storm/blob/3ba3cabb/storm-client/src/jvm/org/apache/storm/trident/windowing/strategy/TumblingDurationWindowStrategy.java
----------------------------------------------------------------------
diff --git 
a/storm-client/src/jvm/org/apache/storm/trident/windowing/strategy/TumblingDurationWindowStrategy.java
 
b/storm-client/src/jvm/org/apache/storm/trident/windowing/strategy/TumblingDurationWindowStrategy.java
index 4478667..8b0a65a 100644
--- 
a/storm-client/src/jvm/org/apache/storm/trident/windowing/strategy/TumblingDurationWindowStrategy.java
+++ 
b/storm-client/src/jvm/org/apache/storm/trident/windowing/strategy/TumblingDurationWindowStrategy.java
@@ -44,7 +44,7 @@ public final class TumblingDurationWindowStrategy<T> extends 
BaseWindowStrategy<
      * @return
      */
     @Override
-    public TriggerPolicy<T> getTriggerPolicy(TriggerHandler triggerHandler, 
EvictionPolicy<T> evictionPolicy) {
+    public TriggerPolicy<T, ?> getTriggerPolicy(TriggerHandler triggerHandler, 
EvictionPolicy<T, ?> evictionPolicy) {
         return new TimeTriggerPolicy<>(windowConfig.getSlidingLength(), 
triggerHandler, evictionPolicy);
     }
 
@@ -54,7 +54,7 @@ public final class TumblingDurationWindowStrategy<T> extends 
BaseWindowStrategy<
      * @return
      */
     @Override
-    public EvictionPolicy<T> getEvictionPolicy() {
+    public EvictionPolicy<T, ?> getEvictionPolicy() {
         return new TimeEvictionPolicy<>(windowConfig.getWindowLength());
     }
 }

http://git-wip-us.apache.org/repos/asf/storm/blob/3ba3cabb/storm-client/src/jvm/org/apache/storm/trident/windowing/strategy/WindowStrategy.java
----------------------------------------------------------------------
diff --git 
a/storm-client/src/jvm/org/apache/storm/trident/windowing/strategy/WindowStrategy.java
 
b/storm-client/src/jvm/org/apache/storm/trident/windowing/strategy/WindowStrategy.java
index 1dfb264..71f90d4 100644
--- 
a/storm-client/src/jvm/org/apache/storm/trident/windowing/strategy/WindowStrategy.java
+++ 
b/storm-client/src/jvm/org/apache/storm/trident/windowing/strategy/WindowStrategy.java
@@ -34,12 +34,12 @@ public interface WindowStrategy<T> {
      * @param evictionPolicy
      * @return
      */
-    public TriggerPolicy<T> getTriggerPolicy(TriggerHandler triggerHandler, 
EvictionPolicy<T> evictionPolicy);
+    public TriggerPolicy<T, ?> getTriggerPolicy(TriggerHandler triggerHandler, 
EvictionPolicy<T, ?> evictionPolicy);
 
     /**
      * Returns an {@code EvictionPolicy} instance for this strategy with the 
given configuration.
      *
      * @return
      */
-    public EvictionPolicy<T> getEvictionPolicy();
+    public EvictionPolicy<T, ?> getEvictionPolicy();
 }

http://git-wip-us.apache.org/repos/asf/storm/blob/3ba3cabb/storm-client/src/jvm/org/apache/storm/windowing/CountEvictionPolicy.java
----------------------------------------------------------------------
diff --git 
a/storm-client/src/jvm/org/apache/storm/windowing/CountEvictionPolicy.java 
b/storm-client/src/jvm/org/apache/storm/windowing/CountEvictionPolicy.java
index 6a9a4f8..9b471c6 100644
--- a/storm-client/src/jvm/org/apache/storm/windowing/CountEvictionPolicy.java
+++ b/storm-client/src/jvm/org/apache/storm/windowing/CountEvictionPolicy.java
@@ -25,7 +25,7 @@ import java.util.concurrent.atomic.AtomicLong;
  *
  * @param <T> the type of event tracked by this policy.
  */
-public class CountEvictionPolicy<T> implements EvictionPolicy<T> {
+public class CountEvictionPolicy<T> implements EvictionPolicy<T, Long> {
     protected final int threshold;
     protected final AtomicLong currentCount;
     private EvictionContext context;
@@ -78,4 +78,19 @@ public class CountEvictionPolicy<T> implements 
EvictionPolicy<T> {
                 ", currentCount=" + currentCount +
                 '}';
     }
+
+    @Override
+    public void reset() {
+        // NOOP
+    }
+
+    @Override
+    public Long getState() {
+        return currentCount.get();
+    }
+
+    @Override
+    public void restoreState(Long state) {
+        currentCount.set(state);
+    }
 }

Reply via email to