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

Spring Kafka发送消息时JsonMessageConverter抛出UnsupportedOperationException

问题:Spring Kafka发送消息时JsonMessageConverter抛出UnsupportedOperationException

配置代码

Java配置类与监听器

@Configuration
public class Config {
    @Bean
    public RecordMessageConverter converter() {
        return new JsonMessageConverter();
    }
    
    @Bean
    public BatchMessagingMessageConverter batchConverter() {
        return new BatchMessagingMessageConverter(converter());
    }
}

@Component
@RequiredArgsConstructor
public class CancelAuthorizationLinkageListener {
    private final KafkaTemplate<String, CancelAuthorizationLinkage> kafkaTemplate;
    
    @KafkaListener(
            id = "${spring.kafka.listener.cancel-auth-linkage.id}",
            topics = "${spring.kafka.listener.cancel-auth-linkage.topic.linkage}",
            autoStartup = "false",
            batch = "true",
            groupId = "cushion")
    public void listen(List<Message<CancelAuthorizationLinkage>> messages) {
        // other operations...
        messages.forEach(message -> kafkaTemplate.send(build));
        kafkaTemplate.send(build);
    }
}

YAML配置

spring:
  ...
  kafka:
    bootstrap-servers: localhost:9092
    producer:
      acks: -1
      transaction-id-prefix: cushion-kafka-tx
      value-serializer: org.springframework.kafka.support.serializer.JsonSerializer
    consumer:
      enable-auto-commit: false
      auto-offset-reset: earliest
      max-poll-records: 100
      value-deserializer: org.apache.kafka.common.serialization.ByteArrayDeserializer
      properties:
        isolation.level: read_committed
        spring.json.trusted.packages: "*"
  ...

异常信息

Caused by: java.lang.UnsupportedOperationException: Select a subclass that creates a ProducerRecord value corresponding to the configured Kafka Serializer
    at org.springframework.kafka.support.converter.JsonMessageConverter.convertPayload(JsonMessageConverter.java:96)
    at org.springframework.kafka.support.converter.MessagingMessageConverter.fromMessage(MessagingMessageConverter.java:251)
    at org.springframework.kafka.core.KafkaTemplate.send(KafkaTemplate.java:583)

问题原因

  1. JsonMessageConverter的定位限制:该类仅负责消费方向的消息转换(将Kafka的ConsumerRecord转为Spring Message),不支持生产者侧的反向转换(Spring Message转为ProducerRecord),其convertPayload方法在生产者场景下直接抛出异常。
  2. 转换器与序列化器不匹配:生产者配置了JsonSerializer,但KafkaTemplate发送Message对象时,会调用自定义的JsonMessageConverter处理Payload,而该转换器无法生成适配JsonSerializer的格式。

解决方案

方案1:替换为生产者友好的消息转换器

将JsonMessageConverter替换为它的子类,根据序列化场景选择:

  • 匹配当前消费者ByteArrayDeserializer配置,选择BytesJsonMessageConverter,它会将Payload转为字节数组,适配JsonSerializer:
@Configuration
public class Config {
    @Bean
    public RecordMessageConverter converter() {
        return new BytesJsonMessageConverter();
    }
    
    @Bean
    public BatchMessagingMessageConverter batchConverter() {
        return new BatchMessagingMessageConverter(converter());
    }
}
  • 若使用字符串序列化,可选择StringJsonMessageConverter。

方案2:直接发送实体对象(跳过Message转换)

如果不需要传递Spring Message的附加属性,直接发送CancelAuthorizationLinkage实体,KafkaTemplate会直接用配置的JsonSerializer序列化,无需经过消息转换器:

// 修改发送逻辑,直接传递实体
messages.forEach(msg -> kafkaTemplate.send("目标Topic名称", msg.getPayload()));

方案3:确保KafkaTemplate绑定正确的转换器

如果是手动创建KafkaTemplate,需手动设置消息转换器;自动配置的KafkaTemplate会自动注入自定义的RecordMessageConverter,也可手动确认绑定:

@Bean
public KafkaTemplate<String, CancelAuthorizationLinkage> kafkaTemplate(ProducerFactory<String, CancelAuthorizationLinkage> producerFactory, RecordMessageConverter converter) {
    KafkaTemplate<String, CancelAuthorizationLinkage> template = new KafkaTemplate<>(producerFactory);
    template.setMessageConverter(converter);
    return template;
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.20 07:50:03