KubeMQ
LearnQueuesTutorials

Batch Operations

Send and receive multiple queue messages in a single operation for higher throughput.

What You Will Build

A batch sender that publishes multiple order messages in one call, and a batch receiver that pulls and acknowledges all of them at once.

One request carries the whole batch in each direction — fewer round-trips, higher throughput.

Prerequisites

Steps

Batch Send

Send multiple messages in a single request. Each message can have its own body, metadata, tags, and policy.

batch_sender.go
orders := []struct {
    ID    string
    Total float64
}{
    {"ORD-001", 29.99},
    {"ORD-002", 149.50},
    {"ORD-003", 75.00},
    {"ORD-004", 210.00},
    {"ORD-005", 15.99},
}

var messages []*kubemq.QueueMessage
for _, order := range orders {
    body := fmt.Sprintf(`{"orderId":"%s","total":%.2f}`, order.ID, order.Total)
    messages = append(messages, kubemq.NewQueueMessage().
        SetChannel("orders.batch").
        SetBody([]byte(body)).
        SetMetadata("order.created"),
    )
}

results, err := client.SendQueueMessages(ctx, messages)
if err != nil {
    log.Fatal(err)
}
for _, r := range results {
    fmt.Printf("Sent: id=%s, error=%v\n", r.MessageID, r.IsError)
}
batch_sender.py
import json

orders = [
    {"orderId": "ORD-001", "total": 29.99},
    {"orderId": "ORD-002", "total": 149.50},
    {"orderId": "ORD-003", "total": 75.00},
    {"orderId": "ORD-004", "total": 210.00},
    {"orderId": "ORD-005", "total": 15.99},
]

messages = [
    QueueMessage(
        channel="orders.batch",
        body=json.dumps(order).encode(),
        metadata="order.created",
    )
    for order in orders
]

results = client.send_queue_messages(messages)
for r in results:
    print(f"Sent: id={r.id}, error={r.is_error}")
batch_sender.ts
const orders = [
  { orderId: 'ORD-001', total: 29.99 },
  { orderId: 'ORD-002', total: 149.5 },
  { orderId: 'ORD-003', total: 75.0 },
  { orderId: 'ORD-004', total: 210.0 },
  { orderId: 'ORD-005', total: 15.99 },
];

const messages = orders.map((order) =>
  createQueueMessage({
    channel: 'orders.batch',
    body: JSON.stringify(order),
    metadata: 'order.created',
  }),
);

const results = await client.sendQueueMessagesBatch(messages);
for (const r of results) {
  console.log(`Sent: id=${r.messageId}, error=${r.isError}`);
}
BatchSender.java
List<QueueMessage> messages = List.of(
    QueueMessage.builder().channel("orders.batch")
        .body("{\"orderId\":\"ORD-001\",\"total\":29.99}".getBytes())
        .metadata("order.created").build(),
    QueueMessage.builder().channel("orders.batch")
        .body("{\"orderId\":\"ORD-002\",\"total\":149.50}".getBytes())
        .metadata("order.created").build(),
    QueueMessage.builder().channel("orders.batch")
        .body("{\"orderId\":\"ORD-003\",\"total\":75.00}".getBytes())
        .metadata("order.created").build()
);

List<SendQueueMessageResult> results = client.sendQueueMessages(messages);
for (SendQueueMessageResult r : results) {
    System.out.printf("Sent: id=%s, error=%b%n", r.getMessageId(), r.isError());
}
BatchSender.cs
var messages = new[]
{
    new QueueMessage { Channel = "orders.batch",
        Body = Encoding.UTF8.GetBytes("{\"orderId\":\"ORD-001\",\"total\":29.99}"),
        Metadata = "order.created" },
    new QueueMessage { Channel = "orders.batch",
        Body = Encoding.UTF8.GetBytes("{\"orderId\":\"ORD-002\",\"total\":149.50}"),
        Metadata = "order.created" },
    new QueueMessage { Channel = "orders.batch",
        Body = Encoding.UTF8.GetBytes("{\"orderId\":\"ORD-003\",\"total\":75.00}"),
        Metadata = "order.created" },
};

