Repository: storm Updated Branches: refs/heads/master 1c2ac2eb4 -> aaebc3b23
http://git-wip-us.apache.org/repos/asf/storm/blob/3ba3cabb/storm-client/test/jvm/org/apache/storm/topology/SimpleWindowPartitionCacheTest.java ---------------------------------------------------------------------- diff --git a/storm-client/test/jvm/org/apache/storm/topology/SimpleWindowPartitionCacheTest.java b/storm-client/test/jvm/org/apache/storm/topology/SimpleWindowPartitionCacheTest.java new file mode 100644 index 0000000..8c19ee8 --- /dev/null +++ b/storm-client/test/jvm/org/apache/storm/topology/SimpleWindowPartitionCacheTest.java @@ -0,0 +1,233 @@ +/** + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * <p> + * http://www.apache.org/licenses/LICENSE-2.0 + * <p> + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.storm.topology; + +import org.apache.storm.utils.Utils; +import org.apache.storm.windowing.persistence.SimpleWindowPartitionCache; +import org.apache.storm.windowing.persistence.WindowPartitionCache; +import org.junit.Assert; +import org.junit.Before; +import org.junit.Test; + +import java.util.ArrayList; +import java.util.Arrays; +import java.util.Collections; +import java.util.List; +import java.util.concurrent.FutureTask; + +/** + * Unit tests for {@link SimpleWindowPartitionCache} + */ +public class SimpleWindowPartitionCacheTest { + + @Before + public void setUp() throws Exception { + } + + @Test(expected = IllegalArgumentException.class) + public void testBuildInvalid1() throws Exception { + SimpleWindowPartitionCache.<Integer, Integer>newBuilder() + .maximumSize(0) + .build(null); + } + + @Test(expected = IllegalArgumentException.class) + public void testBuildInvalid2() throws Exception { + SimpleWindowPartitionCache.<Integer, Integer>newBuilder() + .maximumSize(-1) + .build(null); + } + + @Test(expected = NullPointerException.class) + public void testBuildInvalid3() throws Exception { + SimpleWindowPartitionCache.<Integer, Integer>newBuilder() + .maximumSize(1) + .build(null); + } + + @Test + public void testBuildOk() throws Exception { + SimpleWindowPartitionCache.<Integer, Integer>newBuilder() + .maximumSize(1) + .removalListener((key, val, removalCause) -> { + }) + .build(key -> key); + } + + @Test + public void testGet() throws Exception { + List<Integer> removed = new ArrayList<>(); + List<Integer> loaded = new ArrayList<>(); + SimpleWindowPartitionCache<Integer, Integer> cache = + SimpleWindowPartitionCache.<Integer, Integer>newBuilder() + .maximumSize(2) + .removalListener((key, val, removalCause) -> removed.add(key)) + .build(key -> { + loaded.add(key); + return key; + }); + + cache.get(1); + cache.get(2); + cache.get(3); + Assert.assertEquals(Arrays.asList(1, 2, 3), loaded); + // since 2 is the largest un-pinned entry before 3 is loaded + Assert.assertEquals(Collections.singletonList(2), removed); + } + + @Test(expected = NullPointerException.class) + public void testGetNull() throws Exception { + SimpleWindowPartitionCache<Integer, Integer> cache = + SimpleWindowPartitionCache.<Integer, Integer>newBuilder() + .maximumSize(2) + .build(key -> null); + + cache.get(1); + } + + @Test + public void testEvictNoRemovalListener() throws Exception { + SimpleWindowPartitionCache<Integer, Integer> cache = + SimpleWindowPartitionCache.<Integer, Integer>newBuilder() + .maximumSize(1) + .build(key -> { + return key; + }); + cache.get(1); + cache.get(2); + Assert.assertEquals(Collections.singletonMap(2, 2), cache.asMap()); + cache.invalidate(2); + Assert.assertEquals(Collections.emptyMap(), cache.asMap()); + } + + @Test + public void testPinAndGet() throws Exception { + List<Integer> removed = new ArrayList<>(); + List<Integer> loaded = new ArrayList<>(); + SimpleWindowPartitionCache<Integer, Integer> cache = + SimpleWindowPartitionCache.<Integer, Integer>newBuilder() + .maximumSize(1) + .removalListener(new WindowPartitionCache.RemovalListener<Integer, Integer>() { + @Override + public void onRemoval(Integer key, Integer val, WindowPartitionCache.RemovalCause removalCause) { + removed.add(key); + } + }) + .build(new WindowPartitionCache.CacheLoader<Integer, Integer>() { + @Override + public Integer load(Integer key) { + loaded.add(key); + return key; + } + }); + + cache.get(1); + cache.pinAndGet(2); + cache.get(3); + Assert.assertEquals(Arrays.asList(1, 2, 3), loaded); + Assert.assertEquals(Collections.singletonList(1), removed); + } + + @Test + public void testInvalidate() throws Exception { + List<Integer> removed = new ArrayList<>(); + List<Integer> loaded = new ArrayList<>(); + SimpleWindowPartitionCache<Integer, Integer> cache = + SimpleWindowPartitionCache.<Integer, Integer>newBuilder() + .maximumSize(1) + .removalListener((key, val, removalCause) -> removed.add(key)) + .build(key -> { + loaded.add(key); + return key; + }); + + cache.pinAndGet(1); + cache.invalidate(1); + Assert.assertEquals(Collections.singletonList(1), loaded); + Assert.assertEquals(Collections.emptyList(), removed); + Assert.assertEquals(cache.asMap(), Collections.singletonMap(1, 1)); + + cache.unpin(1); + cache.invalidate(1); + Assert.assertTrue(cache.asMap().isEmpty()); + } + + + @Test(timeout = 10000) + public void testConcurrentGet() throws Exception { + List<Integer> loaded = new ArrayList<>(); + SimpleWindowPartitionCache<Integer, Object> cache = + SimpleWindowPartitionCache.<Integer, Object>newBuilder() + .maximumSize(1) + .build(key -> { + Utils.sleep(1000); + loaded.add(key); + return new Object(); + }); + + FutureTask<Object> ft1 = new FutureTask<>(() -> cache.pinAndGet(1)); + FutureTask<Object> ft2 = new FutureTask<>(() -> cache.pinAndGet(1)); + Thread t1 = new Thread(ft1); + Thread t2 = new Thread(ft2); + t1.start(); + t2.start(); + t1.join(); + t2.join(); + + Assert.assertEquals(Collections.singletonList(1), loaded); + Assert.assertEquals(ft1.get(), ft2.get()); + } + + @Test + public void testConcurrentUnpin() throws Exception { + SimpleWindowPartitionCache<Integer, Object> cache = + SimpleWindowPartitionCache.<Integer, Object>newBuilder() + .maximumSize(1) + .build(key -> new Object()); + + cache.pinAndGet(1); + FutureTask<Boolean> ft1 = new FutureTask<>(() -> cache.unpin(1)); + FutureTask<Boolean> ft2 = new FutureTask<>(() -> cache.unpin(1)); + Thread t1 = new Thread(ft1); + Thread t2 = new Thread(ft2); + t1.start(); + t2.start(); + t1.join(); + t2.join(); + + Assert.assertTrue(ft1.get() || ft2.get()); + Assert.assertFalse(ft1.get() && ft2.get()); + } + + @Test + public void testEviction() throws Exception { + List<Integer> removed = new ArrayList<>(); + SimpleWindowPartitionCache<Integer, Object> cache = + SimpleWindowPartitionCache.<Integer, Object>newBuilder() + .maximumSize(1) + .removalListener((key, val, removalCause) -> removed.add(key)) + .build(key -> new Object()); + + cache.get(0); + cache.pinAndGet(1); + Assert.assertEquals(Collections.singletonList(0), removed); + cache.get(2); + Assert.assertEquals(Collections.singletonList(0), removed); + } +} \ No newline at end of file http://git-wip-us.apache.org/repos/asf/storm/blob/3ba3cabb/storm-client/test/jvm/org/apache/storm/windowing/WindowManagerTest.java ---------------------------------------------------------------------- diff --git a/storm-client/test/jvm/org/apache/storm/windowing/WindowManagerTest.java b/storm-client/test/jvm/org/apache/storm/windowing/WindowManagerTest.java index 178c1bb..d99ecb3 100644 --- a/storm-client/test/jvm/org/apache/storm/windowing/WindowManagerTest.java +++ b/storm-client/test/jvm/org/apache/storm/windowing/WindowManagerTest.java @@ -98,8 +98,8 @@ public class WindowManagerTest { @Test public void testCountBasedWindow() throws Exception { - EvictionPolicy<Integer> evictionPolicy = new CountEvictionPolicy<Integer>(5); - TriggerPolicy<Integer> triggerPolicy = new CountTriggerPolicy<Integer>(2, windowManager, evictionPolicy); + EvictionPolicy<Integer, ?> evictionPolicy = new CountEvictionPolicy<Integer>(5); + TriggerPolicy<Integer, ?> triggerPolicy = new CountTriggerPolicy<Integer>(2, windowManager, evictionPolicy); triggerPolicy.start(); windowManager.setEvictionPolicy(evictionPolicy); windowManager.setTriggerPolicy(triggerPolicy); @@ -141,7 +141,7 @@ public class WindowManagerTest { int threshold = WindowManager.EXPIRE_EVENTS_THRESHOLD; int windowLength = 5; windowManager.setEvictionPolicy(new CountEvictionPolicy<Integer>(5)); - TriggerPolicy<Integer> triggerPolicy = new TimeTriggerPolicy<Integer>(new Duration(1, TimeUnit.HOURS).value, windowManager); + TriggerPolicy<Integer, ?> triggerPolicy = new TimeTriggerPolicy<Integer>(new Duration(1, TimeUnit.HOURS).value, windowManager); triggerPolicy.start(); windowManager.setTriggerPolicy(triggerPolicy); for (int i : seq(1, 5)) { @@ -203,13 +203,13 @@ public class WindowManagerTest { @Test public void testTimeBasedWindow() throws Exception { - EvictionPolicy<Integer> evictionPolicy = new TimeEvictionPolicy<Integer>(new Duration(1, TimeUnit.SECONDS).value); + EvictionPolicy<Integer, ?> evictionPolicy = new TimeEvictionPolicy<Integer>(new Duration(1, TimeUnit.SECONDS).value); windowManager.setEvictionPolicy(evictionPolicy); /* * Don't wait for Timetrigger to fire since this could lead to timing issues in unit tests. * Set it to a large value and trigger manually. */ - TriggerPolicy<Integer> triggerPolicy = new TimeTriggerPolicy<Integer>(new Duration(1, TimeUnit.DAYS).value, windowManager, evictionPolicy); + TriggerPolicy<Integer, ?> triggerPolicy = new TimeTriggerPolicy<Integer>(new Duration(1, TimeUnit.DAYS).value, windowManager, evictionPolicy); triggerPolicy.start(); windowManager.setTriggerPolicy(triggerPolicy); long now = System.currentTimeMillis(); @@ -265,13 +265,13 @@ public class WindowManagerTest { @Test public void testTimeBasedWindowExpiry() throws Exception { - EvictionPolicy<Integer> evictionPolicy = new TimeEvictionPolicy<Integer>(new Duration(100, TimeUnit.MILLISECONDS).value); + EvictionPolicy<Integer, ?> evictionPolicy = new TimeEvictionPolicy<Integer>(new Duration(100, TimeUnit.MILLISECONDS).value); windowManager.setEvictionPolicy(evictionPolicy); /* * Don't wait for Timetrigger to fire since this could lead to timing issues in unit tests. * Set it to a large value and trigger manually. */ - TriggerPolicy<Integer> triggerPolicy = new TimeTriggerPolicy<Integer>(new Duration(1, TimeUnit.DAYS).value, windowManager); + TriggerPolicy<Integer, ?> triggerPolicy = new TimeTriggerPolicy<Integer>(new Duration(1, TimeUnit.DAYS).value, windowManager); triggerPolicy.start(); windowManager.setTriggerPolicy(triggerPolicy); long now = System.currentTimeMillis(); @@ -302,9 +302,9 @@ public class WindowManagerTest { @Test public void testTumblingWindow() throws Exception { - EvictionPolicy<Integer> evictionPolicy = new CountEvictionPolicy<Integer>(3); + EvictionPolicy<Integer, ?> evictionPolicy = new CountEvictionPolicy<Integer>(3); windowManager.setEvictionPolicy(evictionPolicy); - TriggerPolicy<Integer> triggerPolicy = new CountTriggerPolicy<Integer>(3, windowManager, evictionPolicy); + TriggerPolicy<Integer, ?> triggerPolicy = new CountTriggerPolicy<Integer>(3, windowManager, evictionPolicy); triggerPolicy.start(); windowManager.setTriggerPolicy(triggerPolicy); windowManager.add(1); @@ -332,9 +332,9 @@ public class WindowManagerTest { @Test public void testEventTimeBasedWindow() throws Exception { - EvictionPolicy<Integer> evictionPolicy = new WatermarkTimeEvictionPolicy<>(20); + EvictionPolicy<Integer, ?> evictionPolicy = new WatermarkTimeEvictionPolicy<>(20); windowManager.setEvictionPolicy(evictionPolicy); - TriggerPolicy<Integer> triggerPolicy = new WatermarkTimeTriggerPolicy<Integer>(10, windowManager, evictionPolicy, windowManager); + TriggerPolicy<Integer, ?> triggerPolicy = new WatermarkTimeTriggerPolicy<Integer>(10, windowManager, evictionPolicy, windowManager); triggerPolicy.start(); windowManager.setTriggerPolicy(triggerPolicy); @@ -398,9 +398,9 @@ public class WindowManagerTest { @Test public void testCountBasedWindowWithEventTs() throws Exception { - EvictionPolicy<Integer> evictionPolicy = new WatermarkCountEvictionPolicy<>(3); + EvictionPolicy<Integer, ?> evictionPolicy = new WatermarkCountEvictionPolicy<>(3); windowManager.setEvictionPolicy(evictionPolicy); - TriggerPolicy<Integer> triggerPolicy = new WatermarkTimeTriggerPolicy<Integer>(10, windowManager, evictionPolicy, windowManager); + TriggerPolicy<Integer, ?> triggerPolicy = new WatermarkTimeTriggerPolicy<Integer>(10, windowManager, evictionPolicy, windowManager); triggerPolicy.start(); windowManager.setTriggerPolicy(triggerPolicy); @@ -437,9 +437,9 @@ public class WindowManagerTest { @Test public void testCountBasedTriggerWithEventTs() throws Exception { - EvictionPolicy<Integer> evictionPolicy = new WatermarkTimeEvictionPolicy<Integer>(20); + EvictionPolicy<Integer, ?> evictionPolicy = new WatermarkTimeEvictionPolicy<Integer>(20); windowManager.setEvictionPolicy(evictionPolicy); - TriggerPolicy<Integer> triggerPolicy = new WatermarkCountTriggerPolicy<Integer>(3, windowManager, evictionPolicy, windowManager); + TriggerPolicy<Integer, ?> triggerPolicy = new WatermarkCountTriggerPolicy<Integer>(3, windowManager, evictionPolicy, windowManager); triggerPolicy.start(); windowManager.setTriggerPolicy(triggerPolicy); @@ -477,9 +477,9 @@ public class WindowManagerTest { @Test public void testCountBasedTumblingWithSameEventTs() throws Exception { - EvictionPolicy<Integer> evictionPolicy = new WatermarkCountEvictionPolicy<>(2); + EvictionPolicy<Integer, ?> evictionPolicy = new WatermarkCountEvictionPolicy<>(2); windowManager.setEvictionPolicy(evictionPolicy); - TriggerPolicy<Integer> triggerPolicy = new WatermarkCountTriggerPolicy<Integer>(2, windowManager, evictionPolicy, windowManager); + TriggerPolicy<Integer, ?> triggerPolicy = new WatermarkCountTriggerPolicy<Integer>(2, windowManager, evictionPolicy, windowManager); triggerPolicy.start(); windowManager.setTriggerPolicy(triggerPolicy); @@ -505,9 +505,9 @@ public class WindowManagerTest { @Test public void testCountBasedSlidingWithSameEventTs() throws Exception { - EvictionPolicy<Integer> evictionPolicy = new WatermarkCountEvictionPolicy<>(5); + EvictionPolicy<Integer, ?> evictionPolicy = new WatermarkCountEvictionPolicy<>(5); windowManager.setEvictionPolicy(evictionPolicy); - TriggerPolicy<Integer> triggerPolicy = new WatermarkCountTriggerPolicy<Integer>(2, windowManager, evictionPolicy, windowManager); + TriggerPolicy<Integer, ?> triggerPolicy = new WatermarkCountTriggerPolicy<Integer>(2, windowManager, evictionPolicy, windowManager); triggerPolicy.start(); windowManager.setTriggerPolicy(triggerPolicy); @@ -533,9 +533,9 @@ public class WindowManagerTest { } @Test public void testEventTimeLag() throws Exception { - EvictionPolicy<Integer> evictionPolicy = new WatermarkTimeEvictionPolicy<>(20, 5); + EvictionPolicy<Integer, ?> evictionPolicy = new WatermarkTimeEvictionPolicy<>(20, 5); windowManager.setEvictionPolicy(evictionPolicy); - TriggerPolicy<Integer> triggerPolicy = new WatermarkTimeTriggerPolicy<Integer>(10, windowManager, evictionPolicy, windowManager); + TriggerPolicy<Integer, ?> triggerPolicy = new WatermarkTimeTriggerPolicy<Integer>(10, windowManager, evictionPolicy, windowManager); triggerPolicy.start(); windowManager.setTriggerPolicy(triggerPolicy); @@ -560,7 +560,7 @@ public class WindowManagerTest { @Test public void testScanStop() throws Exception { final Set<Integer> eventsScanned = new HashSet<>(); - EvictionPolicy<Integer> evictionPolicy = new WatermarkTimeEvictionPolicy<Integer>(20, 5) { + EvictionPolicy<Integer, ?> evictionPolicy = new WatermarkTimeEvictionPolicy<Integer>(20, 5) { @Override public Action evict(Event<Integer> event) { @@ -570,7 +570,7 @@ public class WindowManagerTest { }; windowManager.setEvictionPolicy(evictionPolicy); - TriggerPolicy<Integer> triggerPolicy = new WatermarkTimeTriggerPolicy<Integer>(10, windowManager, evictionPolicy, windowManager); + TriggerPolicy<Integer, ?> triggerPolicy = new WatermarkTimeTriggerPolicy<Integer>(10, windowManager, evictionPolicy, windowManager); triggerPolicy.start(); windowManager.setTriggerPolicy(triggerPolicy); http://git-wip-us.apache.org/repos/asf/storm/blob/3ba3cabb/storm-client/test/jvm/org/apache/storm/windowing/persistence/WindowStateTest.java ---------------------------------------------------------------------- diff --git a/storm-client/test/jvm/org/apache/storm/windowing/persistence/WindowStateTest.java b/storm-client/test/jvm/org/apache/storm/windowing/persistence/WindowStateTest.java new file mode 100644 index 0000000..587d874 --- /dev/null +++ b/storm-client/test/jvm/org/apache/storm/windowing/persistence/WindowStateTest.java @@ -0,0 +1,246 @@ +/** + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * <p> + * http://www.apache.org/licenses/LICENSE-2.0 + * <p> + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.storm.windowing.persistence; + +import org.apache.storm.state.KeyValueState; +import org.apache.storm.tuple.Tuple; +import org.apache.storm.windowing.Event; +import org.junit.Assert; +import org.junit.Before; +import org.junit.Test; +import org.junit.runner.RunWith; +import org.mockito.ArgumentCaptor; +import org.mockito.Captor; +import org.mockito.Mock; +import org.mockito.Mockito; +import org.mockito.MockitoAnnotations; +import org.mockito.invocation.InvocationOnMock; +import org.mockito.runners.MockitoJUnitRunner; +import org.mockito.stubbing.Answer; + +import java.util.ArrayList; +import java.util.Collections; +import java.util.Deque; +import java.util.HashMap; +import java.util.Iterator; +import java.util.List; +import java.util.Map; +import java.util.Optional; +import java.util.function.Consumer; +import java.util.function.Supplier; + +import static org.mockito.AdditionalAnswers.returnsArgAt; + +/** + * Unit tests for {@link WindowState} + */ +@RunWith(MockitoJUnitRunner.class) +public class WindowStateTest { + + @Mock + private KeyValueState<Long, WindowState.WindowPartition<Integer>> windowState; + @Mock + private KeyValueState<String, Deque<Long>> partitionIdsState; + @Mock + private KeyValueState<String, Optional<?>> systemState; + @Mock + private Supplier<Map<String, Optional<?>>> supplier; + @Captor + private ArgumentCaptor<Long> longCaptor; + @Captor + private ArgumentCaptor<WindowState.WindowPartition<Integer>> windowValuesCaptor; + + @Before + public void setUp() throws Exception { + MockitoAnnotations.initMocks(this); + } + + @Test + public void testAdd() throws Exception { + Mockito.when(partitionIdsState.get(Mockito.any(), Mockito.any())).then(returnsArgAt(1)); + Mockito.when(windowState.get(Mockito.any(), Mockito.any())).then(returnsArgAt(1)); + + WindowState<Integer> ws = getWindowState(10 * WindowState.MAX_PARTITION_EVENTS); + + long partitions = 15; + long numEvents = partitions * WindowState.MAX_PARTITION_EVENTS; + for (int i = 0; i < numEvents; i++) { + ws.add(getEvent(i)); + } + // 5 partitions evicted to window state + Mockito.verify(windowState, Mockito.times(5)).put(longCaptor.capture(), windowValuesCaptor.capture()); + Assert.assertEquals(5, longCaptor.getAllValues().size()); + // each evicted partition has MAX_EVENTS_PER_PARTITION + windowValuesCaptor.getAllValues().forEach(wp -> { + Assert.assertEquals(WindowState.MAX_PARTITION_EVENTS, wp.size()); + }); + // last partition is not evicted + Assert.assertFalse(longCaptor.getAllValues().contains(partitions - 1)); + } + + @Test + public void testIterator() throws Exception { + Map<Long, WindowState.WindowPartition<Event<Tuple>>> partitionMap = new HashMap<>(); + Mockito.when(partitionIdsState.get(Mockito.any(), Mockito.any())).then(returnsArgAt(1)); + Mockito.when(windowState.get(Mockito.any(), Mockito.any())).then(new Answer<Object>() { + @Override + public Object answer(InvocationOnMock invocation) throws Throwable { + Object[] args = invocation.getArguments(); + WindowState.WindowPartition<Event<Tuple>> evicted = partitionMap.get(args[0]); + return evicted != null ? evicted : args[1]; + } + }); + + Mockito.doAnswer(new Answer<Void>() { + @Override + public Void answer(InvocationOnMock invocation) throws Throwable { + Object[] args = invocation.getArguments(); + partitionMap.put((long)args[0], (WindowState.WindowPartition<Event<Tuple>>)args[1]); + return null; + } + }).when(windowState).put(Mockito.any(), Mockito.any()); + + Mockito.doAnswer(new Answer<Void>() { + @Override + public Void answer(InvocationOnMock invocation) throws Throwable { + Object[] args = invocation.getArguments(); + partitionMap.remove(args[0]); + return null; + } + }).when(windowState).delete(Mockito.anyLong()); + + Mockito.when(supplier.get()).thenReturn(Collections.emptyMap()); + + WindowState<Integer> ws = getWindowState(10 * WindowState.MAX_PARTITION_EVENTS); + + long partitions = 15; + + long numEvents = partitions * WindowState.MAX_PARTITION_EVENTS; + List<Event<Integer>> expected = new ArrayList<>(); + for (int i = 0; i < numEvents; i++) { + Event<Integer> event = getEvent(i); + expected.add(event); + ws.add(event); + } + + Assert.assertEquals(5, partitionMap.size()); + Iterator<Event<Integer>> it = ws.iterator(); + List<Event<Integer>> actual = new ArrayList<>(); + it.forEachRemaining(actual::add); + Assert.assertEquals(expected, actual); + + // iterate again + it = ws.iterator(); + actual.clear(); + it.forEachRemaining(actual::add); + Assert.assertEquals(expected, actual); + + // remove + it = ws.iterator(); + while (it.hasNext()) { + it.next(); + it.remove(); + } + + it = ws.iterator(); + actual.clear(); + it.forEachRemaining(actual::add); + Assert.assertEquals(Collections.emptyList(), actual); + } + + @Test + public void testIteratorPartitionNotEvicted() throws Exception { + Map<Long, WindowState.WindowPartition<Event<Tuple>>> partitionMap = new HashMap<>(); + Mockito.when(partitionIdsState.get(Mockito.any(), Mockito.any())).then(returnsArgAt(1)); + Mockito.when(windowState.get(Mockito.any(), Mockito.any())).then(new Answer<Object>() { + @Override + public Object answer(InvocationOnMock invocation) throws Throwable { + Object[] args = invocation.getArguments(); + WindowState.WindowPartition<Event<Tuple>> evicted = partitionMap.get(args[0]); + return evicted != null ? evicted : args[1]; + } + }); + + Mockito.doAnswer(new Answer<Void>() { + @Override + public Void answer(InvocationOnMock invocation) throws Throwable { + Object[] args = invocation.getArguments(); + partitionMap.put((long)args[0], (WindowState.WindowPartition<Event<Tuple>>)args[1]); + return null; + } + }).when(windowState).put(Mockito.any(), Mockito.any()); + + Mockito.when(supplier.get()).thenReturn(Collections.emptyMap()); + + WindowState<Integer> ws = getWindowState(10 * WindowState.MAX_PARTITION_EVENTS); + + long partitions = 10; + + long numEvents = partitions * WindowState.MAX_PARTITION_EVENTS; + List<Event<Integer>> expected = new ArrayList<>(); + for (int i = 0; i < numEvents; i++) { + Event<Integer> event = getEvent(i); + expected.add(event); + ws.add(event); + } + + // Stop iterating in the middle of the 10th partition + Iterator<Event<Integer>> it = ws.iterator(); + for(int i=0; i<9500; i++) { + it.next(); + } + + for (int i = 0; i < numEvents; i++) { + Event<Integer> event = getEvent(i); + expected.add(event); + ws.add(event); + } + + // 10th partition should not have been evicted + Assert.assertFalse(partitionMap.containsKey(9L)); + } + + private Event<Integer> getEvent(int i) { + return getEvent(i, 0); + } + + private Event<Integer> getEvent(int i, long ts) { + return new Event<Integer>() { + @Override + public long getTimestamp() { + return ts; + } + + @Override + public Integer get() { + return i; + } + + @Override + public boolean isWatermark() { + return false; + } + }; + } + + private WindowState<Integer> getWindowState(int maxEvents) { + return new WindowState<>(windowState, partitionIdsState, systemState, + supplier, maxEvents); + } +} \ No newline at end of file
