使用Spring Messaging Message与Protobuf作为Kafka消息时出现类型转换异常
问题原因
你调用kafkaTemplate.send(topic, key, message)时,传入的第三个参数是Spring的GenericMessage对象,但生产者配置的value-serializer是KafkaProtobufSerializer,它仅支持序列化Protobuf原生的Message类型(即你的Event对象),序列化器尝试将GenericMessage强转为Protobuf Message时触发了类型转换异常。
解决方法
根据是否需要传递自定义头,有两种处理方式:
方式1:无需自定义头时,直接发送Protobuf对象
跳过Spring Message包装,直接将Protobuf结构体作为消息value传入send方法:
kafkaTemplate.send(topic, new MessageKey(key).toString(), protobufStruct);
方式2:需要传递自定义头时,操作Kafka原生Header
如果要传递MESSAGE_TIMESTAMP这类头信息,通过Kafka原生ProducerRecord来设置Header:
import org.apache.kafka.clients.producer.ProducerRecord; import java.nio.charset.StandardCharsets; // 创建ProducerRecord,传入topic、key和Protobuf对象 ProducerRecord<String, Event> record = new ProducerRecord<>(topic, new MessageKey(key).toString(), protobufStruct); // 添加自定义Header record.headers().add(MESSAGE_TIMESTAMP, Instant.now().toString().getBytes(StandardCharsets.UTF_8)); // 发送记录 kafkaTemplate.send(record);
若坚持使用Spring Message包装,需为KafkaTemplate配置正确的消息转换器,确保Spring Message的Payload被提取后传递给序列化器,但这种方式复杂度更高,不如直接操作ProducerRecord简洁。
内容的提问来源于stack exchange,提问作者JITHIN_PATHROSE
相关产品推荐
相关产品推荐