var results = await client.SendQueueMessagesBatchAsync(messages);
foreach (var r in results)
{
    Console.WriteLine($"Sent: id={r.MessageId}, error={r.IsError}");
}
BatchSender.kt
val messages = listOf(
    QueueMessage(channel = "orders.batch",
        body = """{"orderId":"ORD-001","total":29.99}""".toByteArray(),
        metadata = "order.created"),
    QueueMessage(channel = "orders.batch",
        body = """{"orderId":"ORD-002","total":149.50}""".toByteArray(),
        metadata = "order.created"),
    QueueMessage(channel = "orders.batch",
        body = """{"orderId":"ORD-003","total":75.00}""".toByteArray(),
        metadata = "order.created"),
)

val results = client.sendQueueMessages(messages)
for (r in results) {
    println("Sent: id=${r.messageId}, error=${r.isError}")
}
batch_sender.cpp
std::vector<kubemq::QueueMessage> messages;
for (const auto& [id, total] : std::vector<std::pair<std::string, double>>{
    {"ORD-001", 29.99}, {"ORD-002", 149.50}, {"ORD-003", 75.00}}) {
    kubemq::QueueMessage msg;
    msg.channel = "orders.batch";
    msg.body = "{\"orderId\":\"" + id + "\",\"total\":" + std::to_string(total) + "}";
    msg.metadata = "order.created";
    messages.push_back(msg);
}

auto results = client.sendQueueMessages(messages);
for (const auto& r : results) {
    std::cout << "Sent: id=" << r.messageId << std::endl;
}
batch_sender.rs
let channel = "orders.batch";

let orders = [
    ("ORD-001", 29.99),
    ("ORD-002", 149.50),
    ("ORD-003", 75.00),
    ("ORD-004", 210.00),
    ("ORD-005", 15.99),
];

