Send & Receive Messages
Learn the complete queue send-receive-acknowledge cycle with metadata, tags, and error handling.
What You Will Build
A producer that sends order processing tasks with metadata and tags, and a consumer that receives, processes, and acknowledges each message. The example includes error handling and message introspection.
The send-receive-acknowledge cycle: the producer sends, the consumer polls and processes, and the ack removes the message from the queue.
Prerequisites
- KubeMQ server running on
localhost:50000 - SDK installed (Getting Started)
Steps
Set Up the Client
package main
import (
"context"
"encoding/json"
"fmt"
"log"
"time"
"github.com/kubemq-io/kubemq-go/v2"
)
func main() {
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
defer cancel()
client, err := kubemq.NewClient(ctx,
kubemq.WithAddress("localhost", 50000),
kubemq.WithClientId("order-processor"),
)
if err != nil {
log.Fatal(err)
}
defer client.Close()from kubemq.queues import Client as QueuesClient
from kubemq import QueueMessage
client = QueuesClient(
address="localhost:50000",
client_id="order-processor",
)import { KubeMQClient, createQueueMessage } from 'kubemq-js';
const client = await KubeMQClient.create({
address: 'localhost:50000',
clientId: 'order-processor',
});QueuesClient client = QueuesClient.builder()
.address("localhost:50000")
.clientId("order-processor")
.build();using KubeMQ.Sdk.Client;
await using var client = new KubeMQClient(new KubeMQClientOptions());
await client.ConnectAsync();val client = QueuesClient("localhost:50000")#include <kubemq/client.h>
auto client = kubemq::QueuesClient("localhost:50000");use kubemq::prelude::*;
use kubemq::QueueMessageBuilder;
use std::collections::HashMap;
#[tokio::main]
async fn main() -> kubemq::Result<()> {
let client = KubemqClient::builder()
.host("localhost")
.port(50000)
.client_id("order-processor")
.build()
.await?;require 'kubemq'
client = KubeMQ::QueuesClient.new(
address: 'localhost:50000',
client_id: 'order-processor',
){:ok, client} = KubeMQ.Client.start_link(
address: "localhost:50000",
client_id: "order-processor"
)Send a Message with Metadata and Tags
Create a queue message with a JSON body, metadata string, and tags for downstream routing.
order := map[string]interface{}{
"orderId": "ORD-5001",
"items": 3,
"total": 149.97,
}
body, _ := json.Marshal(order)
msg := kubemq.NewQueueMessage().
SetChannel("orders").
SetBody(body).
SetMetadata("order.created").
SetTags(map[string]string{
"region": "us-east",
"priority": "high",
})
result, err := client.SendQueueMessage(ctx, msg)
if err != nil {
log.Fatal(err)
}
if result.IsError {
log.Fatalf("Send failed: %s", result.Error)
}
fmt.Printf("Sent: id=%s, sentAt=%d\n", result.MessageID, result.SentAt)import json
order = {"orderId": "ORD-5001", "items": 3, "total": 149.97}
result = client.send_queue_message(
QueueMessage(
channel="orders",
body=json.dumps(order).encode(),
metadata="order.created",
tags={"region": "us-east", "priority": "high"},
)
)
print(f"Sent: id={result.id}, sentAt={result.sent_at}")const result = await client.sendQueueMessage(
createQueueMessage({
channel: 'orders',
body: JSON.stringify({ orderId: 'ORD-5001', items: 3, total: 149.97 }),
metadata: 'order.created',
tags: { region: 'us-east', priority: 'high' },
}),
);
console.log(`Sent: id=${result.messageId}, sentAt=${result.sentAt}`);QueueMessage msg = QueueMessage.builder()
.channel("orders")
.body("{\"orderId\":\"ORD-5001\",\"items\":3,\"total\":149.97}".getBytes())
.metadata("order.created")
.tags(Map.of("region", "us-east", "priority", "high"))
.build();
SendQueueMessageResult result = client.sendQueueMessage(msg);
System.out.printf("Sent: id=%s, sentAt=%d%n", result.getMessageId(), result.getSentAt());var result = await client.SendQueueMessageAsync(new QueueMessage
{
Channel = "orders",
Body = Encoding.UTF8.GetBytes("{\"orderId\":\"ORD-5001\",\"items\":3,\"total\":149.97}"),
Metadata = "order.created",
Tags = new Dictionary<string, string>
{
["region"] = "us-east",
["priority"] = "high"
}
});
Console.WriteLine($"Sent: id={result.MessageId}, sentAt={result.SentAt}");val result = client.sendQueueMessage(QueueMessage(
channel = "orders",
body = """{"orderId":"ORD-5001","items":3,"total":149.97}""".toByteArray(),
metadata = "order.created",
tags = mapOf("region" to "us-east", "priority" to "high")
))
println("Sent: id=${result.messageId}, sentAt=${result.sentAt}")kubemq::QueueMessage msg;
msg.channel = "orders";
msg.body = R"({"orderId":"ORD-5001","items":3,"total":149.97})";
msg.metadata = "order.created";
msg.tags = {{"region", "us-east"}, {"priority", "high"}};
auto result = client.sendQueueMessage(msg);
std::cout << "Sent: id=" << result.messageId << std::endl; let order = r#"{"orderId":"ORD-5001","items":3,"total":149.97}"#;
let mut tags = HashMap::new();
tags.insert("region".to_string(), "us-east".to_string());
tags.insert("priority".to_string(), "high".to_string());
let msg = QueueMessageBuilder::new()
.channel("orders")
.body(order.as_bytes().to_vec())
.metadata("order.created")
.tags(tags)
.build();
let result = client.send_queue_message(msg).await?;
println!("Sent: id={}, sent_at={}", result.message_id, result.sent_at);msg = KubeMQ::Queues::QueueMessage.new(
channel: 'orders',
body: '{"orderId":"ORD-5001","items":3,"total":149.97}',
metadata: 'order.created',
tags: { 'region' => 'us-east', 'priority' => 'high' }
)
result = client.send_queue_message(msg)
puts "Sent: id=#{result.id}, error?=#{result.error?}"msg = KubeMQ.QueueMessage.new(
channel: "orders",
body: ~s({"orderId":"ORD-5001","items":3,"total":149.97}),
metadata: "order.created",
tags: %{"region" => "us-east", "priority" => "high"}
)
{:ok, result} = KubeMQ.Client.send_queue_message(client, msg)
IO.puts("Sent: id=#{result.message_id}")Receive Messages
Poll the queue for available messages. You control how many messages to fetch and how long to wait.
resp, err := client.PollQueue(ctx, &kubemq.PollRequest{
Channel: "orders",
MaxItems: 10,
WaitTimeoutSeconds: 5,
AutoAck: false,
})
if err != nil {
log.Fatal(err)
}
fmt.Printf("Received %d messages\n", len(resp.Messages))response = client.receive_queue_messages(
channel="orders",
max_messages=10,
wait_timeout_in_seconds=5,
)
print(f"Received {len(response.messages)} messages")const messages = await client.receiveQueueMessages({
channel: 'orders',
maxMessages: 10,
waitTimeoutSeconds: 5,
});
console.log(`Received ${messages.length} messages`);ReceiveQueueMessagesResponse response = client.receiveQueueMessages(
ReceiveQueueMessagesRequest.builder()
.channel("orders")
.maxMessages(10)
.waitTimeoutSeconds(5)
.build());
System.out.printf("Received %d messages%n", response.getMessages().size());var response = await client.ReceiveQueueMessagesAsync(new ReceiveQueueMessagesRequest
{
Channel = "orders",
MaxMessages = 10,
WaitTimeoutSeconds = 5,
});
Console.WriteLine($"Received {response.Messages.Count} messages");val response = client.receiveQueueMessages(
channel = "orders",
maxMessages = 10,
waitTimeoutSeconds = 5
)
println("Received ${response.messages.size} messages")auto response = client.receiveQueueMessages("orders", 10, 5);
std::cout << "Received " << response.messages.size() << " messages" << std::endl; // receive_queue_messages(channel, max_items, wait_seconds, auto_ack)
let messages = client
.receive_queue_messages("orders", 10, 5, false)
.await?;
println!("Received {} messages", messages.len());messages = client.receive_queue_messages(
channel: 'orders',
max_messages: 10,
wait_timeout_seconds: 5
)
puts "Received #{messages.size} messages"{:ok, response} =
KubeMQ.Client.receive_queue_messages(client, "orders",
max_messages: 10,
wait_timeout: 5_000
)
IO.puts("Received #{response.messages_received} messages")Process and Acknowledge
Inspect each message's body, metadata, and tags, then acknowledge to remove it from the queue.
for _, dm := range resp.Messages {
fmt.Printf(" ID: %s\n", dm.Message.MessageID)
fmt.Printf(" Body: %s\n", string(dm.Message.Body))
fmt.Printf(" Metadata: %s\n", dm.Message.Metadata)
fmt.Printf(" Tags: %v\n", dm.Message.Tags)
}
if err := resp.AckAll(); err != nil {
log.Fatal(err)
}
fmt.Println("All messages acknowledged")
}for msg in response.messages:
print(f" ID: {msg.id}")
print(f" Body: {msg.body.decode('utf-8')}")
print(f" Metadata: {msg.metadata}")
print(f" Tags: {msg.tags}")
msg.ack()
print(" Acknowledged")
client.close()for (const msg of messages) {
console.log(' ID: ', msg.messageId);
console.log(' Body: ', new TextDecoder().decode(msg.body));
console.log(' Metadata:', msg.metadata);
console.log(' Tags: ', msg.tags);
await msg.ack();
console.log(' Acknowledged');
}
await client.close();for (QueueMessageReceived msg : response.getMessages()) {
System.out.printf(" ID: %s%n", msg.getMessageId());
System.out.printf(" Body: %s%n", new String(msg.getBody()));
System.out.printf(" Metadata: %s%n", msg.getMetadata());
System.out.printf(" Tags: %s%n", msg.getTags());
msg.ack();
System.out.println(" Acknowledged");
}
client.close();foreach (var msg in response.Messages)
{
Console.WriteLine($" ID: {msg.MessageId}");
Console.WriteLine($" Body: {Encoding.UTF8.GetString(msg.Body.Span)}");
Console.WriteLine($" Metadata: {msg.Metadata}");
Console.WriteLine($" Tags: {string.Join(", ", msg.Tags)}");
await msg.AckAsync();
Console.WriteLine(" Acknowledged");
}for (msg in response.messages) {
println(" ID: ${msg.messageId}")
println(" Body: ${String(msg.body)}")
println(" Metadata: ${msg.metadata}")
println(" Tags: ${msg.tags}")
msg.ack()
println(" Acknowledged")
}
client.close()for (const auto& msg : response.messages) {
std::cout << " ID: " << msg.messageId << std::endl;
std::cout << " Body: " << msg.body << std::endl;
std::cout << " Metadata: " << msg.metadata << std::endl;
msg.ack();
std::cout << " Acknowledged" << std::endl;
} // The unary receive returns plain messages; settle them with ack_all.
// For per-message ack/nack, use the queue stream API.
for m in &messages {
println!(" ID: {}", m.id);
println!(" Body: {}", String::from_utf8_lossy(&m.body));
println!(" Metadata: {}", m.metadata);
println!(" Tags: {:?}", m.tags);
}
let ack = AckAllQueueMessagesRequest {
request_id: String::new(),
client_id: String::new(),
channel: "orders".to_string(),
wait_time_seconds: 5,
};
client.ack_all_queue_messages(&ack).await?;
println!("All messages acknowledged");
client.close().await?;
Ok(())
}messages.each do |m|
puts " ID: #{m.id}"
puts " Body: #{m.body}"
puts " Metadata: #{m.metadata}"
puts " Tags: #{m.tags}"
end
# Settle the polled messages. For per-message ack/nack, use the stream receiver.
affected = client.ack_all_queue_messages(channel: 'orders', wait_timeout_seconds: 5)
puts "Acknowledged #{affected} messages"
client.closeEnum.each(response.messages, fn m ->
IO.puts(" ID: #{m.id}")
IO.puts(" Body: #{m.body}")
IO.puts(" Metadata: #{m.metadata}")
IO.puts(" Tags: #{inspect(m.tags)}")
end)
# Settle the received messages. For per-message ack/nack, use poll_queue (stream API).
{:ok, ack} = KubeMQ.Client.ack_all_queue_messages(client, "orders", wait_timeout: 5_000)
IO.puts("Acknowledged #{ack.affected_messages} messages")
KubeMQ.Client.close(client)Next Steps
Was this page helpful?