Purge Queue
Purge a KubeMQ queue channel, permanently clearing all pending messages, using the REST API, CLI, or an SDK admin call.
Purging is irreversible. All messages in the queue are permanently deleted. Use with caution in production environments.
How Purge Works
Purging a queue acknowledges every pending message at once, removing them from the channel. Consumers that poll afterward find the queue empty.
Purge acknowledges all pending messages in a single call, leaving the channel empty.
When to Purge
| Scenario | Description |
|---|---|
| Development | Clear test messages between iterations |
| Testing | Reset queue state before test runs |
| Stuck queues | Remove poison messages blocking consumers |
| Data cleanup | Clear obsolete messages after a schema change |
Purge All Messages
Purging is implemented as an "acknowledge all" operation on the channel — every pending message is acked at once and removed.
resp, err := client.AckAllQueueMessages(ctx, &kubemq.AckAllQueueMessagesRequest{
Channel: "orders",
WaitTimeSeconds: 5,
})
if err != nil {
log.Fatal(err)
}
fmt.Printf("Purged %d messages from 'orders'\n", resp.AffectedMessages)acked = client.ack_all_queue_messages("orders", wait_time_seconds=5)
print(f"Purged {acked} messages from 'orders'")await client.purgeQueue('orders');
console.log("Queue 'orders' purged successfully");// purgeQueue is not yet exposed by the Java v2 SDK.
// Purge by acknowledging all messages on the channel instead.
QueuesPollResponse response = client.receiveQueueMessages(
QueuesPollRequest.builder()
.channel("orders")
.pollMaxMessages(1024)
.pollWaitTimeoutInSeconds(2)
.autoAckMessages(true)
.build());
System.out.printf("Purged %d messages from 'orders'%n", response.getMessages().size());var result = await client.PurgeQueueAsync("orders");
Console.WriteLine($"Purged {result.AffectedMessages} messages from 'orders'");val purged = client.purgeQueuesChannel("orders")
println("Purged $purged messages from 'orders'")kubemq::AckAllQueueMessagesRequest req;
req.channel = "orders";
req.wait_time_seconds = 5;
auto result = client->AckAllQueueMessages(req);
if (result.ok() && !result->is_error) {
std::cout << "Purged " << result->affected_messages << " messages from 'orders'" << std::endl;
}client.ack_all_queue_messages("orders").await?;
println!("Queue 'orders' purged successfully");client.purge_queue_channel(channel_name: 'orders')
puts "Queue 'orders' purged successfully":ok = KubeMQ.Client.purge_queue_channel(client, "orders")
IO.puts("Queue 'orders' purged successfully")Verify the Queue is Empty
After purging, poll the queue to confirm no messages remain.
resp, err := client.PollQueue(ctx, &kubemq.PollRequest{
Channel: "orders",
MaxItems: 10,
WaitTimeoutSeconds: 2,
AutoAck: true,
})
if err != nil {
log.Fatal(err)
}
fmt.Printf("Messages remaining: %d\n", len(resp.Messages))response = client.receive_queue_messages(
channel="orders",
max_messages=10,
wait_timeout_in_seconds=2,
auto_ack=True,
)
print(f"Messages remaining: {len(response.messages)}")const remaining = await client.peekQueueMessages({
channel: 'orders',
maxMessages: 10,
waitTimeoutSeconds: 2,
});
console.log(`Messages remaining: ${remaining.length}`);QueuesPollResponse response = client.receiveQueueMessages(
QueuesPollRequest.builder()
.channel("orders")
.pollMaxMessages(10)
.pollWaitTimeoutInSeconds(2)
.autoAckMessages(true)
.build());
System.out.printf("Messages remaining: %d%n", response.getMessages().size());var response = await client.ReceiveQueueMessagesAsync(new QueuePollRequest
{
Channel = "orders",
MaxMessages = 10,
WaitTimeoutSeconds = 2,
AutoAck = true,
});
Console.WriteLine($"Messages remaining: {response.Messages.Count}");val response = client.peekQueueMessages {
channel = "orders"
maxNumberOfMessages = 10
waitTimeSeconds = 2
}
println("Messages remaining: ${response.messages.size}")kubemq::ReceiveQueueMessagesRequest req;
req.channel = "orders";
req.max_number_of_messages = 10;
req.wait_time_seconds = 2;
req.is_peek = true;
auto response = client->ReceiveQueueMessages(req);
std::cout << "Messages remaining: " << response->messages.size() << std::endl;// receive_queue_messages(channel, max_messages, wait_seconds, is_peek)
let remaining = client.receive_queue_messages("orders", 10, 2, true).await?;
println!("Messages remaining: {}", remaining.len());remaining = client.receive_queue_messages(
channel: 'orders',
max_messages: 10,
wait_timeout_seconds: 2,
peek: true,
)
puts "Messages remaining: #{remaining.size}"{:ok, poll} = KubeMQ.Client.poll_queue(client,
channel: "orders",
max_items: 10,
wait_timeout: 2_000
)
IO.puts("Messages remaining: #{length(poll.messages)}")Next Steps
Was this page helpful?