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
- KubeMQ server running on
localhost:50000 - SDK installed (Getting Started)
Steps
Batch Send
Send multiple messages in a single request. Each message can have its own body, metadata, tags, and policy.
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)
}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}")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}`);
}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());
}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}");
}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}")
}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;
}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);
}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?}"
endchannel = "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.
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")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")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');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");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");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")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;// 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");# 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'# 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
| Approach | Throughput | Network Calls | Use Case |
|---|---|---|---|
| Single send | Lower | 1 per message | Real-time, low volume |
| Batch send | Higher | 1 per batch | Bulk import, high volume |
| Single receive | Lower | 1 per poll | Interactive processing |
| Batch receive | Higher | 1 per poll | Background workers |
The maximum batch size is controlled by the server setting MaxNumberOfMessages (default: 1,024 messages per request).
Next Steps
Was this page helpful?