KubeMQ
LearnQueues

Getting Started with Queues

Send and receive your first KubeMQ queue message with acknowledgment in 5 minutes.

Prerequisites: KubeMQ server running on localhost:50000 and your SDK installed. See Getting Started for setup.

What You Will Build

An order message sender that places a task on a queue, and a receiver that pulls the task, processes it, and acknowledges completion.

Steps

Install the SDK

go get github.com/kubemq-io/kubemq-go/v2
pip install kubemq
npm install kubemq-js
<dependency>
    <groupId>io.kubemq.sdk</groupId>
    <artifactId>kubemq-sdk-Java</artifactId>
    <version>3.1.1</version>
</dependency>
dotnet add package KubeMQ.SDK.CSharp --version 3.0.1
implementation("io.kubemq.sdk:kubemq-sdk-kotlin:1.0.1")
vcpkg install kubemq
[dependencies]
kubemq = "1.0.1"
tokio = { version = "1", features = ["full"] }
gem install kubemq
def deps do
  [{:kubemq, "~> 1.0"}]
end

Create the Sender

Connect to KubeMQ and send an order message to the orders queue.

sender.go
package main

import (
    "context"
    "fmt"
    "log"

    "github.com/kubemq-io/kubemq-go/v2"
)

func main() {
    ctx := context.Background()
    client, err := kubemq.NewClient(ctx,
        kubemq.WithAddress("localhost", 50000),
        kubemq.WithClientId("order-sender"),
    )
    if err != nil {
        log.Fatal(err)
    }
    defer client.Close()

    msg := kubemq.NewQueueMessage().
        SetChannel("orders").
        SetBody([]byte(`{"orderId":"ORD-1234","total":99.99}`)).
        SetMetadata("order.created")

    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\n", result.MessageID)
}
sender.py
from kubemq.queues import Client as QueuesClient
from kubemq import QueueMessage

client = QueuesClient(
    address="localhost:50000",
    client_id="order-sender",
)

result = client.send_queue_message(
    QueueMessage(
        channel="orders",
        body=b'{"orderId":"ORD-1234","total":99.99}',
        metadata="order.created",
    )
)
print(f"Sent: id={result.id}")
client.close()
sender.ts
import { KubeMQClient, createQueueMessage } from 'kubemq-js';

const client = await KubeMQClient.create({
  address: 'localhost:50000',
  clientId: 'order-sender',
});

const result = await client.sendQueueMessage(
  createQueueMessage({
    channel: 'orders',
    body: JSON.stringify({ orderId: 'ORD-1234', total: 99.99 }),
    metadata: 'order.created',
  }),
);
console.log('Sent:', result.messageId);
await client.close();
Sender.java
import io.kubemq.sdk.queues.QueuesClient;
import io.kubemq.sdk.queues.QueueMessage;
import io.kubemq.sdk.queues.QueueSendResult;

QueuesClient client = QueuesClient.builder()
    .address("localhost:50000")
    .clientId("order-sender")
    .build();

QueueMessage msg = QueueMessage.builder()
    .channel("orders")
    .body("{\"orderId\":\"ORD-1234\",\"total\":99.99}".getBytes())
    .metadata("order.created")
    .build();

QueueSendResult result = client.sendQueueMessage(msg);
System.out.println("Sent: id=" + result.getId());
client.close();
Sender.cs
using System.Text;
using KubeMQ.Sdk.Client;
using KubeMQ.Sdk.Queues;

await using var client = new KubeMQClient(new KubeMQClientOptions());
await client.ConnectAsync();

var result = await client.SendQueueMessageAsync(new QueueMessage
{
    Channel = "orders",
    Body = Encoding.UTF8.GetBytes("{\"orderId\":\"ORD-1234\",\"total\":99.99}"),
    Metadata = "order.created"
});
Console.WriteLine($"Sent: id={result.MessageId}");
Sender.kt
import io.kubemq.sdk.client.KubeMQClient
import io.kubemq.sdk.queues.QueueMessage
import kotlinx.coroutines.runBlocking

fun main() = runBlocking {
    val client = KubeMQClient.queues {
        address = "localhost:50000"
        clientId = "order-sender"
    }

    val result = client.sendQueuesMessage(QueueMessage(
        channel = "orders",
        body = """{"orderId":"ORD-1234","total":99.99}""".toByteArray(),
        metadata = "order.created"
    ))
    println("Sent: id=${result.messageId}")
    client.close()
}
sender.cpp
#include <kubemq/client.h>
#include <iostream>

auto client = kubemq::QueuesClient("localhost:50000");

