KubeMQ
LearnQueuesTutorials

Peek Messages

Inspect queue contents without consuming messages.

What You Will Build

A queue inspector that reads messages without removing them from the queue. Peeked messages remain available for normal consumers.

Prerequisites

Peek Queue Messages

Set isPeek to true (or use the peek-specific API) to inspect messages without consuming them. The messages are not hidden from other consumers and no acknowledgment is needed.

Peek reads the queue without removing or hiding messages — a normal consumer still polls, receives, and acks the same messages.

peek.go
resp, err := client.PollQueue(ctx, &kubemq.PollRequest{
    Channel:            "orders",
    MaxItems:           10,
    WaitTimeoutSeconds: 5,
    IsPeek:             true,
})
if err != nil {
    log.Fatal(err)
}

fmt.Printf("Queue has %d messages:\n", len(resp.Messages))
for _, m := range resp.Messages {
    fmt.Printf("  [%d] id=%s body=%s\n",
        m.Message.Attributes.Sequence,
        m.Message.MessageID,
        string(m.Message.Body))
}
peek.py
response = client.receive_queue_messages(
    channel="orders",
    max_messages=10,
    wait_timeout_in_seconds=5,
    is_peek=True,
)

print(f"Queue has {len(response.messages)} messages:")
for msg in response.messages:
    print(f"  [{msg.sequence}] id={msg.id} body={msg.body.decode('utf-8')}")
peek.ts
const messages = await client.receiveQueueMessages({
  channel: 'orders',
  maxMessages: 10,
  waitTimeoutSeconds: 5,
  isPeek: true,
});

console.log(`Queue has ${messages.length} messages:`);
for (const msg of messages) {
  console.log(`  [${msg.sequence}] id=${msg.messageId} body=${new TextDecoder().decode(msg.body)}`);
}
Peek.java
ReceiveQueueMessagesResponse response = client.receiveQueueMessages(
    ReceiveQueueMessagesRequest.builder()
        .channel("orders")
        .maxMessages(10)
        .waitTimeoutSeconds(5)
        .isPeek(true)
        .build());

System.out.printf("Queue has %d messages:%n", response.getMessages().size());
for (QueueMessageReceived msg : response.getMessages()) {
    System.out.printf("  [%d] id=%s body=%s%n",
        msg.getSequence(), msg.getMessageId(), new String(msg.getBody()));
}
Peek.cs
var response = await client.ReceiveQueueMessagesAsync(new ReceiveQueueMessagesRequest
{
    Channel = "orders",
    MaxMessages = 10,
    WaitTimeoutSeconds = 5,
    IsPeek = true,
});

Console.WriteLine($"Queue has {response.Messages.Count} messages:");
foreach (var msg in response.Messages)
{
    Console.WriteLine($"  [{msg.Sequence}] id={msg.MessageId} body={Encoding.UTF8.GetString(msg.Body.Span)}");
}
Peek.kt
val response = client.receiveQueueMessages(
    channel = "orders",
    maxMessages = 10,
    waitTimeoutSeconds = 5,
    isPeek = true
)

println("Queue has ${response.messages.size} messages:")
for (msg in response.messages) {
    println("  [${msg.sequence}] id=${msg.messageId} body=${String(msg.body)}")
}
peek.cpp
auto response = client.receiveQueueMessages("orders", 10, 5, true);

std::cout << "Queue has " << response.messages.size() << " messages:" << std::endl;
for (const auto& msg : response.messages) {
    std::cout << "  [" << msg.sequence << "] id=" << msg.messageId
              << " body=" << msg.body << std::endl;
}
peek.rs
// receive_queue_messages(channel, max_messages, wait_seconds, is_peek)
let peeked = client
    .receive_queue_messages("orders", 10, 5, true)
    .await?;

println!("Queue has {} messages:", peeked.len());
for m in &peeked {
    println!("  id={} body={}", m.id, String::from_utf8_lossy(&m.body));
}

// Messages remain in the queue — still available for a normal receive
let received = client.receive_queue_messages("orders", 10, 5, false).await?;
println!("Received {} messages after peek", received.len());
peek.rb
peeked = client.receive_queue_messages(
  channel: 'orders',
  max_messages: 10,
  wait_timeout_seconds: 5,
  peek: true
)

puts "Queue has #{peeked.size} messages (not consumed):"
peeked.each { |msg| puts "  id=#{msg.id} body=#{msg.body}" }

# Messages remain in the queue — still available for a normal receive
received = client.receive_queue_messages(channel: 'orders', max_messages: 10, wait_timeout_seconds: 5)
puts "Received #{received.size} messages after peek"
peek.exs
{:ok, result} =
  KubeMQ.Client.receive_queue_messages(client, "orders",
    max_messages: 10,
    wait_timeout: 5_000,
    is_peek: true
  )

IO.puts("Queue has #{result.messages_received} messages (is_peek: #{result.is_peek}):")
Enum.each(result.messages, fn msg ->
  IO.puts("  id=#{msg.id} body=#{msg.body}")
end)

# Messages remain in the queue — still available for a normal receive
{:ok, received} = KubeMQ.Client.receive_queue_messages(client, "orders", max_messages: 10, wait_timeout: 5_000)
IO.puts("Received #{received.messages_received} messages after peek")

Use Cases

Use CaseDescription
MonitoringCheck queue depth and message contents without affecting consumers
DebuggingInspect message payloads and metadata to diagnose processing issues
Delayed decisionsPreview messages before deciding whether to consume them
Queue healthVerify messages are being produced correctly before starting consumers

Peek does not change the message state. The receiveCount is not incremented, and messages remain fully available for normal consumers.

Next Steps

Was this page helpful?

On this page