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
- KubeMQ server running on
localhost:50000 - SDK installed (Getting Started)
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.
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))
}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')}")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)}`);
}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()));
}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)}");
}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)}")
}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;
}// 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());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"{: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 Case | Description |
|---|---|
| Monitoring | Check queue depth and message contents without affecting consumers |
| Debugging | Inspect message payloads and metadata to diagnose processing issues |
| Delayed decisions | Preview messages before deciding whether to consume them |
| Queue health | Verify 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?