Consumer Groups
Distribute event processing across multiple consumers with load-balanced groups.
What You Will Build
A publisher sending order events and three consumers in a shared group, where each event is delivered to exactly one consumer (round-robin). You will then compare this with standard fan-out behavior.
A consumer group load-balances each event to exactly one member — one channel, work split across the group.
Prerequisites
- KubeMQ server running on
localhost:50000 - SDK installed (Getting Started)
Steps
Create the Publisher
Send a batch of order events to a channel.
package main
import (
"context"
"fmt"
"log"
"time"
"github.com/kubemq-io/kubemq-go/v2"
)
func main() {
ctx := context.Background()
client, err := kubemq.NewClient(ctx,
kubemq.WithAddress("localhost", 50000),
)
if err != nil {
log.Fatal(err)
}
defer client.Close()
for i := 1; i <= 9; i++ {
body := fmt.Sprintf(`{"orderId":"ORD-%03d","status":"created"}`, i)
err = client.SendEvent(ctx, kubemq.NewEvent().
SetChannel("order-events").
SetBody([]byte(body)),
)
if err != nil {
log.Printf("Failed to send ORD-%03d: %v", i, err)
continue
}
log.Printf("Published ORD-%03d", i)
time.Sleep(200 * time.Millisecond)
}
}import time
from kubemq.pubsub import Client as PubSubClient
from kubemq.pubsub import EventMessage
client = PubSubClient(address="localhost:50000")
for i in range(1, 10):
body = f'{{"orderId":"ORD-{i:03d}","status":"created"}}'
client.send_event(
EventMessage(channel="order-events", body=body.encode("utf-8"))
)
print(f"Published ORD-{i:03d}")
time.sleep(0.2)
client.close()const { KubeMQClient } = require("kubemq-js");
const client = new KubeMQClient({ address: "localhost:50000" });
for (let i = 1; i <= 9; i++) {
const orderId = `ORD-${String(i).padStart(3, "0")}`;
await client.sendEvent({
channel: "order-events",
body: Buffer.from(JSON.stringify({ orderId, status: "created" })),
});
console.log(`Published ${orderId}`);
await new Promise((r) => setTimeout(r, 200));
}PubSubClient client = PubSubClient.builder()
.address("localhost:50000")
.clientId("order-publisher")
.build();
for (int i = 1; i <= 9; i++) {
String body = String.format(
"{\"orderId\":\"ORD-%03d\",\"status\":\"created\"}", i);
client.sendEventsMessage(EventMessage.builder()
.channel("order-events")
.body(body.getBytes())
.build());
System.out.printf("Published ORD-%03d%n", i);
Thread.sleep(200);
}
client.close();await using var client = new KubeMQClient(new KubeMQClientOptions());
await client.ConnectAsync();
for (var i = 1; i <= 9; i++)
{
var body = $"{{\"orderId\":\"ORD-{i:D3}\",\"status\":\"created\"}}";
await client.SendEventAsync(new EventMessage
{
Channel = "order-events",
Body = Encoding.UTF8.GetBytes(body),
});
Console.WriteLine($"Published ORD-{i:D3}");
await Task.Delay(200);
}val client = PubSubClient("localhost:50000")
for (i in 1..9) {
val body = """{"orderId":"ORD-${"%03d".format(i)}","status":"created"}"""
client.sendEvent(EventMessage(
channel = "order-events",
body = body.toByteArray(),
))
println("Published ORD-${"%03d".format(i)}")
Thread.sleep(200)
}
client.close()auto client = kubemq::PubSubClient("localhost:50000");
for (int i = 1; i <= 9; i++) {
kubemq::EventMessage event;
event.channel = "order-events";
event.body = "{\"orderId\":\"ORD-" + std::to_string(i) +
"\",\"status\":\"created\"}";
client.sendEvent(event);
std::cout << "Published ORD-" << i << std::endl;
std::this_thread::sleep_for(std::chrono::milliseconds(200));
}use kubemq::prelude::*;
use kubemq::EventBuilder;
use std::time::Duration;
#[tokio::main]
async fn main() -> kubemq::Result<()> {
let client = KubemqClient::builder()
.host("localhost")
.port(50000)
.build()
.await?;
for i in 1..=9 {
let body = format!(r#"{{"orderId":"ORD-{:03}","status":"created"}}"#, i);
let event = EventBuilder::new()
.channel("order-events")
.body(body.into_bytes())
.build();
client.send_event(event).await?;
println!("Published ORD-{:03}", i);
tokio::time::sleep(Duration::from_millis(200)).await;
}
client.close().await?;
Ok(())
}require 'kubemq'
client = KubeMQ::PubSubClient.new(address: "localhost:50000", client_id: "order-publisher")
(1..9).each do |i|
body = format('{"orderId":"ORD-%03d","status":"created"}', i)
client.send_event(KubeMQ::PubSub::EventMessage.new(
channel: "order-events",
body: body
))
puts format("Published ORD-%03d", i)
sleep 0.2
end
client.close{:ok, client} = KubeMQ.Client.start_link(address: "localhost:50000", client_id: "order-publisher")
for i <- 1..9 do
id = String.pad_leading(Integer.to_string(i), 3, "0")
event = KubeMQ.Event.new(
channel: "order-events",
body: ~s({"orderId":"ORD-#{id}","status":"created"})
)
:ok = KubeMQ.Client.send_event(client, event)
IO.puts("Published ORD-#{id}")
Process.sleep(200)
end
KubeMQ.Client.close(client)Create a Consumer Group
Three subscribers join the same group. KubeMQ distributes events across the group in round-robin fashion.
package main
import (
"context"
"fmt"
"log"
"os"
"github.com/kubemq-io/kubemq-go/v2"
)
func main() {
workerID := os.Getenv("WORKER_ID")
if workerID == "" {
workerID = "worker-1"
}
ctx := context.Background()
client, err := kubemq.NewClient(ctx,
kubemq.WithAddress("localhost", 50000),
)
if err != nil {
log.Fatal(err)
}
defer client.Close()
sub, err := client.SubscribeToEvents(ctx, "order-events", "workers",
kubemq.WithOnEvent(func(event *kubemq.Event) {
fmt.Printf("[%s] Processing: %s\n", workerID,
string(event.Body))
}),
kubemq.WithOnError(func(err error) {
log.Printf("[%s] Error: %v", workerID, err)
}),
)
if err != nil {
log.Fatal(err)
}
defer sub.Unsubscribe()
log.Printf("[%s] Ready in group 'workers'", workerID)
<-ctx.Done()
}import os
import time
from kubemq.pubsub import Client as PubSubClient
from kubemq.pubsub import EventsSubscription, CancellationToken
worker_id = os.environ.get("WORKER_ID", "worker-1")
def on_event(event):
print(f"[{worker_id}] Processing: {event.body.decode('utf-8')}")
client = PubSubClient(address="localhost:50000")
client.subscribe_to_events(
subscription=EventsSubscription(
channel="order-events",
group="workers",
on_receive_event_callback=on_event,
on_error_callback=lambda e: print(f"[{worker_id}] Error: {e}"),
),
cancel=CancellationToken(),
)
print(f"[{worker_id}] Ready in group 'workers'")
time.sleep(300)
client.close()const { KubeMQClient } = require("kubemq-js");
const workerId = process.env.WORKER_ID ?? "worker-1";
const client = new KubeMQClient({ address: "localhost:50000" });
client.subscribeToEvents({
channel: "order-events",
group: "workers",
onEvent: (msg) =>
console.log(
`[${workerId}] Processing: ${Buffer.from(msg.body).toString()}`
),
onError: (err) =>
console.error(`[${workerId}] Error:`, err.message),
});
console.log(`[${workerId}] Ready in group 'workers'`);String workerId = System.getenv().getOrDefault("WORKER_ID", "worker-1");
PubSubClient client = PubSubClient.builder()
.address("localhost:50000")
.clientId(workerId)
.build();
client.subscribeToEvents(EventsSubscription.builder()
.channel("order-events")
.group("workers")
.onReceiveEventCallback(event ->
System.out.printf("[%s] Processing: %s%n", workerId,
new String(event.getBody())))
.onErrorCallback(err ->
System.err.printf("[%s] Error: %s%n", workerId, err.getMessage()))
.build());
System.out.printf("[%s] Ready in group 'workers'%n", workerId);
Thread.sleep(300_000);
client.close();var workerId = Environment.GetEnvironmentVariable("WORKER_ID") ?? "worker-1";
await using var client = new KubeMQClient(new KubeMQClientOptions());
await client.ConnectAsync();
Console.WriteLine($"[{workerId}] Ready in group 'workers'");
await foreach (var msg in client.SubscribeToEventsAsync(
new EventsSubscription { Channel = "order-events", Group = "workers" }))
{
Console.WriteLine($"[{workerId}] Processing: "
+ $"{Encoding.UTF8.GetString(msg.Body.Span)}");
}val workerId = System.getenv("WORKER_ID") ?: "worker-1"
val client = PubSubClient("localhost:50000")
client.subscribeToEvents(
channel = "order-events",
group = "workers",
onEvent = { event ->
println("[$workerId] Processing: ${String(event.body)}")
},
onError = { err ->
System.err.println("[$workerId] Error: ${err.message}")
}
)
println("[$workerId] Ready in group 'workers'")
Thread.sleep(300_000)
client.close()auto workerId = std::getenv("WORKER_ID") ?
std::string(std::getenv("WORKER_ID")) : std::string("worker-1");
auto client = kubemq::PubSubClient("localhost:50000");
client.subscribeToEvents("order-events", "workers",
[&workerId](const kubemq::Event& event) {
std::cout << "[" << workerId << "] Processing: "
<< event.body << std::endl;
},
[&workerId](const std::string& err) {
std::cerr << "[" << workerId << "] Error: " << err << std::endl;
}
);
std::cout << "[" << workerId << "] Ready in group 'workers'" << std::endl;
std::this_thread::sleep_for(std::chrono::seconds(300));use kubemq::prelude::*;
use std::time::Duration;
#[tokio::main]
async fn main() -> kubemq::Result<()> {
let worker_id = std::env::var("WORKER_ID").unwrap_or_else(|_| "worker-1".to_string());
let client = KubemqClient::builder()
.host("localhost")
.port(50000)
.build()
.await?;
// Subscribe with a non-empty group -- each event goes to only one member
let id = worker_id.clone();
let sub = client
.subscribe_to_events(
"order-events",
"workers",
move |event| {
let id = id.clone();
Box::pin(async move {
println!(
"[{}] Processing: {}",
id,
String::from_utf8_lossy(&event.body)
);
})
},
None,
)
.await?;
println!("[{}] Ready in group 'workers'", worker_id);
tokio::time::sleep(Duration::from_secs(300)).await;
sub.unsubscribe().await;
client.close().await?;
Ok(())
}require 'kubemq'
worker_id = ENV.fetch('WORKER_ID', 'worker-1')
client = KubeMQ::PubSubClient.new(address: "localhost:50000", client_id: worker_id)
cancel = KubeMQ::CancellationToken.new
# A non-empty group makes subscribers compete -- each event goes to one member
sub = KubeMQ::PubSub::EventsSubscription.new(channel: "order-events", group: "workers")
client.subscribe_to_events(sub, cancellation_token: cancel, on_error: lambda { |e|
puts "[#{worker_id}] Error: #{e.message}"
}) do |event|
puts "[#{worker_id}] Processing: #{event.body}"
end
puts "[#{worker_id}] Ready in group 'workers'"
sleep 300
cancel.cancel
client.closeworker_id = System.get_env("WORKER_ID", "worker-1")
{:ok, client} = KubeMQ.Client.start_link(address: "localhost:50000", client_id: worker_id)
# A non-empty group makes subscribers compete -- each event goes to one member
{:ok, _sub} =
KubeMQ.Client.subscribe_to_events(client, "order-events",
group: "workers",
on_event: fn event ->
IO.puts("[#{worker_id}] Processing: #{event.body}")
end
)
IO.puts("[#{worker_id}] Ready in group 'workers'")
Process.sleep(300_000)
KubeMQ.Client.close(client)Run three instances with different WORKER_ID values:
WORKER_ID=worker-A ./grouped_worker &
WORKER_ID=worker-B ./grouped_worker &
WORKER_ID=worker-C ./grouped_worker &Verify Load Distribution
Publish 9 events and observe each worker receives approximately 3 events:
Worker A (receives ~3 events):
[worker-A] Processing: {"orderId":"ORD-001","status":"created"}
[worker-A] Processing: {"orderId":"ORD-004","status":"created"}
[worker-A] Processing: {"orderId":"ORD-007","status":"created"}Worker B (receives ~3 events):
[worker-B] Processing: {"orderId":"ORD-002","status":"created"}
[worker-B] Processing: {"orderId":"ORD-005","status":"created"}
[worker-B] Processing: {"orderId":"ORD-008","status":"created"}Worker C (receives ~3 events):
[worker-C] Processing: {"orderId":"ORD-003","status":"created"}
[worker-C] Processing: {"orderId":"ORD-006","status":"created"}
[worker-C] Processing: {"orderId":"ORD-009","status":"created"}Fan-Out vs Consumer Groups
Same channel, two behaviors: fan-out delivers every event to every subscriber; a consumer group splits the load so each event lands on exactly one member.
| Behavior | Fan-Out (no group) | Consumer Group |
|---|---|---|
| Delivery | Every subscriber receives every event | Each event goes to exactly one group member |
| Use case | Multiple independent consumers | Load-balanced processing |
| Group parameter | Empty string "" | Same group name (e.g., "workers") |
| Scaling effect | More subscribers = more total processing | More members = higher throughput |
You can combine both patterns: grouped workers for load-balanced processing and an ungrouped monitor that sees all events. See Scale Subscribers for this pattern.
Next Steps
Was this page helpful?