KubeMQ
LearnQueuesTutorials

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

Steps

Set Up the Client

main.go
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()
main.py
from kubemq.queues import Client as QueuesClient
from kubemq import QueueMessage

client = QueuesClient(
    address="localhost:50000",
    client_id="order-processor",
)
main.ts
import { KubeMQClient, createQueueMessage } from 'kubemq-js';

const client = await KubeMQClient.create({
  address: 'localhost:50000',
  clientId: 'order-processor',
});
Main.java
QueuesClient client = QueuesClient.builder()
    .address("localhost:50000")
    .clientId("order-processor")
    .build();
Program.cs
using KubeMQ.Sdk.Client;

await using var client = new KubeMQClient(new KubeMQClientOptions());
await client.ConnectAsync();
Main.kt
val client = QueuesClient("localhost:50000")
main.cpp
#include <kubemq/client.h>
auto client = kubemq::QueuesClient("localhost:50000");
main.rs
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?;
main.rb
require 'kubemq'

client = KubeMQ::QueuesClient.new(
  address: 'localhost:50000',
  client_id: 'order-processor',
)
main.exs
{: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.close
Enum.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?

On this page