Purge Queue
Purge all pending messages from a KubeMQ Queue channel using the Java SDK administration API to reset its state.
Overview
Purging a queue is a management-plane operation for wiping a channel's backlog without receiving and discarding messages one at a time. It's the tool you'd reach for when a bad producer floods a channel, when you need a clean slate between test runs, or when you're resetting a queue during a maintenance window — all without deleting and recreating the channel itself. Done properly, it's a single server-side call that tells the broker to acknowledge and drop every pending message on a channel, returning an affected-message count so you can confirm how much backlog was cleared.
Not implemented in the Java SDK.
QueuesClient.purgeQueue()always throwsNotImplementedExceptionin the current SDK version (3.x). The KubeMQ server supports the underlyingAckAllQueueMessagesgRPC call, but the Java SDK does not yet expose it as a stable API. Track progress in the SDK feature matrix.
As a workaround, you can drain a queue by receiving all waiting messages with receiveQueueMessages using autoAckMessages(true), or by using the KubeMQ dashboard/CLI to purge channels in the meantime.
Gotchas: a client-side drain-and-ack loop is not atomic like a real purge — messages can be redelivered to other consumers mid-drain, and it competes with production traffic instead of running as a single server-side operation. It also won't reach messages already delivered to and held by another active consumer.
Prerequisites
- KubeMQ server running on
localhost:50000 - Java SDK installed (see Getting Started)
Steps
- Build a
QueuesPollRequestfor the target channel withautoAckMessages(true)so every fetched message is acknowledged immediately — no separateack()/reject()call needed. - Call
client.receiveQueueMessages(request)in a loop until a poll returns zero messages, to drain a backlog larger than a singlepollMaxMessagesbatch. - Track the running total from each response's message count to confirm how many messages were cleared.
QueuesClient client = QueuesClient.builder().address("localhost:50000").clientId("java-drain-client").build();
int totalDrained = 0;
while (true) {
QueuesPollResponse response = client.receiveQueueMessages(QueuesPollRequest.builder()
.channel("orders")
.pollMaxMessages(100)
.pollWaitTimeoutInSeconds(1)
.autoAckMessages(true)
.build());
if (response.getMessages().isEmpty()) break;
totalDrained += response.getMessages().size();
}
System.out.println("Drained " + totalDrained + " messages.");Related
Was this page helpful?