本文整理汇总了Java中org.springframework.messaging.support.ExecutorSubscribableChannel类的典型用法代码示例。如果您正苦于以下问题:Java ExecutorSubscribableChannel类的具体用法?Java ExecutorSubscribableChannel怎么用?Java ExecutorSubscribableChannel使用的例子?那么恭喜您, 这里精选的类代码示例或许可以为您提供帮助。
ExecutorSubscribableChannel类属于org.springframework.messaging.support包,在下文中一共展示了ExecutorSubscribableChannel类的20个代码示例,这些例子默认根据受欢迎程度排序。您可以为喜欢或者感觉有用的代码点赞,您的评价将有助于我们的系统推荐出更棒的Java代码示例。
示例1: run
import org.springframework.messaging.support.ExecutorSubscribableChannel; //导入依赖的package包/类
@Override
public void run(ApplicationArguments args) throws Exception {
logger.info("Consumer running with binder {}", binder);
SubscribableChannel consumerChannel = new ExecutorSubscribableChannel();
consumerChannel.subscribe(new MessageHandler() {
@Override
public void handleMessage(Message<?> message) throws MessagingException {
messagePayload = (String) message.getPayload();
logger.info("Received message: {}", messagePayload);
}
});
String group = null;
if (args.containsOption("group")) {
group = args.getOptionValues("group").get(0);
}
binder.bindConsumer(ConsulBinderTests.BINDING_NAME, group, consumerChannel,
new ConsumerProperties());
isBound = true;
}
开发者ID:spring-cloud,项目名称:spring-cloud-consul,代码行数:22,代码来源:TestConsumer.java
示例2: testGoodConnection
import org.springframework.messaging.support.ExecutorSubscribableChannel; //导入依赖的package包/类
@Test
public void testGoodConnection() throws MqttException
{
StaticApplicationContext applicationContext = getStaticApplicationContext();
MessageChannel inboundMessageChannel = new ExecutorSubscribableChannel();
PahoAsyncMqttClientService service = new PahoAsyncMqttClientService(
BrokerHelper.getProxyUri(), BrokerHelper.getClientId(), MqttClientConnectionType.PUBSUB,
null);
service.setApplicationEventPublisher(applicationContext);
service.setInboundMessageChannel(inboundMessageChannel);
service.subscribe(String.format("client/%s", BrokerHelper.getClientId()),
MqttQualityOfService.QOS_0);
service.getMqttConnectOptions().setCleanSession(true);
Assert.assertTrue(service.start());
Assert.assertTrue(service.isConnected());
Assert.assertTrue(service.isStarted());
Assert.assertEquals(1, clientConnectedCount.get());
Assert.assertEquals(0, clientDisconnectedCount.get());
Assert.assertEquals(0, clientLostConnectionCount.get());
Assert.assertEquals(0, clientFailedConnectionCount.get());
service.stop();
service.close();
applicationContext.close();
}
开发者ID:christophersmith,项目名称:summer-mqtt,代码行数:25,代码来源:AutomaticReconnectTest.java
示例3: clientInboundChannel
import org.springframework.messaging.support.ExecutorSubscribableChannel; //导入依赖的package包/类
@Bean
public SubscribableChannel clientInboundChannel() {
ExecutorSubscribableChannel executorSubscribableChannel = new ExecutorSubscribableChannel(
clientInboundChannelExecutor());
configureClientInboundChannel(executorSubscribableChannel);
for (WampConfigurer wc : this.configurers) {
wc.configureClientInboundChannel(executorSubscribableChannel);
}
return executorSubscribableChannel;
}
开发者ID:ralscha,项目名称:wamp2spring,代码行数:13,代码来源:WampConfiguration.java
示例4: testMqttConnectOptionsAutomaticReconnectDefaultServerAvailableAtStartup
import org.springframework.messaging.support.ExecutorSubscribableChannel; //导入依赖的package包/类
@Test
public void testMqttConnectOptionsAutomaticReconnectDefaultServerAvailableAtStartup()
throws MqttException, InterruptedException
{
StaticApplicationContext applicationContext = getStaticApplicationContext();
MessageChannel inboundMessageChannel = new ExecutorSubscribableChannel();
PahoAsyncMqttClientService service = new PahoAsyncMqttClientService(
BrokerHelper.getProxyUri(), BrokerHelper.getClientId(), MqttClientConnectionType.PUBSUB,
null);
service.setApplicationEventPublisher(applicationContext);
service.setInboundMessageChannel(inboundMessageChannel);
service.subscribe(String.format("client/%s", BrokerHelper.getClientId()),
MqttQualityOfService.QOS_0);
service.getMqttConnectOptions().setCleanSession(true);
Assert.assertTrue(service.start());
Assert.assertTrue(service.isConnected());
Assert.assertTrue(service.isStarted());
// simulate a lost connection
CRUSHER_PROXY.reopen();
Assert.assertFalse(service.isStarted());
Assert.assertFalse(service.isConnected());
Thread.sleep(1000);
Assert.assertFalse(service.isStarted());
Assert.assertFalse(service.isConnected());
Assert.assertEquals(1, clientConnectedCount.get());
Assert.assertEquals(0, clientDisconnectedCount.get());
Assert.assertEquals(1, clientLostConnectionCount.get());
Assert.assertEquals(0, clientFailedConnectionCount.get());
service.stop();
service.close();
applicationContext.close();
}
开发者ID:christophersmith,项目名称:summer-mqtt,代码行数:33,代码来源:AutomaticReconnectTest.java
示例5: testMqttConnectOptionsAutomaticReconnectFalseServerAvailableAtStartup
import org.springframework.messaging.support.ExecutorSubscribableChannel; //导入依赖的package包/类
@Test
public void testMqttConnectOptionsAutomaticReconnectFalseServerAvailableAtStartup()
throws MqttException, InterruptedException
{
StaticApplicationContext applicationContext = getStaticApplicationContext();
MessageChannel inboundMessageChannel = new ExecutorSubscribableChannel();
PahoAsyncMqttClientService service = new PahoAsyncMqttClientService(
BrokerHelper.getProxyUri(), BrokerHelper.getClientId(), MqttClientConnectionType.PUBSUB,
null);
service.setApplicationEventPublisher(applicationContext);
service.setInboundMessageChannel(inboundMessageChannel);
service.subscribe(String.format("client/%s", BrokerHelper.getClientId()),
MqttQualityOfService.QOS_0);
service.getMqttConnectOptions().setCleanSession(true);
service.getMqttConnectOptions().setAutomaticReconnect(false);
Assert.assertTrue(service.start());
Assert.assertTrue(service.isConnected());
Assert.assertTrue(service.isStarted());
// simulate a lost connection
CRUSHER_PROXY.reopen();
Assert.assertFalse(service.isStarted());
Assert.assertFalse(service.isConnected());
Thread.sleep(1000);
Assert.assertFalse(service.isStarted());
Assert.assertFalse(service.isConnected());
Assert.assertEquals(1, clientConnectedCount.get());
Assert.assertEquals(0, clientDisconnectedCount.get());
Assert.assertEquals(1, clientLostConnectionCount.get());
Assert.assertEquals(0, clientFailedConnectionCount.get());
service.stop();
service.close();
applicationContext.close();
}
开发者ID:christophersmith,项目名称:summer-mqtt,代码行数:34,代码来源:AutomaticReconnectTest.java
示例6: testMqttConnectOptionsAutomaticReconnectTrueServerAvailableAtStartup
import org.springframework.messaging.support.ExecutorSubscribableChannel; //导入依赖的package包/类
@Test
public void testMqttConnectOptionsAutomaticReconnectTrueServerAvailableAtStartup()
throws MqttException, InterruptedException
{
StaticApplicationContext applicationContext = getStaticApplicationContext();
MessageChannel inboundMessageChannel = new ExecutorSubscribableChannel();
PahoAsyncMqttClientService service = new PahoAsyncMqttClientService(
BrokerHelper.getProxyUri(), BrokerHelper.getClientId(), MqttClientConnectionType.PUBSUB,
null);
service.setApplicationEventPublisher(applicationContext);
service.setInboundMessageChannel(inboundMessageChannel);
service.subscribe(String.format("client/%s", BrokerHelper.getClientId()),
MqttQualityOfService.QOS_0);
service.getMqttConnectOptions().setCleanSession(true);
service.getMqttConnectOptions().setAutomaticReconnect(true);
Assert.assertTrue(service.start());
Assert.assertTrue(service.isConnected());
Assert.assertTrue(service.isStarted());
// simulate a lost connection
CRUSHER_PROXY.reopen();
Assert.assertFalse(service.isStarted());
Assert.assertFalse(service.isConnected());
Thread.sleep(1100);
Assert.assertTrue(service.isConnected());
Assert.assertTrue(service.isStarted());
Assert.assertEquals(2, clientConnectedCount.get());
Assert.assertEquals(0, clientDisconnectedCount.get());
Assert.assertEquals(1, clientLostConnectionCount.get());
Assert.assertEquals(0, clientFailedConnectionCount.get());
service.stop();
service.close();
applicationContext.close();
}
开发者ID:christophersmith,项目名称:summer-mqtt,代码行数:34,代码来源:AutomaticReconnectTest.java
示例7: testMqttConnectOptionsAutomaticReconnectDefaultServerUnavailableAtStartup
import org.springframework.messaging.support.ExecutorSubscribableChannel; //导入依赖的package包/类
@Test
public void testMqttConnectOptionsAutomaticReconnectDefaultServerUnavailableAtStartup()
throws MqttException, InterruptedException
{
StaticApplicationContext applicationContext = getStaticApplicationContext();
CRUSHER_PROXY.close();
MessageChannel inboundMessageChannel = new ExecutorSubscribableChannel();
PahoAsyncMqttClientService service = new PahoAsyncMqttClientService(
BrokerHelper.getProxyUri(), BrokerHelper.getClientId(), MqttClientConnectionType.PUBSUB,
null);
service.setApplicationEventPublisher(applicationContext);
service.setInboundMessageChannel(inboundMessageChannel);
service.subscribe(String.format("client/%s", BrokerHelper.getClientId()),
MqttQualityOfService.QOS_0);
service.getMqttConnectOptions().setCleanSession(true);
Assert.assertFalse(service.start());
Assert.assertFalse(service.isConnected());
Assert.assertFalse(service.isStarted());
Thread.sleep(1000);
CRUSHER_PROXY.open();
Assert.assertFalse(service.isConnected());
Assert.assertFalse(service.isStarted());
Thread.sleep(1100);
Assert.assertFalse(service.isConnected());
Assert.assertFalse(service.isStarted());
Assert.assertEquals(0, clientConnectedCount.get());
Assert.assertEquals(0, clientDisconnectedCount.get());
Assert.assertEquals(0, clientLostConnectionCount.get());
Assert.assertEquals(1, clientFailedConnectionCount.get());
service.stop();
service.close();
applicationContext.close();
}
开发者ID:christophersmith,项目名称:summer-mqtt,代码行数:34,代码来源:AutomaticReconnectTest.java
示例8: testMqttConnectOptionsAutomaticReconnectFalseServerUnavailableAtStartup
import org.springframework.messaging.support.ExecutorSubscribableChannel; //导入依赖的package包/类
@Test
public void testMqttConnectOptionsAutomaticReconnectFalseServerUnavailableAtStartup()
throws MqttException, InterruptedException
{
StaticApplicationContext applicationContext = getStaticApplicationContext();
CRUSHER_PROXY.close();
MessageChannel inboundMessageChannel = new ExecutorSubscribableChannel();
PahoAsyncMqttClientService service = new PahoAsyncMqttClientService(
BrokerHelper.getProxyUri(), BrokerHelper.getClientId(), MqttClientConnectionType.PUBSUB,
null);
service.setApplicationEventPublisher(applicationContext);
service.setInboundMessageChannel(inboundMessageChannel);
service.subscribe(String.format("client/%s", BrokerHelper.getClientId()),
MqttQualityOfService.QOS_0);
service.getMqttConnectOptions().setCleanSession(true);
service.getMqttConnectOptions().setAutomaticReconnect(false);
Assert.assertFalse(service.start());
Assert.assertFalse(service.isConnected());
Assert.assertFalse(service.isStarted());
Thread.sleep(1000);
CRUSHER_PROXY.open();
Assert.assertFalse(service.isConnected());
Assert.assertFalse(service.isStarted());
Thread.sleep(1100);
Assert.assertFalse(service.isConnected());
Assert.assertFalse(service.isStarted());
Assert.assertEquals(0, clientConnectedCount.get());
Assert.assertEquals(0, clientDisconnectedCount.get());
Assert.assertEquals(0, clientLostConnectionCount.get());
Assert.assertEquals(1, clientFailedConnectionCount.get());
service.stop();
service.close();
applicationContext.close();
}
开发者ID:christophersmith,项目名称:summer-mqtt,代码行数:35,代码来源:AutomaticReconnectTest.java
示例9: testMqttConnectOptionsAutomaticReconnectTrueServerUnavailableAtStartup
import org.springframework.messaging.support.ExecutorSubscribableChannel; //导入依赖的package包/类
@Test
public void testMqttConnectOptionsAutomaticReconnectTrueServerUnavailableAtStartup()
throws MqttException, InterruptedException
{
StaticApplicationContext applicationContext = getStaticApplicationContext();
CRUSHER_PROXY.close();
MessageChannel inboundMessageChannel = new ExecutorSubscribableChannel();
PahoAsyncMqttClientService service = new PahoAsyncMqttClientService(
BrokerHelper.getProxyUri(), BrokerHelper.getClientId(), MqttClientConnectionType.PUBSUB,
null);
service.setApplicationEventPublisher(applicationContext);
service.setInboundMessageChannel(inboundMessageChannel);
service.subscribe(String.format("client/%s", BrokerHelper.getClientId()),
MqttQualityOfService.QOS_0);
service.getMqttConnectOptions().setCleanSession(true);
service.getMqttConnectOptions().setAutomaticReconnect(true);
Assert.assertFalse(service.start());
Assert.assertFalse(service.isConnected());
Assert.assertFalse(service.isStarted());
Thread.sleep(1000);
CRUSHER_PROXY.open();
Assert.assertFalse(service.isConnected());
Assert.assertFalse(service.isStarted());
Thread.sleep(1100);
Assert.assertFalse(service.isStarted());
Assert.assertFalse(service.isConnected());
Thread.sleep(1100);
Assert.assertFalse(service.isConnected());
Assert.assertFalse(service.isStarted());
Assert.assertEquals(0, clientConnectedCount.get());
Assert.assertEquals(0, clientDisconnectedCount.get());
Assert.assertEquals(0, clientLostConnectionCount.get());
Assert.assertEquals(1, clientFailedConnectionCount.get());
service.stop();
service.close();
applicationContext.close();
}
开发者ID:christophersmith,项目名称:summer-mqtt,代码行数:38,代码来源:AutomaticReconnectTest.java
示例10: testReconnectDetailsSetServerAvailableAtStartup
import org.springframework.messaging.support.ExecutorSubscribableChannel; //导入依赖的package包/类
@Test
public void testReconnectDetailsSetServerAvailableAtStartup()
throws MqttException, InterruptedException
{
StaticApplicationContext applicationContext = getStaticApplicationContext();
TaskScheduler taskScheduler = new ConcurrentTaskScheduler();
MessageChannel inboundMessageChannel = new ExecutorSubscribableChannel();
PahoAsyncMqttClientService service = new PahoAsyncMqttClientService(
BrokerHelper.getProxyUri(), BrokerHelper.getClientId(), MqttClientConnectionType.PUBSUB,
null);
service.setApplicationEventPublisher(applicationContext);
service.setInboundMessageChannel(inboundMessageChannel);
service.subscribe(String.format("client/%s", BrokerHelper.getClientId()),
MqttQualityOfService.QOS_0);
service.getMqttConnectOptions().setCleanSession(true);
service.setReconnectDetails(new DefaultReconnectService(), taskScheduler);
Assert.assertTrue(service.start());
Assert.assertTrue(service.isConnected());
Assert.assertTrue(service.isStarted());
// simulate a lost connection
CRUSHER_PROXY.reopen();
Assert.assertFalse(service.isStarted());
Assert.assertFalse(service.isConnected());
Thread.sleep(1100);
Assert.assertTrue(service.isStarted());
Assert.assertTrue(service.isConnected());
Assert.assertEquals(2, clientConnectedCount.get());
Assert.assertEquals(0, clientDisconnectedCount.get());
Assert.assertEquals(1, clientLostConnectionCount.get());
Assert.assertEquals(0, clientFailedConnectionCount.get());
service.stop();
service.close();
applicationContext.close();
}
开发者ID:christophersmith,项目名称:summer-mqtt,代码行数:35,代码来源:AutomaticReconnectTest.java
示例11: testReconnectDetailsSetServerUnavailableAtStartup
import org.springframework.messaging.support.ExecutorSubscribableChannel; //导入依赖的package包/类
@Test
public void testReconnectDetailsSetServerUnavailableAtStartup()
throws MqttException, InterruptedException
{
StaticApplicationContext applicationContext = getStaticApplicationContext();
CRUSHER_PROXY.close();
TaskScheduler taskScheduler = new ConcurrentTaskScheduler();
MessageChannel inboundMessageChannel = new ExecutorSubscribableChannel();
PahoAsyncMqttClientService service = new PahoAsyncMqttClientService(
BrokerHelper.getProxyUri(), BrokerHelper.getClientId(), MqttClientConnectionType.PUBSUB,
null);
service.setApplicationEventPublisher(applicationContext);
service.setInboundMessageChannel(inboundMessageChannel);
service.subscribe(String.format("client/%s", BrokerHelper.getClientId()),
MqttQualityOfService.QOS_0);
service.getMqttConnectOptions().setCleanSession(true);
service.setReconnectDetails(new DefaultReconnectService(), taskScheduler);
Assert.assertFalse(service.start());
Assert.assertFalse(service.isConnected());
Assert.assertFalse(service.isStarted());
CRUSHER_PROXY.open();
Thread.sleep(3100);
Assert.assertTrue(service.isConnected());
Assert.assertTrue(service.isStarted());
Assert.assertEquals(1, clientConnectedCount.get());
Assert.assertEquals(0, clientDisconnectedCount.get());
Assert.assertEquals(0, clientLostConnectionCount.get());
Assert.assertEquals(1, clientFailedConnectionCount.get());
service.stop();
service.close();
applicationContext.close();
}
开发者ID:christophersmith,项目名称:summer-mqtt,代码行数:33,代码来源:AutomaticReconnectTest.java
示例12: clientInboundChannel
import org.springframework.messaging.support.ExecutorSubscribableChannel; //导入依赖的package包/类
@Bean
public AbstractSubscribableChannel clientInboundChannel() {
ExecutorSubscribableChannel channel = new ExecutorSubscribableChannel(clientInboundChannelExecutor());
ChannelRegistration reg = getClientInboundChannelRegistration();
channel.setInterceptors(reg.getInterceptors());
return channel;
}
开发者ID:langtianya,项目名称:spring4-understanding,代码行数:8,代码来源:AbstractMessageBrokerConfiguration.java
示例13: clientOutboundChannel
import org.springframework.messaging.support.ExecutorSubscribableChannel; //导入依赖的package包/类
@Bean
public AbstractSubscribableChannel clientOutboundChannel() {
ExecutorSubscribableChannel channel = new ExecutorSubscribableChannel(clientOutboundChannelExecutor());
ChannelRegistration reg = getClientOutboundChannelRegistration();
channel.setInterceptors(reg.getInterceptors());
return channel;
}
开发者ID:langtianya,项目名称:spring4-understanding,代码行数:8,代码来源:AbstractMessageBrokerConfiguration.java
示例14: brokerChannel
import org.springframework.messaging.support.ExecutorSubscribableChannel; //导入依赖的package包/类
@Bean
public AbstractSubscribableChannel brokerChannel() {
ChannelRegistration reg = getBrokerRegistry().getBrokerChannelRegistration();
ExecutorSubscribableChannel channel = reg.hasTaskExecutor() ?
new ExecutorSubscribableChannel(brokerChannelExecutor()) : new ExecutorSubscribableChannel();
reg.setInterceptors(new ImmutableMessageChannelInterceptor());
channel.setInterceptors(reg.getInterceptors());
return channel;
}
开发者ID:langtianya,项目名称:spring4-understanding,代码行数:10,代码来源:AbstractMessageBrokerConfiguration.java
示例15: setUp
import org.springframework.messaging.support.ExecutorSubscribableChannel; //导入依赖的package包/类
@Before
public void setUp() throws Exception {
logger.debug("Setting up before '" + this.testName.getMethodName() + "'");
this.port = SocketUtils.findAvailableTcpPort(61613);
this.responseChannel = new ExecutorSubscribableChannel();
this.responseHandler = new TestMessageHandler();
this.responseChannel.subscribe(this.responseHandler);
this.eventPublisher = new TestEventPublisher();
startActiveMqBroker();
createAndStartRelay();
}
开发者ID:langtianya,项目名称:spring4-understanding,代码行数:12,代码来源:StompBrokerRelayMessageHandlerIntegrationTests.java
示例16: sendAndReceive
import org.springframework.messaging.support.ExecutorSubscribableChannel; //导入依赖的package包/类
@Test
public void sendAndReceive() {
SubscribableChannel channel = new ExecutorSubscribableChannel(this.executor);
channel.subscribe(new MessageHandler() {
@Override
public void handleMessage(Message<?> message) throws MessagingException {
MessageChannel replyChannel = (MessageChannel) message.getHeaders().getReplyChannel();
replyChannel.send(new GenericMessage<String>("response"));
}
});
String actual = this.template.convertSendAndReceive(channel, "request", String.class);
assertEquals("response", actual);
}
开发者ID:langtianya,项目名称:spring4-understanding,代码行数:16,代码来源:GenericMessagingTemplateTests.java
示例17: clientInboundChannel
import org.springframework.messaging.support.ExecutorSubscribableChannel; //导入依赖的package包/类
/**
* Channel for inbound messages between {@link WampSubProtocolHandler},
* {@link #brokerMessageHandler()} and {@link #annotationMethodMessageHandler()}
*/
@Bean
public SubscribableChannel clientInboundChannel() {
ExecutorSubscribableChannel executorSubscribableChannel = new ExecutorSubscribableChannel(
clientInboundChannelExecutor());
configureClientInboundChannel(executorSubscribableChannel);
return executorSubscribableChannel;
}
开发者ID:ralscha,项目名称:wampspring,代码行数:12,代码来源:DefaultWampConfiguration.java
示例18: configureClientInboundChannel
import org.springframework.messaging.support.ExecutorSubscribableChannel; //导入依赖的package包/类
protected void configureClientInboundChannel(
@SuppressWarnings("unused") ExecutorSubscribableChannel executorSubscribableChannel) {
// nothing here
}
开发者ID:ralscha,项目名称:wamp2spring,代码行数:5,代码来源:WampConfiguration.java
示例19: clientOutboundChannel
import org.springframework.messaging.support.ExecutorSubscribableChannel; //导入依赖的package包/类
@Bean
public SubscribableChannel clientOutboundChannel() {
return new ExecutorSubscribableChannel(clientOutboundChannelExecutor());
}
开发者ID:ralscha,项目名称:wamp2spring,代码行数:5,代码来源:WampConfiguration.java
示例20: brokerChannel
import org.springframework.messaging.support.ExecutorSubscribableChannel; //导入依赖的package包/类
/**
* Channel from the {@link WampPublisher} to the {@link PubSubMessageHandler}
*/
@Bean
public SubscribableChannel brokerChannel() {
return new ExecutorSubscribableChannel(brokerChannelExecutor());
}
开发者ID:ralscha,项目名称:wamp2spring,代码行数:8,代码来源:WampConfiguration.java
注:本文中的org.springframework.messaging.support.ExecutorSubscribableChannel类示例整理自Github/MSDocs等源码及文档管理平台,相关代码片段筛选自各路编程大神贡献的开源项目,源码版权归原作者所有,传播和使用请参考对应项目的License;未经允许,请勿转载。 |
请发表评论