I am trying to create a topic and produce an avro record to a kafka topic, but I am facing an issue while send the record to the producer. I am getting an error.
org.apache.kafka.common.errors.SerializationException: Error registering Avro schema: {"type":"record","name":"MyClass","namespace":"com.test.avro","fields":[{"name":"event_type","type":"string"},{"name":"time_stamp","type":"long"},{"name":"context","type":{"type":"record","name":"context","fields":[{"name":"device","type":"string"},{"name":"type","type":"string"}]}}]}
Caused by: io.confluent.kafka.schemaregistry.client.rest.exceptions.RestClientException: null; error code: 0
at io.confluent.kafka.schemaregistry.client.rest.RestService.sendHttpRequest(RestService.java:294)
Caused by: io.confluent.kafka.schemaregistry.client.rest.exceptions.RestClientException: null; error code: 0
I am not able to fix the issue.
I tried changing the version of kafka avro seriaizer from 5.3.0 version to 7.1.1
implementation group: 'io.confluent', name: 'kafka-avro-serializer', version: '5.3.0'
Facing the same issue in all the version.
My code was able to create the topic was throwing an error when I was sending/submitting a record to the topic
producer.send(record);
Sharing my code as well.
import org.apache.avro.Schema;
import org.apache.avro.generic.GenericData;
import org.apache.avro.generic.GenericRecord;
import org.apache.kafka.clients.producer.KafkaProducer;
import org.apache.kafka.clients.producer.ProducerConfig;
import org.apache.kafka.clients.producer.ProducerRecord;
import java.util.Properties;
public class AvroSerialization {
public static void main(String[] args) {
Properties props = new Properties();
props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG,
org.apache.kafka.common.serialization.StringSerializer.class);
props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG,
io.confluent.kafka.serializers.KafkaAvroSerializer.class);
props.put("schema.registry.url", "http://localhost:8081");
KafkaProducer<Object, Object> producer = new KafkaProducer<>(props);
String key = "key1";
String userSchema = "{\n" +
" \"type\" : \"record\",\n" +
" \"name\" : \"MyClass\",\n" +
" \"namespace\" : \"com.test.avro\",\n" +
" \"fields\" : [ {\n" +
" \"name\" : \"event_type\",\n" +
" \"type\" : \"string\"\n" +
" }, {\n" +
" \"name\" : \"time_stamp\",\n" +
" \"type\" : \"long\"\n" +
" }, {\n" +
" \"name\" : \"context\",\n" +
" \"type\" : {\n" +
" \"type\" : \"record\",\n" +
" \"name\" : \"context\",\n" +
" \"fields\" : [ {\n" +
" \"name\" : \"device\",\n" +
" \"type\" : \"string\"\n" +
" }, {\n" +
" \"name\" : \"type\",\n" +
" \"type\" : \"string\"\n" +
" }]\n" +
" }\n" +
" } ]\n" +
"}";
Schema.Parser parser = new Schema.Parser();
Schema schema = parser.parse(userSchema);
GenericRecord avroRecord = new GenericData.Record(schema);
GenericRecord context = new GenericData.Record(schema.getField("context").schema());
avroRecord.put("event_type", "visit");
avroRecord.put("time_stamp", System.currentTimeMillis());
context.put("device", "ios");
context.put("type","iphone12max");
avroRecord.put("context", context);
ProducerRecord<Object, Object> record = new ProducerRecord<>("avrotopic", avroRecord);
try {
producer.send(record);
} catch (Throwable e) {
e.printStackTrace();
// may need to do something with it
}
// When you're finished producing records, you can flush the producer to ensure it has all been written to Kafka and
// then close the producer to free its resources.
finally {
producer.flush();
producer.close();
}
}
}