You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Spring Integration入站/出站通道适配器区别及MQTT解码实现咨询

问题背景与疑问

我的项目采用MQTT协议、RabbitMQ作为Broker,基于Spring Integration实现,流程如下:
RabbitMQ(数据源) --> 入站流(将Base64数据转为JsonObject) --> 出站流(发送数据至另一topic)

RabbitMQ中的示例数据为eyJrZXkiOiJIZWxsbyBXb3JsZCJ9,解码后内容为{"key":"Hello World"}

我对入站和出站通道适配器的理解是:入站适配器用于从外部系统获取数据,数据的终点是出站适配器。现存在以下问题:

  1. 此场景中RabbitMQ是否属于外部系统?
  2. 如何通过IntegrationFlow.transform实现Base64数据与JsonObject的编解码?
  3. 我对Spring Integration的上述理解是否正确?

附相关代码:

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


    @Value("${mqtt.username}")
    private String mqttBrokerUsername;

    @Value("${mqtt.password}")
    private String mqttBrokerPassword;
    @Bean
    MqttPahoClientFactory clientFactory(@Value("${mqtt.brokerUrl}") String host){
        var factory = new DefaultMqttPahoClientFactory();
        var options = new MqttConnectOptions();
        options.setServerURIs(new String[]{host});
        options.setUserName(mqttBrokerUsername);
        options.setPassword(mqttBrokerPassword.toCharArray());
        factory.setConnectionOptions(options);
        return factory;
    }

    @Bean
    MessageChannel integrationMessageChannels(){
        return MessageChannels.direct().getObject();
    }

    @Bean
    MqttPahoMessageDrivenChannelAdapter inboundAdapter(@Value("${mqtt.decodedTopic}") String topic, MqttPahoClientFactory factory){
       var adapter = new MqttPahoMessageDrivenChannelAdapter("consumerClientID", factory, topic);
        adapter.setCompletionTimeout(5000);
        adapter.setConverter(new DefaultPahoMessageConverter());
        adapter.setQos(1);
        adapter.setOutputChannel(integrationMessageChannels());

        return adapter;
    }
    @Bean
    IntegrationFlow inboundFlow(MqttPahoMessageDrivenChannelAdapter inboundAdapter){
        return IntegrationFlow
                .from(inboundAdapter)
                .transform( ... )  // How to implement decoder
                .handle((payload, headers) -> {
                    System.out.println(payload);
                    return null;
                })
                .get();
    }

    @Bean
    MqttPahoMessageHandler outboundAdapter(@Value("${mqtt.encodedTopic}") String topic, MqttPahoClientFactory factory){
        var mh = new MqttPahoMessageHandler("producerClientID", factory);
        mh.setDefaultTopic(topic);
        return mh;
    }

    @Bean
    IntegrationFlow outboundFlow(MessageChannel integrationMessageChannels, MqttPahoMessageHandler outboundAdapter){
        return IntegrationFlow
                .from(integrationMessageChannels)
                .handle(outboundAdapter)
                .get();
    }


}

问题解答

1. RabbitMQ是否属于外部系统?

是的,此场景中RabbitMQ属于外部系统。Spring Integration的入站适配器负责连接独立于应用之外的系统并拉取/接收数据,RabbitMQ作为独立部署的消息中间件,完全符合外部系统的定义。

2. 如何通过IntegrationFlow.transform实现Base64与JsonObject的编解码?

可以借助JDK自带的Base64工具类处理编解码,结合Jackson的ObjectMapper实现JSON字符串与JsonObject的转换,修改后的完整代码如下:

import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
import org.springframework.context.annotation.Bean;
import org.springframework.integration.dsl.IntegrationFlow;
import org.springframework.integration.dsl.MessageChannels;
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.MessageChannel;
import org.springframework.beans.factory.annotation.Value;

import java.util.Base64;

@SpringBootApplication
public class DeepDiveIntegrationApplication {
    private final ObjectMapper objectMapper = new ObjectMapper();

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

