kafka producer throwing error when sending an Avro record to topic:io.confluent.kafka.schemaregistry.client.rest.exceptions.RestClientException: null

Viewed 421

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();
        }
    }

}

0 Answers
Related