Kafka Distribution
The following sample code demonstrates how a Kafka client consumes messages from a topic and decodes multiple AVRO-encoded Cboe Options Lite messages, such as Trade, Quote, and Refresh, using a binary decoder and message-type dispatch logic.
template <class T>
T from_payload(avro::Decoder& decoder, const std::vector<std::uint8_t>& payload) {
// returns a decoded AVRO object of type T from the given payload
auto stream = avro::memoryInputStream(
reinterpret_cast<const std::uint8_t*>(payload.data()), payload.size());
decoder.init(*stream);
T result;
avro::decode(decoder, result);
return result;
}
void kafka_consumer_loop(rd_kafka_t* rk, std::atomic<bool>& running) {
auto header_decoder = avro::binaryDecoder();
auto payload_decoder = avro::binaryDecoder();
CboeGlobalCloud::msg_payload payload;
while(running) {
// grab a message from Kafka
auto m = rd_kafka_consumer_poll(rk, 100);
if(m && !m->err) {
// frame header has a frame sequence number
// (per topic/partition) and version number
CboeGlobalCloud::frame_header frame_header;
auto frame_buffer = avro::memoryInputStream(
reinterpret_cast<const std::uint8_t*>(m->payload), m->len);
header_decoder->init(*frame_buffer);
avro::decode(*header_decoder, frame_header);
header_decoder->drain();
// use frame_header.sequence to detect gaps per topic/partition.
// sequence numbers are monotonically increasing integers starting at 1.
// 1 indicates a reset of the stream
// A Kafka message can contain more than one AVRO message...
while(frame_buffer->byteCount() < m->len) {
avro::decode(*header_decoder, payload);
switch(payload.message_type) {
case vod::client::VodDictionary::MID_NBBO: {
auto qte = from_payload<CboeGlobalCloud::truesize>(
*payload_decoder, payload.payload);
// process quote message (qte)
break;
}
case vod::client::VodDictionary::MID_TRADE: {
auto trd = from_payload<CboeGlobalCloud::trade>(
*payload_decoder, payload.payload);
// process trade message (trd)
break;
}
case vod::client::VodDictionary::MID_SECURITY_DEFINITION: {
auto r = from_payload<CboeGlobalCloud::refresh>(
*payload_decoder, payload.payload);
// process refresh message (r)
break;
}
case vod::client::VodDictionary::MID_OPTION_CHAIN: {
auto r = from_payload<CboeGlobalCloud::reference>(
*payload_decoder, payload.payload);
// process reference message (r)
break;
}
case vod::client::VodDictionary::MID_DELETE: {
auto del = from_payload<std::string>(
*payload_decoder, payload.payload);
// process delete message (del)
break;
}
default:
// LOG, unknown message type
break;
}
// return any non-decoded bytes to the AVRO decoder
header_decoder->drain();
}
rd_kafka_message_destroy(m);
}
}
}