let messages: Vec<QueueMessage> = orders
    .iter()
    .map(|(id, total)| {
        let body = format!(r#"{{"orderId":"{}","total":{}}}"#, id, total);
        QueueMessageBuilder::new()
            .channel(channel)
            .body(body.into_bytes())
            .metadata("order.created")
            .build()
    })
    .collect();

let results = client.send_queue_messages(messages).await?;
for r in &results {
    println!("Sent: id={}, error={}", r.message_id, r.is_error);
}
batch_sender.rb
channel = 'orders.batch'

orders = [
  { orderId: 'ORD-001', total: 29.99 },
  { orderId: 'ORD-002', total: 149.50 },
  { orderId: 'ORD-003', total: 75.00 },
  { orderId: 'ORD-004', total: 210.00 },
  { orderId: 'ORD-005', total: 15.99 },
]

messages = orders.map do |order|
  KubeMQ::Queues::QueueMessage.new(
    channel: channel,
    metadata: 'order.created',
    body: order.to_json,
  )
end

results = client.send_queue_messages_batch(messages)
results.each do |r|
  puts "Sent: id=#{r.id}, error?=#{r.error?}"
end
batch_sender.exs
channel = "orders.batch"

orders = [
  %{orderId: "ORD-001", total: 29.99},
  %{orderId: "ORD-002", total: 149.50},
  %{orderId: "ORD-003", total: 75.00},
  %{orderId: "ORD-004", total: 210.00},
  %{orderId: "ORD-005", total: 15.99}
]

messages =
  for order <- orders do
    KubeMQ.QueueMessage.new(
      channel: channel,
      body: Jason.encode!(order),
      metadata: "order.created"
    )
  end

{:ok, result} = KubeMQ.Client.send_queue_messages(client, messages)
Enum.each(result.results, fn r ->
  IO.puts("Sent: id=#{r.message_id}, error=#{r.is_error}")
end)

Batch Receive

Receive multiple messages in one call. The maxMessages parameter controls how many messages to fetch.

batch_receiver.go
resp, err := client.PollQueue(ctx, &kubemq.PollRequest{
    Channel:            "orders.batch",
    MaxItems:           10,
    WaitTimeoutSeconds: 5,
    AutoAck:            false,
})
if err != nil {
    log.Fatal(err)
}
fmt.Printf("Received %d messages\n", len(resp.Messages))
for _, m := range resp.Messages {
    fmt.Printf("  %s: %s\n", m.Message.MessageID, string(m.Message.Body))
}
if err := resp.AckAll(); err != nil {
    log.Fatal(err)
}
fmt.Println("All messages acknowledged")
batch_receiver.py
response = client.receive_queue_messages(
    channel="orders.batch",
    max_messages=10,
    wait_timeout_in_seconds=5,
)
print(f"Received {len(response.messages)} messages")
for msg in response.messages:
    print(f"  {msg.id}: {msg.body.decode('utf-8')}")
    msg.ack()
print("All messages acknowledged")
batch_receiver.ts
const messages = await client.receiveQueueMessages({
  channel: 'orders.batch',
  maxMessages: 10,
  waitTimeoutSeconds: 5,
});
console.log(`Received ${messages.length} messages`);
for (const msg of messages) {
  console.log(`  ${msg.messageId}: ${new TextDecoder().decode(msg.body)}`);
  await msg.ack();
}
console.log('All messages acknowledged');
BatchReceiver.java
ReceiveQueueMessagesResponse response = client.receiveQueueMessages(
    ReceiveQueueMessagesRequest.builder()
        .channel("orders.batch")
        .maxMessages(10)
        .waitTimeoutSeconds(5)
        .build());

System.out.printf("Received %d messages%n", response.getMessages().size());
for (QueueMessageReceived msg : response.getMessages()) {
    System.out.printf("  %s: %s%n", msg.getMessageId(), new String(msg.getBody()));
    msg.ack();
}
System.out.println("All messages acknowledged");
BatchReceiver.cs
var response = await client.ReceiveQueueMessagesAsync(new ReceiveQueueMessagesRequest
{
    Channel = "orders.batch",
    MaxMessages = 10,
    WaitTimeoutSeconds = 5,
});
Console.WriteLine($"Received {response.Messages.Count} messages");
foreach (var msg in response.Messages)
{
    Console.WriteLine($"  {msg.MessageId}: {Encoding.UTF8.GetString(msg.Body.Span)}");
    await msg.AckAsync();
}
Console.WriteLine("All messages acknowledged");
BatchReceiver.kt
val response = client.receiveQueueMessages(
    channel = "orders.batch",
    maxMessages = 10,
    waitTimeoutSeconds = 5
)
println("Received ${response.messages.size} messages")
for (msg in response.messages) {
    println("  ${msg.messageId}: ${String(msg.body)}")
    msg.ack()
}
println("All messages acknowledged")
batch_receiver.cpp
auto response = client.receiveQueueMessages("orders.batch", 10, 5);
std::cout << "Received " << response.messages.size() << " messages" << std::endl;
for (const auto& msg : response.messages) {
    std::cout << "  " << msg.messageId << ": " << msg.body << std::endl;
    msg.ack();
}
std::cout << "All messages acknowledged" << std::endl;
batch_receiver.rs
// The simple queues client acks on receive; the fourth arg is auto-requeue.
let messages = client
    .receive_queue_messages("orders.batch", 10, 5, false)
    .await?;
println!("Received {} messages", messages.len());
for m in &messages {
    println!("  {}: {}", m.id, String::from_utf8_lossy(&m.body));
}
println!("All messages acknowledged");
batch_receiver.rb
# The simple queues client acks on receive; clear any leftover with ack_all.
messages = client.receive_queue_messages(
  channel: 'orders.batch',
  max_messages: 10,
  wait_timeout_seconds: 5,
)
puts "Received #{messages.size} messages"
messages.each do |m|
  puts "  #{m.id}: #{m.body}"
end
client.ack_all_queue_messages(channel: 'orders.batch', wait_timeout_seconds: 5)
puts 'All messages acknowledged'
batch_receiver.exs
# The simple queues client acks on receive; clear any leftover with ack_all.
{:ok, result} =
  KubeMQ.Client.receive_queue_messages(client, "orders.batch",
    max_messages: 10,
    wait_timeout: 5_000
  )

IO.puts("Received #{result.messages_received} messages")
Enum.each(result.messages, fn m ->
  IO.puts("  #{m.message_id}: #{m.body}")
end)

KubeMQ.Client.ack_all_queue_messages(client, "orders.batch", wait_timeout: 5_000)
IO.puts("All messages acknowledged")

Performance: Single vs Batch

ApproachThroughputNetwork CallsUse Case
Single sendLower1 per messageReal-time, low volume
Batch sendHigher1 per batchBulk import, high volume
Single receiveLower1 per pollInteractive processing
Batch receiveHigher1 per pollBackground workers

The maximum batch size is controlled by the server setting MaxNumberOfMessages (default: 1,024 messages per request).

Next Steps

Was this page helpful?

On this page