I'm building a Spring Boot 2.4.4 application to produce message to Kafka (Confluent Platform), using schema validation.
The schema is set up on the schema registry for the topic, and I've created the POJO (LoanInitiate) from the associated AVSC file using the Maven Avro plugin. When sending using the following, there are no issues:
import org.springframework.messaging.Message;
import org.springframework.messaging.support.MessageBuilder;
...
@Autowired
private KafkaTemplate<String, LoanInitiate> kafkaTemplate;
@Value("${kafka.topic}")
private String topic;
public void regularSend(){
LoanInitiate li = initializeLoanInitiate();
kafkaTemplate.send(topic, li); //ok
}
However, if I attempt the following, I will get Unsupported Avro type.
@Autowired
private KafkaTemplate<String, Message<LoanInitiate>> kafkaTemplate2;
public void sendMessage(){
Message<LoanInitiate> message = MessageBuilder
.withPayload(li)
.setHeader(KafkaHeaders.CORRELATION_ID, "Correlation ID")
.build();
kafkaTemplate2.send(topic, message); //"Unsupported Avro type"
}
My producer configuration is as follows:
#Producer values
spring.kafka.producer.bootstrap-servers=localhost:9092
spring.kafka.producer.key-serializer=org.apache.kafka.common.serialization.StringSerializer
spring.kafka.producer.value-serializer=io.confluent.kafka.serializers.KafkaAvroSerializer
I understand that the kafkaTemplate2 is trying to send a Message<LoanInitiate> and not a LoanInitiate but it's not clear to me what I am missing to provide for writing the message and header.
How should I approach sending my message along with header information?