KubeMQ
Client SDKsC++How-to guidesQueues

Requeue All

Re-queue all received KubeMQ Queue messages to another channel with the C++ SDK to redirect undelivered work.

Overview

Requeue all moves an entire batch of polled messages to a different channel in one server-side operation, without republishing them from the client. Reach for it when you need to make a routing decision after looking at a batch — shovel a stuck batch into a review queue, redirect it to a priority pipeline, or migrate messages off a channel that's being retired, all while the source queue is cleared atomically.

It works against the result returned by a manual poll: after receiving messages with auto_ack = false, call poll_result->ReQueueAll(dst_channel) to move every message held in that result to the destination channel, removing them from the source at the same instant. The messages keep their original body, tags, and policies — the broker relocates them, it doesn't recreate them.

Gotchas: requeuing is all-or-nothing for the batch — there's no per-message filter, so split the batch yourself first if only some messages should move. The destination channel is an ordinary queue with no special semantics; nothing consumes it automatically. And it only affects messages still held in that poll result — anything already acked or expired beforehand is gone before ReQueueAll runs. Compare with NackAll(), which returns messages to the same queue instead of moving them elsewhere.

Prerequisites

  • KubeMQ server running on localhost:50000
  • C++ SDK installed (vcpkg or CMake FetchContent)
  • C++17 compiler (GCC 9+, Clang 9+, MSVC 2019+)

Code

main.cc
// Example: queues_stream/requeue_all
//
// Demonstrates moving all messages from one queue to another using ReQueueAll.
// This is useful for routing messages to different processing pipelines.
//
// Channel: cpp-queues.requeue-all
// Client ID: cpp-queues-requeue-all-client
//
// Run with a KubeMQ server on localhost:50000
// (see https://docs.kubemq.io/deploy).

#include <kubemq/kubemq.h>

#include <iostream>

int main() {
    std::cout << "[1] Connecting to localhost:50000" << std::endl;

    kubemq::ClientOptions options;
    options.set_address("localhost", 50000);
    options.set_client_id("cpp-queues-requeue-all-client");

    auto client_result = kubemq::Client::Create(options);
    if (!client_result.ok()) {
        std::cerr << "[ERROR] Failed to create client: " << client_result.status().message()
                  << std::endl;
        return 1;
    }
    auto& client = *client_result;

    std::string src_channel = "cpp-queues.requeue-all";
    std::string dst_channel = "cpp-queues.requeue-all.dest";

    // Send a message to the source queue with a receive policy.
    std::cout << "[2] Sending message to source queue " << src_channel << std::endl;
    auto msg_result = kubemq::QueueMessage::Builder()
                          .SetChannel(src_channel)
                          .SetBody("will be requeued")
                          .SetMaxReceiveCount(3)
                          .Build();
    if (!msg_result.ok()) {
        std::cerr << "[ERROR] Build message: " << msg_result.status().message() << std::endl;
        return 1;
    }
    auto send_result = client->SendQueueMessage(*msg_result);
    if (!send_result.ok()) {
        std::cerr << "[ERROR] SendQueueMessage: " << send_result.status().message() << std::endl;
        return 1;
    }
    std::cout << "[3] Message sent to source queue" << std::endl;

    // Receive from source queue via downstream receiver.
    std::cout << "[4] Opening downstream receiver" << std::endl;
    auto receiver_result = client->NewQueueDownstreamReceiver();
    if (!receiver_result.ok()) {
        std::cerr << "[ERROR] NewQueueDownstreamReceiver: " << receiver_result.status().message()
                  << std::endl;
        return 1;
    }
    auto& receiver = *receiver_result;

    kubemq::PollRequest poll_req;
    poll_req.channel = src_channel;
    poll_req.max_items = 10;
    poll_req.wait_timeout_seconds = 5;
    poll_req.auto_ack = false;

    auto poll_result = receiver->Poll(poll_req);
    if (!poll_result.ok()) {
        std::cerr << "[ERROR] Poll: " << poll_result.status().message() << std::endl;
        return 1;
    }
    if (poll_result->is_error()) {
        std::cerr << "[ERROR] Poll response: " << poll_result->error() << std::endl;
        return 1;
    }

    std::cout << "[5] Received " << poll_result->messages().size() << " messages" << std::endl;
    for (const auto& dm : poll_result->messages()) {
        std::cout << "  body=" << dm.message().body() << " tx=" << dm.transaction_id() << std::endl;
    }

    // ReQueueAll: move all messages to the destination queue.
    if (!poll_result->messages().empty()) {
        std::cout << "[6] ReQueueAll: moving messages to " << dst_channel << std::endl;
        auto requeue_status = poll_result->ReQueueAll(dst_channel);
        if (!requeue_status.ok()) {
            std::cerr << "[ERROR] ReQueueAll: " << requeue_status.message() << std::endl;
            return 1;
        }
        std::cout << "[7] Messages requeued to destination" << std::endl;
    }

    // Close resources explicitly.
    // Note: Destructors also handle cleanup, but explicit calls are shown
    // for clarity and to match Go's defer pattern.
    receiver->Close();

    auto close_status = client->Close();
    if (!close_status.ok()) {
        std::cerr << "[ERROR] Close failed: " << close_status.message() << std::endl;
        return 1;
    }
    std::cout << "[8] Client closed" << std::endl;

    return 0;
}

How It Works

  • Receives messages and re-queues them to a different channel with poll_result->ReQueueAll(channel).
  • Useful for routing messages to retry queues, overflow queues, or alternate processors.
  • The original messages are removed from the source queue.
  • Compare with NackAll() which returns messages to the same queue.

Was this page helpful?

On this page