4.11 Deserializers
① KafkaAvroDeserializer deserializes the Avro messages.
"It should be obvious that the serializer used to produce events to Kafka must match the deserializer used when consuming events. Serializing with
IntSerializerand then deserializing withStringDeserializerwill not end well. This means that, as a developer, you need to keep track of which serializers were used to write into each topic and make sure each topic only contains data that your deserializers can interpret."
Why Avro + Schema Registry fixes this class of problem:
"the
AvroSerializercan make sure that all the data written to a specific topic is compatible with the schema of the topic, which means it can be deserialized with the matching deserializer and schema. Any errors in compatibility — on the producer or the consumer side — will be caught easily with an appropriate error message, which means you will not need to try to debug byte arrays for serialization errors."
Custom deserializer (shown to argue against it)
public class CustomerDeserializer implements Deserializer<Customer> {
@Override
public void configure(Map configs, boolean isKey) { /* nothing */ }
@Override
public Customer deserialize(String topic, byte[] data) {
int id; int nameSize; String name;
try {
if (data == null) return null;
if (data.length < 8)
throw new SerializationException("Size of data received " +
"by deserializer is shorter than expected");
ByteBuffer buffer = ByteBuffer.wrap(data);
id = buffer.getInt();
nameSize = buffer.getInt();
byte[] nameBytes = new byte[nameSize];
buffer.get(nameBytes);
name = new String(nameBytes, "UTF-8");
return new Customer(id, name);
} catch (Exception e) {
throw new SerializationException(
"Error when deserializing byte[] to Customer " + e);
}
}
@Override
public void close() { /* nothing */ }
}"The consumer also needs the implementation of the
Customerclass, and both the class and the serializer need to match on the producing and consuming applications. In a large organization with many consumers and producers sharing access to the data, this can become challenging.""implementing a custom serializer and deserializer is not recommended. It tightly couples producers and consumers and is fragile and error prone. A better solution would be to use a standard message format, such as JSON, Thrift, Protobuf, or Avro."
Avro deserialization
Properties props = new Properties();
props.put("bootstrap.servers", "broker1:9092,broker2:9092");
props.put("group.id", "CountryCounter");
props.put("key.deserializer",
"org.apache.kafka.common.serialization.StringDeserializer");
props.put("value.deserializer",
"io.confluent.kafka.serializers.KafkaAvroDeserializer"); // ①
props.put("specific.avro.reader","true"); // ③
props.put("schema.registry.url", schemaUrl); // ②
String topic = "customerContacts";
KafkaConsumer<String, Customer> consumer = new KafkaConsumer<>(props);
consumer.subscribe(Collections.singletonList(topic));
while (true) {
ConsumerRecords<String, Customer> records = consumer.poll(timeout);
for (ConsumerRecord<String, Customer> record: records) {
System.out.println("Current customer name is: " +
record.value().getName()); // ④
}
consumer.commitSync();
}① KafkaAvroDeserializer deserializes the Avro messages.
② schema.registry.url — "points to where we store the schemas. This way, the consumer can use the schema that was registered by the producer to deserialize the message."
③ specific.avro.reader=true + generated class Customer as the value type.
④ record.value() is a Customer instance.