【问题标题】:Spring MQTT JAva Config example issueSpring MQTT JAva Config 示例问题
【发布时间】:2016-03-05 08:58:13
【问题描述】:

在哪里可以找到如何使用 MQTT + JAva Config 的示例?

这对我不起作用:http://docs.spring.io/spring-integration/reference/html/mqtt.html

【问题讨论】:

  • 什么不起作用?

标签: spring spring-integration mqtt


【解决方案1】:

使用 Spring Boot 解决

@Configuration
@ComponentScan
@EnableAutoConfiguration
@IntegrationComponentScan
public class Application {

  public static void main(String[] args) {
    SpringApplication.run(Application.class, args);
  }

  @Bean
  public MessageChannel mqttInputChannel() {
    return new DirectChannel();
  }

  @Bean
  public MqttPahoClientFactory mqttClientFactory() {
    DefaultMqttPahoClientFactory factory = new DefaultMqttPahoClientFactory();
    factory.setServerURIs("tcp://url:10423");
    factory.setUserName("username");
    factory.setPassword("password");
    return factory;
  }

  @Bean
  public MessageProducer inbound() {
    MqttPahoMessageDrivenChannelAdapter adapter =
            new MqttPahoMessageDrivenChannelAdapter("testMqtt", mqttClientFactory(),
                    "test");
    adapter.setCompletionTimeout(5000);
    adapter.setConverter(new DefaultPahoMessageConverter());
    adapter.setQos(1);
    adapter.setOutputChannel(mqttInputChannel());
    return adapter;
  }

  @Bean
  @ServiceActivator(inputChannel = "mqttInputChannel")
  public MessageHandler handler() {
    return new MessageHandler() {

      @Override
      public void handleMessage(Message<?> message) throws MessagingException {
        System.out.println("!!!!!!!!!!!!!!!!!!!" + message.getPayload());
      }

    };
  }

  @Bean
  @ServiceActivator(inputChannel = "mqttOutboundChannel")
  public MessageHandler mqttOutbound() {
    MqttPahoMessageHandler messageHandler =
            new MqttPahoMessageHandler("testClient", mqttClientFactory());
    messageHandler.setAsync(true);
    messageHandler.setDefaultTopic("test");
    return messageHandler;
  }

  @Bean
  public MessageChannel mqttOutboundChannel() {
    return new DirectChannel();
  }

  @MessagingGateway(defaultRequestChannel = "mqttOutboundChannel")
  public interface MyGateway {

    void sendToMqtt(String data);

  }

}

