KubeMQ
LearnEventsTutorials

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

Steps

Create the Publisher

Send a batch of order events to a channel.

order_publisher.go
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)
    }
}
order_publisher.py
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()
order_publisher.js
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));
}
OrderPublisher.java
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();
OrderPublisher.cs
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);
}
OrderPublisher.kt
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()
order_publisher.cpp
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));
}
order_publisher.rs
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(())
}
order_publisher.rb
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
order_publisher.exs
{: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.

grouped_worker.go
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()
}
grouped_worker.py
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()
grouped_worker.js
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'`);
GroupedWorker.java
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();
GroupedWorker.cs
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)}");
}
GroupedWorker.kt
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()
grouped_worker.cpp
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));
grouped_worker.rs
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(())
}
grouped_worker.rb
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.close
grouped_worker.exs
worker_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.

BehaviorFan-Out (no group)Consumer Group
DeliveryEvery subscriber receives every eventEach event goes to exactly one group member
Use caseMultiple independent consumersLoad-balanced processing
Group parameterEmpty string ""Same group name (e.g., "workers")
Scaling effectMore subscribers = more total processingMore 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?

On this page