kubemq::QueueMessage msg;
msg.channel = "orders";
msg.body = R"({"orderId":"ORD-1234","total":99.99})";
msg.metadata = "order.created";

auto result = client.sendQueueMessage(msg);
std::cout << "Sent: id=" << result.messageId << std::endl;
sender.rs
use kubemq::prelude::*;
use kubemq::QueueMessageBuilder;

#[tokio::main]
async fn main() -> kubemq::Result<()> {
    let client = KubemqClient::builder()
        .host("localhost")
        .port(50000)
        .client_id("order-sender")
        .build()
        .await?;

    let msg = QueueMessageBuilder::new()
        .channel("orders")
        .body(br#"{"orderId":"ORD-1234","total":99.99}"#.to_vec())
        .metadata("order.created")
        .build();

    let result = client.send_queue_message(msg).await?;
    println!("Sent: id={}", result.message_id);

    client.close().await?;
    Ok(())
}
sender.rb
require 'kubemq'

client = KubeMQ::QueuesClient.new(
  address: 'localhost:50000',
  client_id: 'order-sender'
)

msg = KubeMQ::Queues::QueueMessage.new(
  channel: 'orders',
  body: '{"orderId":"ORD-1234","total":99.99}',
  metadata: 'order.created'
)
result = client.send_queue_message(msg)
puts "Sent: id=#{result.id}"

client.close
sender.exs
{:ok, client} =
  KubeMQ.Client.start_link(address: "localhost:50000", client_id: "order-sender")

msg =
  KubeMQ.QueueMessage.new(
    channel: "orders",
    body: ~s({"orderId":"ORD-1234","total":99.99}),
    metadata: "order.created"
  )

{:ok, result} = KubeMQ.Client.send_queue_message(client, msg)
IO.puts("Sent: id=#{result.message_id}")

KubeMQ.Client.close(client)

Create the Receiver

In a separate terminal, receive the message and acknowledge it.

receiver.go
package main

import (
    "context"
    "fmt"
    "log"

    "github.com/kubemq-io/kubemq-go/v2"
)

func main() {
    ctx := context.Background()
    client, err := kubemq.NewClient(ctx,
        kubemq.WithAddress("localhost", 50000),
        kubemq.WithClientId("order-receiver"),
    )
    if err != nil {
        log.Fatal(err)
    }
    defer client.Close()

    resp, err := client.PollQueue(ctx, &kubemq.PollRequest{
        Channel:            "orders",
        MaxItems:           1,
        WaitTimeoutSeconds: 10,
        AutoAck:            false,
    })
    if err != nil {
        log.Fatal(err)
    }

    for _, m := range resp.Messages {
        fmt.Printf("Received: %s\n", string(m.Message.Body))
        fmt.Printf("Metadata: %s\n", m.Message.Metadata)
    }
    if err := resp.AckAll(); err != nil {
        log.Fatal(err)
    }
    fmt.Println("All messages acknowledged")
}
receiver.py
from kubemq.queues import Client as QueuesClient

client = QueuesClient(
    address="localhost:50000",
    client_id="order-receiver",
)

response = client.receive_queue_messages(
    channel="orders",
    max_messages=1,
    wait_timeout_in_seconds=10,
)
for msg in response.messages:
    print(f"Received: {msg.body.decode('utf-8')}")
    print(f"Metadata: {msg.metadata}")
    msg.ack()
    print("Message acknowledged")

client.close()
receiver.ts
import { KubeMQClient } from 'kubemq-js';

const client = await KubeMQClient.create({
  address: 'localhost:50000',
  clientId: 'order-receiver',
});

const messages = await client.receiveQueueMessages({
  channel: 'orders',
  maxMessages: 1,
  waitTimeoutSeconds: 10,
});

for (const msg of messages) {
  console.log('Received:', new TextDecoder().decode(msg.body));
  console.log('Metadata:', msg.metadata);
  await msg.ack();
  console.log('Message acknowledged');
}

await client.close();
Receiver.java
import io.kubemq.sdk.queues.QueuesClient;
import io.kubemq.sdk.queues.QueuesPollRequest;
import io.kubemq.sdk.queues.QueuesPollResponse;
import io.kubemq.sdk.queues.QueueMessageReceived;

QueuesClient client = QueuesClient.builder()
    .address("localhost:50000")
    .clientId("order-receiver")
    .build();

QueuesPollResponse response = client.receiveQueueMessages(
    QueuesPollRequest.builder()
        .channel("orders")
        .pollMaxMessages(1)
        .pollWaitTimeoutInSeconds(10)
        .autoAckMessages(false)
        .build());

for (QueueMessageReceived msg : response.getMessages()) {
    System.out.println("Received: " + new String(msg.getBody()));
    System.out.println("Metadata: " + msg.getMetadata());
    msg.ack();
    System.out.println("Message acknowledged");
}

client.close();
Receiver.cs
using System.Text;
using KubeMQ.Sdk.Client;
using KubeMQ.Sdk.Queues;

await using var client = new KubeMQClient(new KubeMQClientOptions());
await client.ConnectAsync();

var response = await client.ReceiveQueueMessagesAsync(new QueuePollRequest
{
    Channel = "orders",
    MaxMessages = 1,
    WaitTimeoutSeconds = 10,
    AutoAck = false,
});

foreach (var msg in response.Messages)
{
    Console.WriteLine($"Received: {Encoding.UTF8.GetString(msg.Body.Span)}");
    Console.WriteLine($"Metadata: {msg.Metadata}");
    await msg.AckAsync();
    Console.WriteLine("Message acknowledged");
}
Receiver.kt
import io.kubemq.sdk.client.KubeMQClient
import kotlinx.coroutines.runBlocking

fun main() = runBlocking {
    val client = KubeMQClient.queues {
        address = "localhost:50000"
        clientId = "order-receiver"
    }

    val response = client.receiveQueuesMessages {
        channel = "orders"
        maxItems = 1
        waitTimeoutMs = 10_000
        autoAck = false
    }
    for (msg in response.messages) {
        println("Received: ${String(msg.body)}")
        println("Metadata: ${msg.metadata}")
        msg.ack()
        println("Message acknowledged")
    }

    client.close()
}
receiver.cpp
#include <kubemq/client.h>
#include <iostream>

auto client = kubemq::QueuesClient("localhost:50000");

auto response = client.receiveQueueMessages("orders", 1, 10);
for (const auto& msg : response.messages) {
    std::cout << "Received: " << msg.body << std::endl;
    std::cout << "Metadata: " << msg.metadata << std::endl;
    msg.ack();
    std::cout << "Message acknowledged" << std::endl;
}
receiver.rs
use kubemq::prelude::*;

#[tokio::main]
async fn main() -> kubemq::Result<()> {
    let client = KubemqClient::builder()
        .host("localhost")
        .port(50000)
        .client_id("order-receiver")
        .build()
        .await?;

    let (response, _receiver) = client
        .poll_queue(PollRequest {
            channel: "orders".to_string(),
            max_items: 1,
            wait_timeout_seconds: 10,
            auto_ack: false,
        })
        .await?;

    for m in &response.messages {
        println!("Received: {}", String::from_utf8_lossy(&m.message.body));
        println!("Metadata: {}", m.message.metadata);
    }
    response.ack_all().await?;
    println!("All messages acknowledged");

    client.close().await?;
    Ok(())
}
receiver.rb
require 'kubemq'

client = KubeMQ::QueuesClient.new(
  address: 'localhost:50000',
  client_id: 'order-receiver'
)

messages = client.receive_queue_messages(
  channel: 'orders',
  max_messages: 1,
  wait_timeout_seconds: 10
)
messages.each do |msg|
  puts "Received: #{msg.body}"
  puts "Metadata: #{msg.metadata}"
end

client.close
receiver.exs
{:ok, client} =
  KubeMQ.Client.start_link(address: "localhost:50000", client_id: "order-receiver")

{:ok, result} =
  KubeMQ.Client.receive_queue_messages(client, "orders",
    max_messages: 1,
    wait_timeout: 10_000
  )

Enum.each(result.messages, fn msg ->
  IO.puts("Received: #{msg.body}")
  IO.puts("Metadata: #{msg.metadata}")
end)

KubeMQ.Client.close(client)

Verify

Run the sender first, then the receiver. You should see:

Sent: id=<message-id>
Received: {"orderId":"ORD-1234","total":99.99}
Metadata: order.created
Message acknowledged

What Just Happened

Send → durable queue → poll & deliver to one receiver → acknowledge to remove.

The full round-trip, step by step:

  1. The sender connected to KubeMQ and placed an order message on the orders queue
  2. The message was stored durably — it persists even if no receiver is connected yet
  3. The receiver polled the queue and received the message (hidden from other consumers)
  4. After processing, the receiver acknowledged the message, permanently removing it from the queue

Unlike Events, queue messages persist until acknowledged. If the receiver crashes before acknowledging, the message becomes available again after the visibility timeout expires.

Next Steps

Was this page helpful?

On this page