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)
问题原因
JsonMessageConverter的定位限制:该类仅负责消费方向的消息转换(将Kafka的ConsumerRecord转为SpringMessage),不支持生产者侧的反向转换(SpringMessage转为ProducerRecord),其convertPayload方法在生产者场景下直接抛出异常。- 转换器与序列化器不匹配:生产者配置了
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
相关产品推荐
相关产品推荐

