Spring Integration入站/出站通道适配器区别及MQTT解码实现咨询
问题背景与疑问
我的项目采用MQTT协议、RabbitMQ作为Broker,基于Spring Integration实现,流程如下:RabbitMQ(数据源) --> 入站流(将Base64数据转为JsonObject) --> 出站流(发送数据至另一topic)
RabbitMQ中的示例数据为eyJrZXkiOiJIZWxsbyBXb3JsZCJ9,解码后内容为{"key":"Hello World"}
我对入站和出站通道适配器的理解是:入站适配器用于从外部系统获取数据,数据的终点是出站适配器。现存在以下问题:
- 此场景中RabbitMQ是否属于外部系统?
- 如何通过
IntegrationFlow.transform实现Base64数据与JsonObject的编解码? - 我对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
相关产品推荐
相关产品推荐