    @Value("${mqtt.username}")
    private String mqttBrokerUsername;

    @Value("${mqtt.password}")
    private String mqttBrokerPassword;

    @Bean
    MqttPahoClientFactory clientFactory(@Value("${mqtt.brokerUrl}") String host){
        var factory = new DefaultMqttPahoClientFactory();
        var options = new MqttConnectOptions();
        options.setServerURIs(new String[]{host});
        options.setUserName(mqttBrokerUsername);
        options.setPassword(mqttBrokerPassword.toCharArray());
        factory.setConnectionOptions(options);
        return factory;
    }

    @Bean
    MessageChannel integrationMessageChannels(){
        return MessageChannels.direct().getObject();
    }

    @Bean
    MqttPahoMessageDrivenChannelAdapter inboundAdapter(@Value("${mqtt.sourceTopic}") String topic, MqttPahoClientFactory factory){
       var adapter = new MqttPahoMessageDrivenChannelAdapter("consumerClientID", factory, topic);
        adapter.setCompletionTimeout(5000);
        adapter.setConverter(new DefaultPahoMessageConverter());
        adapter.setQos(1);
        adapter.setOutputChannel(integrationMessageChannels());
        return adapter;
    }

    @Bean
    IntegrationFlow inboundFlow(MqttPahoMessageDrivenChannelAdapter inboundAdapter){
        return IntegrationFlow
                .from(inboundAdapter)
                // 解码:Base64字符串 -> JsonNode(JsonObject)
                .transform(payload -> {
                    String base64Str = (String) payload;
                    byte[] decodedBytes = Base64.getDecoder().decode(base64Str);
                    String jsonStr = new String(decodedBytes);
                    return objectMapper.readTree(jsonStr);
                })
                .handle((payload, headers) -> {
                    System.out.println("解码后的JsonObject: " + payload);
                    return payload; // 将解码后的数据传递到出站通道
                })
                .channel(integrationMessageChannels())
                .get();
    }

    @Bean
    MqttPahoMessageHandler outboundAdapter(@Value("${mqtt.targetTopic}") String topic, MqttPahoClientFactory factory){
        var mh = new MqttPahoMessageHandler("producerClientID", factory);
        mh.setDefaultTopic(topic);
        return mh;
    }

    @Bean
    IntegrationFlow outboundFlow(MessageChannel integrationMessageChannels, MqttPahoMessageHandler outboundAdapter){
        return IntegrationFlow
                .from(integrationMessageChannels)
                // 编码:JsonNode(JsonObject)-> Base64字符串
                .transform(payload -> {
                    JsonNode jsonNode = (JsonNode) payload;
                    String jsonStr = objectMapper.writeValueAsString(jsonNode);
                    return Base64.getEncoder().encodeToString(jsonStr.getBytes());
                })
                .handle(outboundAdapter)
                .get();
    }
}

关键说明:

  • 入站流通过transform将接收到的Base64字符串解码为JSON字符串,再转为JsonNode类型的JsonObject
  • 出站流通过transform将处理后的JsonObject转回JSON字符串,再编码为Base64格式发送到目标topic
  • 修正了原代码中handle方法返回null的问题,确保数据能正常流转到出站适配器

3. 对Spring Integration适配器的理解是否正确?

你的理解基本正确,但需要补充核心细节:

  • 入站适配器(如MqttPahoMessageDrivenChannelAdapter)的作用是从外部系统接收/拉取数据,并将数据投递到Spring Integration的消息通道中
  • 出站适配器(如MqttPahoMessageHandler)的作用是从Spring Integration的消息通道获取数据,并将数据发送到外部系统
  • 完整的数据流转路径是:外部系统 → 入站适配器 → 消息通道 → 处理逻辑(编解码等) → 消息通道 → 出站适配器 → 外部系统

内容的提问来源于stack exchange,提问作者Romillion

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.29 12:23:24