serialize Kafka messages with confluent registry under Flink 1.9.1

Viewed 507

is it possible to publish message to Kafka serialized with KafkaAvroSerializer by Confluent. I'm using Flink 1.9.1 have saw that some development is going on newer version of flink-avro (1.11.0) but I’m stick to the version.

I would like to use the newly introduced KafkaSerializationSchema for serializing the message to Confluent schema-registry and Kakfa.

Here I have currently a class that is converting a class type T to avro but I want to use the confluent serialization.

public class KafkaMessageSerialization<T extends SpecificRecordBase> implements KafkaSerializationSchema<T> {
    public static final Logger LOG = LoggerFactory.getLogger(KafkaMessageSerialization.class);

    final private String topic;

    public KafkaMessageSerialization(String topic) {
        this.topic = topic;
    }

    @Override
    public ProducerRecord<byte[], byte[]> serialize(T event, Long timestamp) {
        final ByteArrayOutputStream outputStream = new ByteArrayOutputStream();
        final Schema schema = event.getSchema();
        final DatumWriter<T> writer = new ReflectDatumWriter<>(schema);
        final BinaryEncoder binEncoder = EncoderFactory.get().binaryEncoder(outputStream, null);

        try {
            writer.write(event, binEncoder);
            binEncoder.flush();
        } catch (final Exception e) {
            LOG.error("serialization error", e);
            throw new RuntimeException(e);
        }

        return new ProducerRecord<>(topic, outputStream.toByteArray());
    }
}

The usage is quite convenient .addSink(new FlinkKafkaProducer<>(SINK_TOPIC, new KafkaMessageSerialization<>(SINK_TOPIC), producerProps, Semantic.AT_LEAST_ONCE))

1 Answers

I am in the same situation and based on your solution i wrote this class. I have tested it with Flink 1.10.1.

public class ConfluentAvroMessageSerialization<T extends SpecificRecordBase> implements KafkaSerializationSchema<T> {

    public static final org.slf4j.Logger LOG = LoggerFactory.getLogger(ConfluentAvroMessageSerialization.class);

    final private String topic;
    final private int schemaId;
    final private int magicByte;

    public ConfluentAvroMessageSerialization(String topic, String schemaRegistryUrl) throws IOException, RestClientException {
        magicByte = 0;
        this.topic = topic;

        SchemaRegistryClient schemaRegistry = new CachedSchemaRegistryClient(schemaRegistryUrl, 1000);
        SchemaMetadata schemaMetadata = schemaRegistry.getLatestSchemaMetadata(topic + "-value");
        schemaId = schemaMetadata.getId();

        LOG.info("Confluent Schema ID {} for topic {} found", schemaId, topic);
    }

    @Override
    public ProducerRecord<byte[], byte[]> serialize(T event, Long timestamp) {
        final ByteArrayOutputStream outputStream = new ByteArrayOutputStream();
        final Schema schema = event.getSchema();
        final DatumWriter<T> writer = new ReflectDatumWriter<>(schema);
        final BinaryEncoder binEncoder = EncoderFactory.get().binaryEncoder(outputStream, null);

        try {
            byte[] schemaIdBytes = ByteBuffer.allocate(4).putInt(schemaId).array();
            outputStream.write(magicByte); // Confluent Magic Byte
            outputStream.write(schemaIdBytes); // Confluent Schema ID (4 Byte Format)
            writer.write(event, binEncoder); // Avro data
            binEncoder.flush();
        } catch (final Exception e) {
            LOG.error("Schema Registry Serialization Error", e);
            throw new RuntimeException(e);
        }

        return new ProducerRecord<>(topic, outputStream.toByteArray());
    }
}

Confluent has a propertary format with a magic Byte and the schema id (4 Byte). For more information please check https://docs.confluent.io/current/schema-registry/serdes-develop/index.html#wire-format

Related