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

使用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直接发送原始字符串即可:

  1. 定义JSON转换处理器:
@Bean
@Transformer(inputChannel = "preMqttChannel", outputChannel = MQTT_OUTBOUND_CHANNEL)
public GenericTransformer<Object, String> jsonTransformer(ObjectMapper objectMapper) {
    return source -> objectMapper.writeValueAsString(source);
}
  1. 简化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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.29 14:00:55