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); + } }
