KubeMQ
LearnQueuesHow-To Guides

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

ScenarioDescription
DevelopmentClear test messages between iterations
TestingReset queue state before test runs
Stuck queuesRemove poison messages blocking consumers
Data cleanupClear 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.

purge.go
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)
purge.py
acked = client.ack_all_queue_messages("orders", wait_time_seconds=5)
print(f"Purged {acked} messages from 'orders'")
purge.ts
await client.purgeQueue('orders');
console.log("Queue 'orders' purged successfully");
Purge.java
// 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());
Purge.cs
var result = await client.PurgeQueueAsync("orders");
Console.WriteLine($"Purged {result.AffectedMessages} messages from 'orders'");
Purge.kt
val purged = client.purgeQueuesChannel("orders")
println("Purged $purged messages from 'orders'")
purge.cpp
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;
}
purge.rs
client.ack_all_queue_messages("orders").await?;
println!("Queue 'orders' purged successfully");
purge.rb
client.purge_queue_channel(channel_name: 'orders')
puts "Queue 'orders' purged successfully"
purge.exs
: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.

verify_purge.go
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))
verify_purge.py
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)}")
verify_purge.ts
const remaining = await client.peekQueueMessages({
  channel: 'orders',
  maxMessages: 10,
  waitTimeoutSeconds: 2,
});
console.log(`Messages remaining: ${remaining.length}`);
VerifyPurge.java
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());
VerifyPurge.cs
var response = await client.ReceiveQueueMessagesAsync(new QueuePollRequest
{
    Channel = "orders",
    MaxMessages = 10,
    WaitTimeoutSeconds = 2,
    AutoAck = true,
});
Console.WriteLine($"Messages remaining: {response.Messages.Count}");
VerifyPurge.kt
val response = client.peekQueueMessages {
    channel = "orders"
    maxNumberOfMessages = 10
    waitTimeSeconds = 2
}
println("Messages remaining: ${response.messages.size}")
verify_purge.cpp
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;
verify_purge.rs
// 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());
verify_purge.rb
remaining = client.receive_queue_messages(
  channel: 'orders',
  max_messages: 10,
  wait_timeout_seconds: 2,
  peek: true,
)
puts "Messages remaining: #{remaining.size}"
verify_purge.exs
{: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?

On this page