diff --git a/spring-geode/src/main/java/org/springframework/geode/cache/AsyncInlineCachingRegionConfigurer.java b/spring-geode/src/main/java/org/springframework/geode/cache/AsyncInlineCachingRegionConfigurer.java index 6b5dfb84..d2002464 100644 --- a/spring-geode/src/main/java/org/springframework/geode/cache/AsyncInlineCachingRegionConfigurer.java +++ b/spring-geode/src/main/java/org/springframework/geode/cache/AsyncInlineCachingRegionConfigurer.java @@ -514,9 +514,10 @@ public class AsyncInlineCachingRegionConfigurer implements RegionConfigur * Builder method used to enable all {@link AsyncEventQueue AEQs} attached to {@link Region Regions} hosted * and distributed across the cache cluster to process cache events. * - * Default is {@literal false}. + * Default is {@literal false}, or {@literal serial}. * * @return this {@link AsyncInlineCachingRegionConfigurer}. + * @see #withSerialQueue() */ public AsyncInlineCachingRegionConfigurer withParallelQueue() { this.parallel = true; @@ -738,4 +739,18 @@ public class AsyncInlineCachingRegionConfigurer implements RegionConfigur return this; } + + /** + * Builder method used to enable a single {@link AsyncEventQueue AEQ} attached to a {@link Region Region} + * (possibly) hosted and distributed across the cache cluster to process cache events. + * + * Default is {@literal false}, or {@literal serial}. + * + * @return this {@link AsyncInlineCachingRegionConfigurer}. + * @see #withParallelQueue() + */ + public AsyncInlineCachingRegionConfigurer withSerialQueue() { + this.parallel = false; + return this; + } } diff --git a/spring-geode/src/test/java/org/springframework/geode/cache/AsyncInlineCachingRegionConfigurerUnitTests.java b/spring-geode/src/test/java/org/springframework/geode/cache/AsyncInlineCachingRegionConfigurerUnitTests.java index 2e6c47d6..30a74219 100644 --- a/spring-geode/src/test/java/org/springframework/geode/cache/AsyncInlineCachingRegionConfigurerUnitTests.java +++ b/spring-geode/src/test/java/org/springframework/geode/cache/AsyncInlineCachingRegionConfigurerUnitTests.java @@ -372,10 +372,8 @@ public class AsyncInlineCachingRegionConfigurerUnitTests { assertThat(regionConfigurer.withQueueDiskSynchronizationEnabled()).isSameAs(regionConfigurer); assertThat(regionConfigurer.withQueueDispatcherThreadCount(8)).isSameAs(regionConfigurer); assertThat(regionConfigurer.withQueueEventDispatchingPaused()).isSameAs(regionConfigurer); - assertThat(regionConfigurer.withQueueEventFilters(Arrays.asList(mockEventFilterOne, mockEventFilterTwo))) - .isSameAs(regionConfigurer); - assertThat(regionConfigurer.withQueueEventSubstitutionFilter(mockEventSubstitutionFilter)) - .isSameAs(regionConfigurer); + assertThat(regionConfigurer.withQueueEventFilters(Arrays.asList(mockEventFilterOne, mockEventFilterTwo))).isSameAs(regionConfigurer); + assertThat(regionConfigurer.withQueueEventSubstitutionFilter(mockEventSubstitutionFilter)).isSameAs(regionConfigurer); assertThat(regionConfigurer.withQueueForwardedExpirationDestroyEvents()).isSameAs(regionConfigurer); assertThat(regionConfigurer.withQueueMaxMemory(51)).isSameAs(regionConfigurer); assertThat(regionConfigurer.withQueueOrderPolicy(GatewaySender.OrderPolicy.THREAD)).isSameAs(regionConfigurer); @@ -422,6 +420,50 @@ public class AsyncInlineCachingRegionConfigurerUnitTests { mockEventFilterOne, mockEventFilterTwo, mockEventSubstitutionFilter); } + @Test + public void newAsyncEventQueueIsParallel() { + + AsyncEventQueueFactory mockAsyncEventQueueFactory = mock(AsyncEventQueueFactory.class); + + Cache mockCache = mock(Cache.class); + + doReturn(mockAsyncEventQueueFactory).when(mockCache).createAsyncEventQueueFactory(); + + CrudRepository mockRepository = mock(CrudRepository.class); + + AsyncInlineCachingRegionConfigurer regionConfigurer = + spy(AsyncInlineCachingRegionConfigurer.create(mockRepository, Predicate.isEqual("TestRegion"))); + + assertThat(regionConfigurer).isNotNull(); + assertThat(regionConfigurer.withParallelQueue()).isSameAs(regionConfigurer); + + regionConfigurer.newAsyncEventQueue(mockCache, "TestRegion"); + + verify(mockAsyncEventQueueFactory, times(1)).setParallel(eq(true)); + } + + @Test + public void newAsyncEventQueueIsSerial() { + + AsyncEventQueueFactory mockAsyncEventQueueFactory = mock(AsyncEventQueueFactory.class); + + Cache mockCache = mock(Cache.class); + + doReturn(mockAsyncEventQueueFactory).when(mockCache).createAsyncEventQueueFactory(); + + CrudRepository mockRepository = mock(CrudRepository.class); + + AsyncInlineCachingRegionConfigurer regionConfigurer = + spy(AsyncInlineCachingRegionConfigurer.create(mockRepository, Predicate.isEqual("TestRegion"))); + + assertThat(regionConfigurer).isNotNull(); + assertThat(regionConfigurer.withSerialQueue()).isSameAs(regionConfigurer); + + regionConfigurer.newAsyncEventQueue(mockCache, "TestRegion"); + + verify(mockAsyncEventQueueFactory, times(1)).setParallel(eq(false)); + } + @Test public void newAsyncEventQueueCreatesQueueFromFactoryWithIdAndListener() {