Update thread pool settings in STOMP/WebSocket config
The clientInboundChannel and clientOutboundChannel now use twice the number of available processors by default to accomodate for some degree of blocking in task execution on average. In practice these settings still need to be configured explicitly in applications but these should serve as better default values than the default values in ThreadPoolTaskExecutor. Issue: SPR-11450
This commit is contained in:
@@ -31,6 +31,7 @@ import org.springframework.context.annotation.Bean;
|
||||
import org.springframework.context.annotation.Configuration;
|
||||
import org.springframework.context.support.StaticApplicationContext;
|
||||
import org.springframework.messaging.Message;
|
||||
import org.springframework.messaging.MessageChannel;
|
||||
import org.springframework.messaging.MessageHandler;
|
||||
import org.springframework.messaging.converter.*;
|
||||
import org.springframework.messaging.handler.annotation.MessageMapping;
|
||||
@@ -66,24 +67,30 @@ import static org.junit.Assert.*;
|
||||
*/
|
||||
public class MessageBrokerConfigurationTests {
|
||||
|
||||
private AnnotationConfigApplicationContext simpleContext;
|
||||
private AnnotationConfigApplicationContext simpleBrokerContext;
|
||||
|
||||
private AnnotationConfigApplicationContext brokerRelayContext;
|
||||
|
||||
private AnnotationConfigApplicationContext defaultContext;
|
||||
|
||||
private AnnotationConfigApplicationContext customChannelContext;
|
||||
|
||||
|
||||
@Before
|
||||
public void setupOnce() {
|
||||
|
||||
this.simpleContext = new AnnotationConfigApplicationContext();
|
||||
this.simpleContext.register(SimpleConfig.class);
|
||||
this.simpleContext.refresh();
|
||||
this.simpleBrokerContext = new AnnotationConfigApplicationContext();
|
||||
this.simpleBrokerContext.register(SimpleBrokerConfig.class);
|
||||
this.simpleBrokerContext.refresh();
|
||||
|
||||
this.brokerRelayContext = new AnnotationConfigApplicationContext();
|
||||
this.brokerRelayContext.register(BrokerRelayConfig.class);
|
||||
this.brokerRelayContext.refresh();
|
||||
|
||||
this.defaultContext = new AnnotationConfigApplicationContext();
|
||||
this.defaultContext.register(DefaultConfig.class);
|
||||
this.defaultContext.refresh();
|
||||
|
||||
this.customChannelContext = new AnnotationConfigApplicationContext();
|
||||
this.customChannelContext.register(CustomChannelConfig.class);
|
||||
this.customChannelContext.refresh();
|
||||
@@ -93,13 +100,13 @@ public class MessageBrokerConfigurationTests {
|
||||
@Test
|
||||
public void clientInboundChannel() {
|
||||
|
||||
TestChannel channel = this.simpleContext.getBean("clientInboundChannel", TestChannel.class);
|
||||
TestChannel channel = this.simpleBrokerContext.getBean("clientInboundChannel", TestChannel.class);
|
||||
Set<MessageHandler> handlers = channel.getSubscribers();
|
||||
|
||||
assertEquals(3, handlers.size());
|
||||
assertTrue(handlers.contains(simpleContext.getBean(SimpAnnotationMethodMessageHandler.class)));
|
||||
assertTrue(handlers.contains(simpleContext.getBean(UserDestinationMessageHandler.class)));
|
||||
assertTrue(handlers.contains(simpleContext.getBean(SimpleBrokerMessageHandler.class)));
|
||||
assertTrue(handlers.contains(simpleBrokerContext.getBean(SimpAnnotationMethodMessageHandler.class)));
|
||||
assertTrue(handlers.contains(simpleBrokerContext.getBean(UserDestinationMessageHandler.class)));
|
||||
assertTrue(handlers.contains(simpleBrokerContext.getBean(SimpleBrokerMessageHandler.class)));
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -130,8 +137,8 @@ public class MessageBrokerConfigurationTests {
|
||||
|
||||
@Test
|
||||
public void clientOutboundChannelUsedByAnnotatedMethod() {
|
||||
TestChannel channel = this.simpleContext.getBean("clientOutboundChannel", TestChannel.class);
|
||||
SimpAnnotationMethodMessageHandler messageHandler = this.simpleContext.getBean(SimpAnnotationMethodMessageHandler.class);
|
||||
TestChannel channel = this.simpleBrokerContext.getBean("clientOutboundChannel", TestChannel.class);
|
||||
SimpAnnotationMethodMessageHandler messageHandler = this.simpleBrokerContext.getBean(SimpAnnotationMethodMessageHandler.class);
|
||||
|
||||
StompHeaderAccessor headers = StompHeaderAccessor.create(StompCommand.SUBSCRIBE);
|
||||
headers.setSessionId("sess1");
|
||||
@@ -151,8 +158,8 @@ public class MessageBrokerConfigurationTests {
|
||||
|
||||
@Test
|
||||
public void clientOutboundChannelUsedBySimpleBroker() {
|
||||
TestChannel channel = this.simpleContext.getBean("clientOutboundChannel", TestChannel.class);
|
||||
SimpleBrokerMessageHandler broker = this.simpleContext.getBean(SimpleBrokerMessageHandler.class);
|
||||
TestChannel channel = this.simpleBrokerContext.getBean("clientOutboundChannel", TestChannel.class);
|
||||
SimpleBrokerMessageHandler broker = this.simpleBrokerContext.getBean(SimpleBrokerMessageHandler.class);
|
||||
|
||||
StompHeaderAccessor headers = StompHeaderAccessor.create(StompCommand.SUBSCRIBE);
|
||||
headers.setSessionId("sess1");
|
||||
@@ -197,12 +204,14 @@ public class MessageBrokerConfigurationTests {
|
||||
|
||||
@Test
|
||||
public void brokerChannel() {
|
||||
TestChannel channel = this.simpleContext.getBean("brokerChannel", TestChannel.class);
|
||||
TestChannel channel = this.simpleBrokerContext.getBean("brokerChannel", TestChannel.class);
|
||||
Set<MessageHandler> handlers = channel.getSubscribers();
|
||||
|
||||
assertEquals(2, handlers.size());
|
||||
assertTrue(handlers.contains(simpleContext.getBean(UserDestinationMessageHandler.class)));
|
||||
assertTrue(handlers.contains(simpleContext.getBean(SimpleBrokerMessageHandler.class)));
|
||||
assertTrue(handlers.contains(simpleBrokerContext.getBean(UserDestinationMessageHandler.class)));
|
||||
assertTrue(handlers.contains(simpleBrokerContext.getBean(SimpleBrokerMessageHandler.class)));
|
||||
|
||||
assertNull(channel.getExecutor());
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -217,8 +226,9 @@ public class MessageBrokerConfigurationTests {
|
||||
|
||||
@Test
|
||||
public void brokerChannelUsedByAnnotatedMethod() {
|
||||
TestChannel channel = this.simpleContext.getBean("brokerChannel", TestChannel.class);
|
||||
SimpAnnotationMethodMessageHandler messageHandler = this.simpleContext.getBean(SimpAnnotationMethodMessageHandler.class);
|
||||
TestChannel channel = this.simpleBrokerContext.getBean("brokerChannel", TestChannel.class);
|
||||
SimpAnnotationMethodMessageHandler messageHandler =
|
||||
this.simpleBrokerContext.getBean(SimpAnnotationMethodMessageHandler.class);
|
||||
|
||||
StompHeaderAccessor headers = StompHeaderAccessor.create(StompCommand.SEND);
|
||||
headers.setDestination("/foo");
|
||||
@@ -236,10 +246,10 @@ public class MessageBrokerConfigurationTests {
|
||||
|
||||
@Test
|
||||
public void brokerChannelUsedByUserDestinationMessageHandler() {
|
||||
TestChannel channel = this.simpleContext.getBean("brokerChannel", TestChannel.class);
|
||||
UserDestinationMessageHandler messageHandler = this.simpleContext.getBean(UserDestinationMessageHandler.class);
|
||||
TestChannel channel = this.simpleBrokerContext.getBean("brokerChannel", TestChannel.class);
|
||||
UserDestinationMessageHandler messageHandler = this.simpleBrokerContext.getBean(UserDestinationMessageHandler.class);
|
||||
|
||||
this.simpleContext.getBean(UserSessionRegistry.class).registerSessionId("joe", "s1");
|
||||
this.simpleBrokerContext.getBean(UserSessionRegistry.class).registerSessionId("joe", "s1");
|
||||
|
||||
StompHeaderAccessor headers = StompHeaderAccessor.create(StompCommand.SEND);
|
||||
headers.setDestination("/user/joe/foo");
|
||||
@@ -285,6 +295,24 @@ public class MessageBrokerConfigurationTests {
|
||||
assertEquals(MimeTypeUtils.APPLICATION_JSON, ((DefaultContentTypeResolver) resolver).getDefaultMimeType());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void threadPoolSizeDefault() {
|
||||
|
||||
String name = "clientInboundChannelExecutor";
|
||||
ThreadPoolTaskExecutor executor = this.defaultContext.getBean(name, ThreadPoolTaskExecutor.class);
|
||||
assertEquals(Runtime.getRuntime().availableProcessors() * 2, executor.getCorePoolSize());
|
||||
// No way to verify queue capacity
|
||||
|
||||
name = "clientOutboundChannelExecutor";
|
||||
executor = this.defaultContext.getBean(name, ThreadPoolTaskExecutor.class);
|
||||
assertEquals(Runtime.getRuntime().availableProcessors() * 2, executor.getCorePoolSize());
|
||||
|
||||
name = "brokerChannelExecutor";
|
||||
executor = this.defaultContext.getBean(name, ThreadPoolTaskExecutor.class);
|
||||
assertEquals(0, executor.getCorePoolSize());
|
||||
assertEquals(1, executor.getMaxPoolSize());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void configureMessageConvertersCustom() {
|
||||
final MessageConverter testConverter = Mockito.mock(MessageConverter.class);
|
||||
@@ -360,7 +388,7 @@ public class MessageBrokerConfigurationTests {
|
||||
@Test
|
||||
public void simpValidatorInjected() {
|
||||
SimpAnnotationMethodMessageHandler messageHandler =
|
||||
this.simpleContext.getBean(SimpAnnotationMethodMessageHandler.class);
|
||||
this.simpleBrokerContext.getBean(SimpAnnotationMethodMessageHandler.class);
|
||||
|
||||
assertThat(messageHandler.getValidator(), Matchers.notNullValue(Validator.class));
|
||||
}
|
||||
@@ -381,8 +409,9 @@ public class MessageBrokerConfigurationTests {
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@Configuration
|
||||
static class SimpleConfig extends AbstractMessageBrokerConfiguration {
|
||||
static class SimpleBrokerConfig extends AbstractMessageBrokerConfiguration {
|
||||
|
||||
@Bean
|
||||
public TestController subscriptionController() {
|
||||
@@ -409,7 +438,7 @@ public class MessageBrokerConfigurationTests {
|
||||
}
|
||||
|
||||
@Configuration
|
||||
static class BrokerRelayConfig extends SimpleConfig {
|
||||
static class BrokerRelayConfig extends SimpleBrokerConfig {
|
||||
|
||||
@Override
|
||||
public void configureMessageBroker(MessageBrokerRegistry registry) {
|
||||
@@ -417,6 +446,10 @@ public class MessageBrokerConfigurationTests {
|
||||
}
|
||||
}
|
||||
|
||||
@Configuration
|
||||
static class DefaultConfig extends AbstractMessageBrokerConfiguration {
|
||||
}
|
||||
|
||||
@Configuration
|
||||
static class CustomChannelConfig extends AbstractMessageBrokerConfiguration {
|
||||
|
||||
|
||||
Reference in New Issue
Block a user