C++ Message Queues: RabbitMQ vs Kafka, Serialization and Error Handling
Introduction: “Sync calls are the bottleneck”
If service A calls B over HTTP and waits, A pays B’s full latency (DB, external APIs, heavy CPU). Message queues let A publish work and return quickly—decoupling producers and consumers and buffering spikes.
The price is that you trade a simple failure model for a more complicated one. With a synchronous call, either B answered or A got an error. With a queue, the message can be accepted by the broker and then processed once, twice, late, or never, and A no longer knows which. Most of this article is about the second half of that sentence: acknowledgments, offset commits, retries and dead letter queues are the tools that turn “maybe delivered” into something you can reason about.
Topics:
- RabbitMQ (AMQP) — exchanges, queues, acks (e.g. SimpleAmqpClient patterns)
- Kafka — topics, partitions, consumer groups (e.g. librdkafka)
- Serialization — JSON, Protobuf
- Errors, tuning, production patterns Environment: C++17+.
When queues help
Use cases
- Order pipeline — user gets immediate ACK; payment, stock, email, analytics run asynchronously.
- Traffic spikes — queue absorbs bursts; workers drain at sustainable rate.
- Microservices — publish OrderCreated; downstream services subscribe independently.
- Centralized logs/metrics — Kafka-style log aggregation.
- Heavy jobs — image resize, encoding, fan-out to worker pools.
Architecture comparison
flowchart LR
subgraph sync["Synchronous (HTTP)"]
A1[Service A] -->|wait| B1[Service B]
B1 -->|wait| C1[Service C]
end
subgraph async["Asynchronous (Queue)"]
A2[Service A] -->|publish| Q[Queue]
Q --> B2[Service B]
Q --> C2[Service C]
end
Benefits:
- Decoupling: Services don’t need to know about each other
- Buffering: Handle traffic spikes
- Reliability: Retry failed messages
- Scalability: Add more consumers
A queue is not free, though. It adds a broker you must run, monitor and upgrade, it makes end-to-end latency harder to see, and it turns “the request failed” into “the request is somewhere in a pipeline”. If the caller genuinely needs the answer before it can continue (checking a password, reading a price for the page it is rendering), a queue only adds a hop. Queues pay off when the caller can accept “accepted, will be done” as the response.
RabbitMQ vs Kafka
| RabbitMQ | Kafka | |
|---|---|---|
| Model | Classic broker, routing | Durable log / streaming |
| Retention | Often delete after ack | Time/size retention |
| Throughput | High; per-message acks and routing cost more per message | Higher for bulk streams; batching and sequential disk writes |
| Ordering | Per queue | Per partition |
| Consumer state | Broker tracks each message’s ack | Consumer tracks an offset per partition |
| Use case | RPC-ish, work queues | Event logs, analytics |
| Replay | Limited (Streams add it) | Full replay within retention |
| Routing | Flexible (exchanges) | Simple (topics) |
The row that matters most in practice is “consumer state”. RabbitMQ remembers, per message, whether it was acknowledged; an unacked message goes back to the queue when the consumer’s channel closes. Kafka does not track individual messages at all: a consumer group stores one committed offset per partition, and “acknowledging” message 105 implicitly acknowledges 0 through 104. That single difference explains most of the behavioral differences below, including why Kafka cannot skip one bad message without a DLQ and why RabbitMQ cannot replay yesterday’s traffic.
When to use RabbitMQ
- Traditional messaging patterns
- Complex routing (topic, fanout, direct)
- RPC-style request/response
- Per-message acknowledgment and redelivery matter
- Messages are work to be done, not history to keep
When to use Kafka
- Event sourcing
- Log aggregation
- Sustained high-volume streams
- Message replay needed
- Stream processing
- Multiple consumers reading same data
Throughput numbers you see quoted for either broker depend heavily on message size, persistence settings, replication factor, acks level and batching, so measure with your own payloads before letting throughput decide. For most C++ services the deciding question is simpler: do several independent systems need to read the same events, possibly later? If yes, Kafka. If each message is a job that one worker should do once, RabbitMQ is usually the more natural fit.
RabbitMQ implementation
Installation
# Ubuntu
sudo apt-get install rabbitmq-server
# Docker
docker run -d --name rabbitmq \
-p 5672:5672 \
-p 15672:15672 \
rabbitmq:3-management
# C++ client
sudo apt-get install librabbitmq-dev
# Or use AMQP-CPP: https://github.com/CopernicaMarketingSoftware/AMQP-CPP
The examples below use rabbitmq-c (librabbitmq), the C client, wrapped in small C++ classes. It is synchronous and blocking, which keeps the code readable. AMQP-CPP is the usual alternative when you already run an event loop (libev, libuv, Boost.Asio), because it only parses the protocol and lets you supply the I/O.
Producer example
#include <amqp.h>
#include <amqp_tcp_socket.h>
#include <iostream>
#include <string>
class RabbitMQProducer {
amqp_connection_state_t conn_;
amqp_socket_t* socket_;
std::string exchange_;
public:
RabbitMQProducer(const std::string& host, int port, const std::string& exchange)
: exchange_(exchange) {
conn_ = amqp_new_connection();
socket_ = amqp_tcp_socket_new(conn_);
if (amqp_socket_open(socket_, host.c_str(), port) != 0) {
throw std::runtime_error("Cannot open socket");
}
amqp_rpc_reply_t reply = amqp_login(conn_, "/", 0, 131072, 0,
AMQP_SASL_METHOD_PLAIN, "guest", "guest");
if (reply.reply_type != AMQP_RESPONSE_NORMAL) {
throw std::runtime_error("Login failed");
}
amqp_channel_open(conn_, 1);
amqp_get_rpc_reply(conn_);
}
~RabbitMQProducer() {
amqp_channel_close(conn_, 1, AMQP_REPLY_SUCCESS);
amqp_connection_close(conn_, AMQP_REPLY_SUCCESS);
amqp_destroy_connection(conn_);
}
void publish(const std::string& routing_key, const std::string& message) {
amqp_basic_properties_t props;
props._flags = AMQP_BASIC_CONTENT_TYPE_FLAG | AMQP_BASIC_DELIVERY_MODE_FLAG;
props.content_type = amqp_cstring_bytes("application/json");
props.delivery_mode = 2; // Persistent
amqp_bytes_t message_bytes = amqp_cstring_bytes(message.c_str());
int result = amqp_basic_publish(
conn_,
1,
amqp_cstring_bytes(exchange_.c_str()),
amqp_cstring_bytes(routing_key.c_str()),
0,
0,
&props,
message_bytes
);
if (result != 0) {
throw std::runtime_error("Publish failed");
}
}
};
A few details in this producer are easy to get wrong:
delivery_mode = 2marks the message persistent, but that only helps if the queue itself was declared durable. A persistent message in a non-durable queue is still gone after a broker restart.amqp_basic_publishreturning 0 does not mean the broker stored the message. It only means the frame was written to the socket. If the exchange does not exist, the broker closes the channel asynchronously and you discover it on the next call. For real delivery guarantees, enable publisher confirms (amqp_confirm_select) and wait for the broker’sbasic.ack.- Unroutable messages are dropped silently. Publishing to an exchange with a routing key that matches no binding succeeds and the message disappears. Setting the
mandatoryflag (the first0after the routing key) makes the broker return such messages instead. - The destructor ignores errors, which is fine for teardown, but the constructor’s
amqp_channel_openresult is never checked either; in production, checkamqp_get_rpc_replyafter every RPC-style call, because that is where “access refused” and “vhost not found” show up.
The first time I wired up a producer like this, everything “worked” until someone renamed an exchange in the deployment config. Publishes kept returning success, nothing arrived, and there was no error anywhere. That is the standard way to learn about publisher confirms and the mandatory flag.
Consumer example
class RabbitMQConsumer {
amqp_connection_state_t conn_;
std::string queue_;
public:
RabbitMQConsumer(const std::string& host, int port, const std::string& queue)
: queue_(queue) {
conn_ = amqp_new_connection();
amqp_socket_t* socket = amqp_tcp_socket_new(conn_);
amqp_socket_open(socket, host.c_str(), port);
amqp_login(conn_, "/", 0, 131072, 0, AMQP_SASL_METHOD_PLAIN, "guest", "guest");
amqp_channel_open(conn_, 1);
amqp_get_rpc_reply(conn_);
// Set prefetch count
amqp_basic_qos(conn_, 1, 0, 10, 0);
}
~RabbitMQConsumer() {
amqp_channel_close(conn_, 1, AMQP_REPLY_SUCCESS);
amqp_connection_close(conn_, AMQP_REPLY_SUCCESS);
amqp_destroy_connection(conn_);
}
void consume(std::function<bool(const std::string&)> handler) {
amqp_basic_consume(conn_, 1, amqp_cstring_bytes(queue_.c_str()),
amqp_empty_bytes, 0, 0, 0, amqp_empty_table);
while (true) {
amqp_envelope_t envelope;
amqp_maybe_release_buffers(conn_);
amqp_rpc_reply_t result = amqp_consume_message(conn_, &envelope, NULL, 0);
if (result.reply_type == AMQP_RESPONSE_NORMAL) {
std::string message(
static_cast<char*>(envelope.message.body.bytes),
envelope.message.body.len
);
bool success = handler(message);
if (success) {
amqp_basic_ack(conn_, 1, envelope.delivery_tag, 0);
} else {
amqp_basic_nack(conn_, 1, envelope.delivery_tag, 0, 1);
}
amqp_destroy_envelope(&envelope);
}
}
}
};
The important lines are the QoS call and the ack/nack pair. amqp_basic_qos(conn_, 1, 0, 10, 0) sets a prefetch count of 10: the broker will push at most ten unacknowledged messages to this consumer. Without it, RabbitMQ pushes as many messages as the network allows, and a slow consumer ends up holding thousands of messages in memory while other consumers sit idle. Prefetch is the main backpressure knob on the consumer side; a small value spreads work evenly across workers, a larger one reduces round trips when handlers are fast.
The fourth 0 in amqp_basic_consume is no_ack = false, meaning manual acknowledgment. The message stays “unacked” on the broker until amqp_basic_ack runs, and if the process dies first, the broker redelivers it to another consumer. That is at-least-once delivery, so the handler must tolerate duplicates.
The amqp_basic_nack(..., requeue = 1) branch is the one that bites. A message that always fails (malformed JSON, a bug in the handler) is requeued, delivered again immediately, fails again, and the consumer spins at full CPU on one “poison” message. Either nack with requeue = 0 and configure a dead-letter exchange on the queue (x-dead-letter-exchange), or count attempts yourself. The loop also silently ignores non-normal replies from amqp_consume_message; a closed connection returns AMQP_RESPONSE_LIBRARY_EXCEPTION and this loop would spin forever instead of reconnecting.
Kafka implementation
Installation
# Install librdkafka
sudo apt-get install librdkafka-dev
# Or build from source
git clone https://github.com/edenhill/librdkafka.git
cd librdkafka
./configure
make
sudo make install
Producer example
#include <librdkafka/rdkafkacpp.h>
#include <iostream>
#include <string>
class KafkaProducer {
std::unique_ptr<RdKafka::Producer> producer_;
std::unique_ptr<RdKafka::Topic> topic_;
public:
KafkaProducer(const std::string& brokers, const std::string& topic) {
std::string errstr;
RdKafka::Conf* conf = RdKafka::Conf::create(RdKafka::Conf::CONF_GLOBAL);
conf->set("bootstrap.servers", brokers, errstr);
conf->set("client.id", "cpp-producer", errstr);
producer_.reset(RdKafka::Producer::create(conf, errstr));
if (!producer_) {
throw std::runtime_error("Failed to create producer: " + errstr);
}
RdKafka::Conf* tconf = RdKafka::Conf::create(RdKafka::Conf::CONF_TOPIC);
topic_.reset(RdKafka::Topic::create(producer_.get(), topic, tconf, errstr));
delete conf;
delete tconf;
}
void produce(const std::string& key, const std::string& message) {
RdKafka::ErrorCode resp = producer_->produce(
topic_.get(),
RdKafka::Topic::PARTITION_UA,
RdKafka::Producer::RK_MSG_COPY,
const_cast<char*>(message.c_str()),
message.size(),
&key,
nullptr
);
if (resp != RdKafka::ERR_NO_ERROR) {
throw std::runtime_error("Produce failed: " + RdKafka::err2str(resp));
}
producer_->poll(0);
}
void flush(int timeout_ms = 10000) {
producer_->flush(timeout_ms);
}
};
produce() in librdkafka is asynchronous: it copies the payload into an internal queue (RK_MSG_COPY) and returns immediately. A background thread batches messages per partition and sends them. Two consequences follow. First, ERR_NO_ERROR from produce() only means “queued locally”; a delivery failure is reported later through a delivery report callback (dr_cb), which you only receive if you register one and keep calling poll(). Second, messages still in the local queue are lost if the process exits without flush(). ERR__QUEUE_FULL from produce() is librdkafka telling you the local queue (queue.buffering.max.messages) is full; the correct reaction is to poll() for a while and retry, not to drop the message.
PARTITION_UA (“unassigned”) lets the partitioner choose a partition from the key, so all messages with the same key land in the same partition and stay ordered. Passing an empty key spreads messages across partitions and gives up ordering. For durability, set acks=all (the default since librdkafka 1.x is all) and consider enable.idempotence=true, which prevents duplicates caused by producer retries.
Consumer example
class KafkaConsumer {
std::unique_ptr<RdKafka::KafkaConsumer> consumer_;
public:
KafkaConsumer(const std::string& brokers, const std::string& group_id,
const std::vector<std::string>& topics) {
std::string errstr;
RdKafka::Conf* conf = RdKafka::Conf::create(RdKafka::Conf::CONF_GLOBAL);
conf->set("bootstrap.servers", brokers, errstr);
conf->set("group.id", group_id, errstr);
conf->set("auto.offset.reset", "earliest", errstr);
conf->set("enable.auto.commit", "false", errstr);
consumer_.reset(RdKafka::KafkaConsumer::create(conf, errstr));
if (!consumer_) {
throw std::runtime_error("Failed to create consumer: " + errstr);
}
RdKafka::ErrorCode err = consumer_->subscribe(topics);
if (err) {
throw std::runtime_error("Subscribe failed: " + RdKafka::err2str(err));
}
delete conf;
}
void consume(std::function<bool(const std::string&, const std::string&)> handler) {
while (true) {
RdKafka::Message* message = consumer_->consume(1000);
if (message->err() == RdKafka::ERR_NO_ERROR) {
std::string key = message->key() ? *message->key() : "";
std::string payload(
static_cast<const char*>(message->payload()),
message->len()
);
bool success = handler(key, payload);
if (success) {
consumer_->commitSync(message);
}
}
delete message;
}
}
};
With enable.auto.commit=false, the consumer commits its position only after the handler succeeds, which again gives at-least-once delivery. Note what happens when the handler returns false: nothing is committed, but the consumer still moves on to the next message, because Kafka consumers track their read position in memory independently of the committed offset. The failed message is only re-read after a restart or rebalance, and when the next successful message is committed, its offset covers the failed one too. In Kafka you cannot “leave one message unacked”; failures have to be handled explicitly with retries or a DLQ topic.
commitSync per message is also slow, since each call is a round trip to the group coordinator. Committing every N messages or every few hundred milliseconds with commitAsync is the usual compromise; the cost is more duplicates after a crash. Finally, consume(1000) returns a message object even on timeout, with err() == ERR__TIMED_OUT, which is why the object is deleted unconditionally. Errors such as ERR__PARTITION_EOF also arrive this way and should be logged, not treated as payloads.
Serialization strategies
JSON serialization
#include <nlohmann/json.hpp>
struct Order {
std::string id;
std::string customer;
double amount;
std::string to_json() const {
nlohmann::json j;
j["id"] = id;
j["customer"] = customer;
j["amount"] = amount;
return j.dump();
}
static Order from_json(const std::string& json_str) {
auto j = nlohmann::json::parse(json_str);
Order order;
order.id = j["id"];
order.customer = j["customer"];
order.amount = j["amount"];
return order;
}
};
JSON is readable in the broker’s management UI and every language can parse it, which makes debugging a pipeline much easier. Its weaknesses are size, parsing cost, and the lack of an enforced schema: j["amount"] throws nlohmann::json::type_error if a producer starts sending the amount as a string, and j["id"] on a missing key inserts a null value, whose conversion to std::string then throws the same kind of error. Use j.at("id") and catch the exception so a malformed message goes to the DLQ instead of crashing the consumer.
Protobuf serialization
// order.proto
syntax = "proto3";
message OrderProto {
string id = 1;
string customer = 2;
double amount = 3;
int64 timestamp = 4;
}
#include "order.pb.h"
std::string serialize_order(const Order& order) {
OrderProto proto;
proto.set_id(order.id);
proto.set_customer(order.customer);
proto.set_amount(order.amount);
proto.set_timestamp(std::time(nullptr));
std::string serialized;
proto.SerializeToString(&serialized);
return serialized;
}
Order deserialize_order(const std::string& data) {
OrderProto proto;
proto.ParseFromString(data);
Order order;
order.id = proto.id();
order.customer = proto.customer();
order.amount = proto.amount();
return order;
}
The proto message is named OrderProto so it does not collide with the C++ Order struct in the same namespace. Protobuf gives you compact binary messages and, more importantly, rules for evolving the schema: new fields get new numbers, old consumers ignore fields they do not know, and new consumers see default values for fields old producers never set. Never reuse a field number after removing a field; mark it reserved instead. Also check the return value of ParseFromString: it returns false for bytes that are not a valid message, and ignoring it gives you an empty order with default values rather than an error. The trade-off versus JSON is that binary payloads are unreadable in broker tools and every producer and consumer must share the generated code.
Real-world examples
Example 1: Order processing system
class OrderProcessor {
KafkaProducer producer_;
public:
OrderProcessor() : producer_("localhost:9092", "orders") {}
void create_order(const Order& order) {
// Publish to Kafka
std::string message = order.to_json();
producer_.produce(order.id, message);
producer_.flush();
std::cout << "Order published: " << order.id << "\n";
}
};
class PaymentService {
KafkaConsumer consumer_;
public:
PaymentService()
: consumer_("localhost:9092", "payment-service", {"orders"}) {}
void start() {
consumer_.consume([this](const std::string& key, const std::string& payload) {
try {
Order order = Order::from_json(payload);
process_payment(order);
return true;
} catch (const std::exception& e) {
std::cerr << "Payment failed: " << e.what() << "\n";
return false;
}
});
}
private:
void process_payment(const Order& order) {
std::cout << "Processing payment for order: " << order.id << "\n";
// Payment logic...
}
};
Two things in this example are deliberate simplifications. create_order calls flush() after every message, which makes each publish synchronous and removes most of the batching benefit; it is acceptable for low-volume order creation where you want the publish confirmed before replying, but not for high-volume streams. And using order.id as the key means all events for one order go to one partition, so a later OrderCancelled can never overtake OrderCreated.
A harder problem hides in the flow itself: the order is usually written to a database and published to Kafka, and those are two separate systems. If the process crashes between the two, you have an order nobody hears about, or an event for an order that was never stored. The common fix is the transactional outbox: write the event into an outbox table in the same database transaction as the order, and let a separate relay publish rows from that table.
Example 2: Log aggregation
class LogAggregator {
KafkaProducer producer_;
public:
LogAggregator() : producer_("localhost:9092", "logs") {}
void log(const std::string& level, const std::string& message) {
nlohmann::json log_entry;
log_entry["timestamp"] = std::time(nullptr);
log_entry["level"] = level;
log_entry["message"] = message;
log_entry["service"] = "myapp";
producer_.produce("", log_entry.dump());
}
};
Performance comparison
Rather than quote benchmark tables that depend on hardware, payload size and configuration, it is more useful to know which settings move performance and what they cost:
- Batching (Kafka
linger.ms,batch.size): waiting a few milliseconds to fill a batch raises throughput considerably and adds that much latency. librdkafka’s defaultlinger.msis 5 ms. - Compression (
compression.type=lz4orzstd): works on whole batches, so it helps most when batching is effective and payloads are repetitive, as JSON usually is. - Durability (
acks=all, RabbitMQ publisher confirms, quorum queues): waiting for replicas or confirms costs latency per message; batching confirms rather than waiting for each one recovers most of the throughput. - Consumer prefetch and commit frequency: small prefetch and per-message commits are safe and slow; larger values are fast and increase redelivery after a crash.
Measure with kafka-producer-perf-test.sh or RabbitMQ’s PerfTest using your real message sizes. In most C++ services I have seen, the broker was not the bottleneck; the handler’s database calls were.
Error handling
Retry logic
class RetryHandler {
int max_retries_;
std::chrono::milliseconds backoff_;
public:
RetryHandler(int max_retries, std::chrono::milliseconds backoff)
: max_retries_(max_retries), backoff_(backoff) {}
bool handle_message(const std::string& message) {
for (int attempt = 0; attempt < max_retries_; ++attempt) {
try {
process(message);
return true;
} catch (const std::exception& e) {
std::cerr << "Attempt " << attempt + 1 << " failed: " << e.what() << "\n";
if (attempt < max_retries_ - 1) {
std::this_thread::sleep_for(backoff_ * (attempt + 1));
}
}
}
// Send to dead letter queue
send_to_dlq(message);
return false;
}
};
This handler retries in-process with linear backoff (backoff * (attempt + 1)). That is fine for short, transient failures such as a brief database hiccup, but it blocks the consumer thread for the whole backoff. On Kafka, blocking longer than max.poll.interval.ms (5 minutes by default) makes the broker consider the consumer dead and triggers a rebalance, after which the same message is handed to another consumer and the retry storm moves with it. For long waits, publish the message to a separate retry topic or queue with a delay and move on. Also distinguish failures: retrying a message that fails validation will never help, so send it to the DLQ immediately.
Dead letter queue
class DeadLetterQueue {
KafkaProducer dlq_producer_;
public:
DeadLetterQueue() : dlq_producer_("localhost:9092", "dlq") {}
void send(const std::string& original_topic, const std::string& message,
const std::string& error) {
nlohmann::json dlq_message;
dlq_message["original_topic"] = original_topic;
dlq_message["message"] = message;
dlq_message["error"] = error;
dlq_message["timestamp"] = std::time(nullptr);
dlq_producer_.produce("", dlq_message.dump());
}
};
A DLQ is only useful if someone looks at it. Store enough context to act on (original topic, partition, offset, error text, attempt count), alert when it grows, and build a way to replay messages back to the original topic once the bug is fixed. An unmonitored DLQ just moves data loss to a place nobody checks.
Common mistakes
Mistake 1: Not handling backpressure
// ❌ BAD: Producer overwhelms consumer
while (true) {
producer.produce(message); // No rate limiting
}
// ✅ GOOD: Rate limiting
RateLimiter limiter(1000); // 1000 msg/s
while (true) {
limiter.wait();
producer.produce(message);
}
Rate limiting the producer is one option, but in a queue system the broker is supposed to absorb bursts. The real backpressure problem is usually unbounded memory: librdkafka’s local queue or an in-process buffer grows until the process is killed. Bound the local queue, treat ERR__QUEUE_FULL as a signal to slow down, and watch consumer lag so you add consumers before the backlog becomes hours long.
Mistake 2: Not committing offsets
// ❌ BAD: Auto-commit may lose messages
consumer.consume([](const std::string& msg) {
process(msg); // If crash here, message lost
});
// ✅ GOOD: Manual commit after processing
consumer.consume([&](const std::string& msg) {
process(msg);
consumer.commit(); // Commit after success
});
Auto-commit in Kafka commits the offsets returned by previous polls on a timer (auto.commit.interval.ms, 5 seconds by default). If the consumer crashes after the commit but before processing finished, those messages are skipped. Manual commit after processing flips the failure mode to duplicates, which idempotent handlers can absorb.
Mistake 3: Assuming exactly-once
Neither broker gives end-to-end exactly-once processing for side effects outside the broker. Kafka transactions make “read, process, write back to Kafka” atomic, but a handler that charges a card or sends an email can still run twice. Give each message a unique ID and record processed IDs in the same transaction as the side effect, or make the side effect naturally idempotent (an upsert keyed by order ID).
Best practices
- Idempotent consumers: Handle duplicate messages
- Schema evolution: Use Protobuf or Avro
- Monitoring: Track lag, throughput, errors
- Partitioning: Design partition keys carefully
- Backpressure: Implement rate limiting
- Dead letter queues: Handle poison messages
- Graceful shutdown: Flush producers, close consumers
- Testing: Use embedded brokers for tests
Production patterns
Pattern 1: Circuit breaker
class CircuitBreaker {
enum State { CLOSED, OPEN, HALF_OPEN };
State state_ = CLOSED;
int failure_count_ = 0;
int threshold_ = 5;
public:
bool allow_request() {
if (state_ == OPEN) {
// Check if should try again
return false;
}
return true;
}
void on_success() {
failure_count_ = 0;
state_ = CLOSED;
}
void on_failure() {
failure_count_++;
if (failure_count_ >= threshold_) {
state_ = OPEN;
}
}
};
This circuit breaker is intentionally minimal: once it opens, allow_request() returns false forever, because the HALF_OPEN transition is left as a comment. A real implementation records when the circuit opened, moves to HALF_OPEN after a cool-down, lets one trial request through, and closes or reopens depending on the result. It also needs a mutex or atomics if several consumer threads share it. In a consumer, an open circuit should pause consumption (pause() on the Kafka partitions, or cancel the RabbitMQ consumer) rather than fail every message into the DLQ while the downstream service is down.
RabbitMQ or Kafka for a C++ service
Choose by how consumers read the data. If each message is a task that one worker should process and then forget, and you want per-message acknowledgements, routing by key or pattern, priorities or delayed retries, RabbitMQ fits. If several independent consumers need to read the same stream of events, new consumers must be able to start from the beginning, or you need to replay history after fixing a bug, Kafka’s retained, partitioned log fits better.
Whichever you pick, plan for duplicates. Both are normally operated with at-least-once delivery: a consumer that crashes after doing the work but before acknowledging or committing its offset will see the message again. Make handlers idempotent, for example by recording processed message IDs, rather than trying to configure duplicates away.
Next: Caching strategy (#50-8) Previous: Production deployment (#50-5)
Frequently Asked Questions (FAQ)
Q. If I commit offsets only after processing, can a message be processed twice?
A. Yes. Committing after processing gives at-least-once delivery: if the consumer crashes after handling a message but before the commit, the message is delivered again on restart. That is why the article recommends idempotent consumers, for example by storing a message ID and skipping IDs already processed, or by making the side effect an upsert. Committing before processing avoids duplicates but loses the message on a crash, which is usually worse.