Queues
Queue messaging methods (Simple and Stream APIs) -- KubeMQ C++ SDK reference.
Guaranteed message delivery with acknowledgment. Two APIs are available: Simple (one-shot operations) and Stream (persistent streams with per-message settlement).
Simple API
SendQueueMessage
[[nodiscard]] StatusOr<QueueSendResult> SendQueueMessage(const QueueMessage& message);Send a single message to a queue. Returns message ID and timestamps on success.
SendQueueMessages
[[nodiscard]] StatusOr<std::vector<QueueSendResult>> SendQueueMessages(
const std::vector<QueueMessage>& messages);Send a batch of messages. Returns a result for each message.
ReceiveQueueMessages
[[nodiscard]] StatusOr<ReceiveQueueMessagesResponse> ReceiveQueueMessages(
const ReceiveQueueMessagesRequest& request);Receive messages from a queue (simple, one-shot).
AckAllQueueMessages
[[nodiscard]] StatusOr<AckAllQueueMessagesResponse> AckAllQueueMessages(
const AckAllQueueMessagesRequest& request);Acknowledge all pending messages in a queue.
SendQueueMessageSimple
[[nodiscard]] StatusOr<QueueSendResult> SendQueueMessageSimple(
const std::string& channel, const std::string& body,
int expiration_seconds = 0, int delay_seconds = 0);Convenience method with minimal parameters.
Stream API
QueueUpstream
[[nodiscard]] StatusOr<std::unique_ptr<QueueUpstreamHandle>> QueueUpstream(
std::function<void(const QueueUpstreamResult&)> on_result,
std::function<void(const Status&)> on_error);Open an upstream stream for batch sending. Use handle->Send(request_id, messages) to send batches.
NewQueueDownstreamReceiver
[[nodiscard]] StatusOr<std::unique_ptr<QueueDownstreamReceiver>> NewQueueDownstreamReceiver();Create a downstream receiver for polling and settling messages.
PollQueue
[[nodiscard]] StatusOr<PollResponse> PollQueue(const PollRequest& request);Poll a queue for messages (convenience wrapper around downstream receiver).
Types
QueueMessage
Built via QueueMessage::Builder.
| Method | Type | Required | Description |
|---|---|---|---|
SetChannel(s) | string | Yes | Target queue channel |
SetBody(s) | string | No | Message body |
SetMetadata(s) | string | No | String metadata |
SetTags(m) | map<string,string> | No | Key-value tags |
SetExpirationSeconds(n) | int | No | Message TTL (0 = no expiry) |
SetDelaySeconds(n) | int | No | Delay before delivery |
SetMaxReceiveCount(n) | int | No | Max receives before DLQ |
SetMaxReceiveQueue(s) | string | No | Dead-letter queue channel |
PollRequest
| Field | Type | Default | Description |
|---|---|---|---|
channel | string | -- | Queue channel to poll |
max_items | int | 1 | Max messages to receive |
wait_timeout_seconds | int | 5 | Wait timeout |
auto_ack | bool | false | Auto-acknowledge on receipt |
PollResponse
| Method | Description |
|---|---|
messages() | Vector of QueueDownstreamMessage |
is_error() | Whether the response is an error |
error() | Error message |
AckAll() | Acknowledge all messages |
NackAll() | Negative-acknowledge all |
ReQueueAll(ch) | Re-queue all to another channel |
QueueDownstreamMessage
| Method | Description |
|---|---|
message() | The underlying QueueMessage |
sequence() | Sequence number |
transaction_id() | Transaction identifier |
Ack() | Acknowledge this message |
Nack() | Negative-acknowledge |
ReQueue(ch) | Re-queue to another channel |
Quick Usage
// Simple send
auto msg_or = kubemq::QueueMessage::Builder()
.SetChannel("tasks")
.SetBody("process-order")
.Build();
auto result = client->SendQueueMessage(*msg_or);
// Poll with auto-ack
kubemq::PollRequest req;
req.channel = "tasks";
req.max_items = 10;
req.wait_timeout_seconds = 5;
req.auto_ack = true;
auto poll = client->PollQueue(req);
// Stream upstream
auto upstream = client->QueueUpstream(on_result, on_error);
(*upstream)->Send("req-1", messages);See Also
Was this page helpful?