使用Spring Integration将对象转JSON作为MQTT负载报错求解决方案
解决方案
方法一:直接配置BytesMessageMapper处理JSON序列化
不需要使用DefaultPahoMessageConverter,直接将ConvertingBytesMessageMapper绑定到Mqttv5PahoMessageHandler即可,避开断言限制的同时完成JSON序列化:
@Bean @ServiceActivator(inputChannel = MQTT_OUTBOUND_CHANNEL) public MessageHandler mqttOutbound(final ClientManager<IMqttAsyncClient, MqttConnectionOptions> clientManager, final MqttClientPersistence persistence, final MqttHeaderMapper mqttHeaderMapper, final JacksonJsonMessageConverter jacksonJsonMessageConverter) { final var messageHandler = new Mqttv5PahoMessageHandler(clientManager); messageHandler.setHeaderMapper(mqttHeaderMapper); messageHandler.setPersistence(persistence); messageHandler.setAsync(true); messageHandler.setAsyncEvents(false); // 直接设置JSON序列化的BytesMessageMapper final var bytesMessageMapper = new ConvertingBytesMessageMapper(jacksonJsonMessageConverter); messageHandler.setBytesMessageMapper(bytesMessageMapper); // 设置默认QoS(根据业务需求调整) messageHandler.setDefaultQos(MqttQoS.AT_LEAST_ONCE.value()); return messageHandler; }
方法二:提前将对象转为JSON字符串发送
在消息进入MQTT出站通道前,通过转换器将Java对象序列化为JSON字符串,Mqttv5PahoMessageHandler直接发送原始字符串即可:
- 定义JSON转换处理器:
@Bean @Transformer(inputChannel = "preMqttChannel", outputChannel = MQTT_OUTBOUND_CHANNEL) public GenericTransformer<Object, String> jsonTransformer(ObjectMapper objectMapper) { return source -> objectMapper.writeValueAsString(source); }
- 简化MQTT出站处理器:
@Bean @ServiceActivator(inputChannel = MQTT_OUTBOUND_CHANNEL) public MessageHandler mqttOutbound(final ClientManager<IMqttAsyncClient, MqttConnectionOptions> clientManager, final MqttClientPersistence persistence, final MqttHeaderMapper mqttHeaderMapper) { final var messageHandler = new Mqttv5PahoMessageHandler(clientManager); messageHandler.setHeaderMapper(mqttHeaderMapper); messageHandler.setPersistence(persistence); messageHandler.setAsync(true); messageHandler.setAsyncEvents(false); messageHandler.setDefaultQos(MqttQoS.AT_LEAST_ONCE.value()); return messageHandler; }
方法三:改用无ClientManager的构造器(若业务允许)
如果不需要依赖ClientManager管理MQTT客户端,可以使用基于客户端ID和服务端地址的构造器,此时允许使用DefaultPahoMessageConverter:
@Bean @ServiceActivator(inputChannel = MQTT_OUTBOUND_CHANNEL) public MessageHandler mqttOutbound(final MqttConnectionOptions connectionOptions, final MqttClientPersistence persistence, final MqttHeaderMapper mqttHeaderMapper, final JacksonJsonMessageConverter jacksonJsonMessageConverter) { // 用客户端ID和MQTT服务地址构造处理器 final var messageHandler = new Mqttv5PahoMessageHandler("your-client-id", "tcp://mqtt-server:1883", connectionOptions); messageHandler.setHeaderMapper(mqttHeaderMapper); messageHandler.setPersistence(persistence); messageHandler.setAsync(true); messageHandler.setAsyncEvents(false); final var defaultPahoMessageConverter = new DefaultPahoMessageConverter(MqttQoS.AT_LEAST_ONCE.value(), false, StandardCharsets.UTF_8.name()); final var bytesMessageMapper = new ConvertingBytesMessageMapper(jacksonJsonMessageConverter); defaultPahoMessageConverter.setBytesMessageMapper(bytesMessageMapper); messageHandler.setConverter(defaultPahoMessageConverter); return messageHandler; }
内容的提问来源于stack exchange,提问作者wilx
相关产品推荐
相关产品推荐

