diff --git a/org.springframework.integration/src/main/java/org/springframework/integration/channel/PublishSubscribeChannel.java b/org.springframework.integration/src/main/java/org/springframework/integration/channel/PublishSubscribeChannel.java index 98630a6888..e07a9a0ee6 100644 --- a/org.springframework.integration/src/main/java/org/springframework/integration/channel/PublishSubscribeChannel.java +++ b/org.springframework.integration/src/main/java/org/springframework/integration/channel/PublishSubscribeChannel.java @@ -36,6 +36,8 @@ public class PublishSubscribeChannel extends AbstractSubscribableChannel impleme private volatile ErrorHandler errorHandler; + private volatile boolean applySequence; + /** * Create a PublishSubscribeChannel that will use a {@link TaskExecutor} @@ -61,6 +63,7 @@ public class PublishSubscribeChannel extends AbstractSubscribableChannel impleme } public void setApplySequence(boolean applySequence) { + this.applySequence = applySequence; this.getDispatcher().setApplySequence(applySequence); } @@ -73,6 +76,7 @@ public class PublishSubscribeChannel extends AbstractSubscribableChannel impleme this.taskExecutor = new ErrorHandlingTaskExecutor(this.taskExecutor, this.errorHandler); } this.dispatcher = new BroadcastingDispatcher(this.taskExecutor); + this.dispatcher.setApplySequence(this.applySequence); } } diff --git a/org.springframework.integration/src/test/java/org/springframework/integration/config/PublishSubscribeChannelParserTests.java b/org.springframework.integration/src/test/java/org/springframework/integration/config/PublishSubscribeChannelParserTests.java index 9edefc9b89..9b2a80c6e4 100644 --- a/org.springframework.integration/src/test/java/org/springframework/integration/config/PublishSubscribeChannelParserTests.java +++ b/org.springframework.integration/src/test/java/org/springframework/integration/config/PublishSubscribeChannelParserTests.java @@ -81,6 +81,25 @@ public class PublishSubscribeChannelParserTests { assertEquals(context.getBean("pool"), innerExecutor); } + @Test + public void applySequenceEnabledWithTaskExecutor() { + ClassPathXmlApplicationContext context = new ClassPathXmlApplicationContext( + "publishSubscribeChannelParserTests.xml", this.getClass()); + PublishSubscribeChannel channel = (PublishSubscribeChannel) + context.getBean("channelWithApplySequenceEnabledAndTaskExecutor"); + DirectFieldAccessor accessor = new DirectFieldAccessor(channel); + BroadcastingDispatcher dispatcher = (BroadcastingDispatcher) + accessor.getPropertyValue("dispatcher"); + DirectFieldAccessor dispatcherAccessor = new DirectFieldAccessor(dispatcher); + assertTrue((Boolean) dispatcherAccessor.getPropertyValue("applySequence")); + TaskExecutor executor = (TaskExecutor) dispatcherAccessor.getPropertyValue("taskExecutor"); + assertNotNull(executor); + assertEquals(ErrorHandlingTaskExecutor.class, executor.getClass()); + DirectFieldAccessor executorAccessor = new DirectFieldAccessor(executor); + TaskExecutor innerExecutor = (TaskExecutor) executorAccessor.getPropertyValue("taskExecutor"); + assertEquals(context.getBean("pool"), innerExecutor); + } + @Test public void channelWithErrorHandler() { ClassPathXmlApplicationContext context = new ClassPathXmlApplicationContext( diff --git a/org.springframework.integration/src/test/java/org/springframework/integration/config/publishSubscribeChannelParserTests.xml b/org.springframework.integration/src/test/java/org/springframework/integration/config/publishSubscribeChannelParserTests.xml index 2d448e84b5..4ab0788254 100644 --- a/org.springframework.integration/src/test/java/org/springframework/integration/config/publishSubscribeChannelParserTests.xml +++ b/org.springframework.integration/src/test/java/org/springframework/integration/config/publishSubscribeChannelParserTests.xml @@ -13,6 +13,8 @@ + +