Add 'withSerialQueue()' builder method to configure the AEQ as a serial queue.

This commit is contained in:
John Blum
2020-12-01 15:50:47 -08:00
parent cd28eb19a8
commit c25100ebe7
2 changed files with 62 additions and 5 deletions

View File

@@ -514,9 +514,10 @@ public class AsyncInlineCachingRegionConfigurer<T, ID> 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<T, ID> withParallelQueue() {
this.parallel = true;
@@ -738,4 +739,18 @@ public class AsyncInlineCachingRegionConfigurer<T, ID> 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<T, ID> withSerialQueue() {
this.parallel = false;
return this;
}
}

View File

@@ -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() {