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);
        }
    }
}
    
Cboe Titanium U.S. Options Lite Feed Specification - Kafka Distribution | Cboe