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/v2pip install kubemqnpm 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.1implementation("io.kubemq.sdk:kubemq-sdk-kotlin:1.0.1")vcpkg install kubemq[dependencies]
kubemq = "1.0.1"
tokio = { version = "1", features = ["full"] }gem install kubemqdef deps do
[{:kubemq, "~> 1.0"}]
endCreate the Sender
Connect to KubeMQ and send an order message to the orders queue.
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)
}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()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();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();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}");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()
}#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;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(())
}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{: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.
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")
}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()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();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();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");
}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()
}#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;
}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(())
}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{: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 acknowledgedWhat Just Happened
Send → durable queue → poll & deliver to one receiver → acknowledge to remove.
The full round-trip, step by step:
- The sender connected to KubeMQ and placed an order message on the
ordersqueue - The message was stored durably — it persists even if no receiver is connected yet
- The receiver polled the queue and received the message (hidden from other consumers)
- 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?