【讨论】:

    【解决方案2】:

    很高兴您找到了解决方案。

    我创建了一个示例应用程序,它使用 Java DSL 从标准输入读取、发送到 MQTT、接收和记录。

    以下是相关的位:

    // publisher
    
    @Bean
    public IntegrationFlow mqttOutFlow() {
        return IntegrationFlows.from(CharacterStreamReadingMessageSource.stdin(),
                        e -> e.poller(Pollers.fixedDelay(1000)))
                .transform(p -> p + " sent to MQTT")
                .handle(mqttOutbound())
                .get();
    }
    
    @Bean
    public MessageHandler mqttOutbound() {
        MqttPahoMessageHandler messageHandler = new MqttPahoMessageHandler("siSamplePublisher", mqttClientFactory());
        messageHandler.setAsync(true);
        messageHandler.setDefaultTopic("siSampleTopic");
        return messageHandler;
    }
    
    // consumer
    
    @Bean
    public IntegrationFlow mqttInFlow() {
        return IntegrationFlows.from(mqttInbound())
                .transform(p -> p + ", received from MQTT")
                .handle(logger())
                .get();
    }
    
    private LoggingHandler logger() {
        LoggingHandler loggingHandler = new LoggingHandler("INFO");
        loggingHandler.setLoggerName("siSample");
        return loggingHandler;
    }
    
    @Bean
    public MessageProducerSupport mqttInbound() {
        MqttPahoMessageDrivenChannelAdapter adapter = new MqttPahoMessageDrivenChannelAdapter("siSampleConsumer",
                mqttClientFactory(), "siSampleTopic");
        adapter.setCompletionTimeout(5000);
        adapter.setConverter(new DefaultPahoMessageConverter());
        adapter.setQos(1);
        return adapter;
    }
    

    .

    foo
    14:40:56.770 [MQTT Call: siSampleConsumer] INFO  siSample - foo sent to MQTT, received from MQTT
    

    编辑

    带有注释和 DSL 配置的官方 Spring Integration MQTT 示例位于:https://github.com/spring-projects/spring-integration-samples/tree/master/basic/mqtt

    【讨论】:

    • 在 Gary 的回答中查看我的编辑。
    【解决方案3】:

    我正在尝试。 http://docs.spring.io/spring-integration/reference/html/mqtt.html
    运行良好。我的来源是这个

    import org.eclipse.paho.client.mqttv3.MqttConnectOptions;
    import org.springframework.boot.SpringApplication;
    import org.springframework.boot.autoconfigure.SpringBootApplication;
    import org.springframework.context.ConfigurableApplicationContext;
    import org.springframework.context.annotation.Bean;
    import org.springframework.integration.annotation.MessagingGateway;
    import org.springframework.integration.annotation.ServiceActivator;
    import org.springframework.integration.channel.DirectChannel;
    import org.springframework.integration.core.MessageProducer;
    import org.springframework.integration.mqtt.core.DefaultMqttPahoClientFactory;
    import org.springframework.integration.mqtt.core.MqttPahoClientFactory;
    import org.springframework.integration.mqtt.inbound.MqttPahoMessageDrivenChannelAdapter;
    import org.springframework.integration.mqtt.outbound.MqttPahoMessageHandler;
    import org.springframework.integration.mqtt.support.DefaultPahoMessageConverter;
    import org.springframework.messaging.Message;
    import org.springframework.messaging.MessageChannel;
    import org.springframework.messaging.MessageHandler;
    import org.springframework.messaging.MessagingException;
    
    @SpringBootApplication
    public class MqttApplication {
    
        public static void main(String[] args) {
            ConfigurableApplicationContext context = SpringApplication.run(MqttApplication.class, args);
    
            MyGateway gateway = context.getBean(MyGateway.class);
            gateway.sendToMqtt("foo");
        }
    
        @Bean
        public MessageChannel mqttInputChannel() {
            return new DirectChannel();
        }
    
        @Bean
        public MqttPahoClientFactory mqttClientFactory() {
            DefaultMqttPahoClientFactory factory = new DefaultMqttPahoClientFactory();
            MqttConnectOptions options = new MqttConnectOptions();
            options.setServerURIs(new String[] { "tcp://localhost:1883" });
            factory.setConnectionOptions(options);
            return factory;
        }
    
        @Bean
        public MessageProducer inbound() {
            MqttPahoMessageDrivenChannelAdapter adapter = new MqttPahoMessageDrivenChannelAdapter("testMqtt",
                mqttClientFactory(), "testTopic");
            adapter.setCompletionTimeout(5000);
            adapter.setConverter(new DefaultPahoMessageConverter());
            adapter.setQos(1);
            adapter.setOutputChannel(mqttInputChannel());
            return adapter;
        }
    
        @Bean
        @ServiceActivator(inputChannel = "mqttInputChannel")
        public MessageHandler handler() {
            return new MessageHandler() {
    
                @Override
                public void handleMessage(Message<?> message) throws MessagingException {
                    System.out.println(message.getPayload());
                }
    
            };
        }
    
        @Bean
        @ServiceActivator(inputChannel = "mqttOutboundChannel")
        public MessageHandler mqttOutbound() {
            MqttPahoMessageHandler messageHandler = new MqttPahoMessageHandler("testClient", mqttClientFactory());
            messageHandler.setAsync(true);
            messageHandler.setDefaultTopic("testTopic");
            return messageHandler;
        }
    
        @Bean
        public MessageChannel mqttOutboundChannel() {
            return new DirectChannel();
        }
    
        @MessagingGateway(defaultRequestChannel = "mqttOutboundChannel")
        public interface MyGateway {
    
            void sendToMqtt(String data);
    
        }
    
    }
    

    【讨论】:

      猜你喜欢
      • 2014-07-05
      • 1970-01-01
      • 1970-01-01
      • 2016-09-29
      • 2015-05-18
      • 1970-01-01
      • 2020-09-05
      • 2019-01-02
      • 2022-01-14
      相关资源
      最近更新 更多