# Consumer Groups (/learn/events/tutorials/consumer-groups)



## What You Will Build [#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.

<Mermaid
  chart="graph LR
  P[&#x22;Publisher&#x22;]
  CH{{&#x22;Events channel<br/>order-events&#x22;}}
  W1[&#x22;Worker A<br/>group: workers&#x22;]
  W2[&#x22;Worker B<br/>group: workers&#x22;]
  W3[&#x22;Worker C<br/>group: workers&#x22;]

  P -- publish --> CH
  CH -- &#x22;1 of 3&#x22; --> W1
  CH -- &#x22;1 of 3&#x22; --> W2
  CH -- &#x22;1 of 3&#x22; --> W3

  class CH events
  class P,W1,W2,W3 queue"
/>

*A consumer group load-balances each event to exactly one member — one channel, work split across the group.*

## Prerequisites [#prerequisites]

* KubeMQ server running on `localhost:50000`
* SDK installed ([Getting Started](/learn/events/getting-started))

## Steps [#steps]

<Steps>
  <Step>
    ### Create the Publisher [#create-the-publisher]

    Send a batch of order events to a channel.

    <Tabs groupId="language" items="['Go', 'Python', 'Node.js', 'Java', 'C#', 'Kotlin', 'C++', 'Rust', 'Ruby', 'Elixir']">
      <Tab value="Go">
        ```go title="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)
            }
        }
        ```
      </Tab>

      <Tab value="Python">
        ```python title="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()
        ```
      </Tab>

      <Tab value="Node.js">
        ```javascript title="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));
        }
        ```
      </Tab>

      <Tab value="Java">
        ```java title="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();
        ```
      </Tab>

      <Tab value="C#">
        ```csharp title="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);
        }
        ```
      </Tab>

      <Tab value="Kotlin">
        ```kotlin title="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()
        ```
      </Tab>

      <Tab value="C++">
        ```cpp title="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));
        }
        ```
      </Tab>

      <Tab value="Rust">
        ```rust title="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(())
        }
        ```
      </Tab>

      <Tab value="Ruby">
        ```ruby title="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
        ```
      </Tab>

      <Tab value="Elixir">
        ```elixir title="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)
        ```
      </Tab>
    </Tabs>
  </Step>

  <Step>
    ### Create a Consumer Group [#create-a-consumer-group]

    Three subscribers join the same group. KubeMQ distributes events across the group in round-robin fashion.

    <Tabs groupId="language" items="['Go', 'Python', 'Node.js', 'Java', 'C#', 'Kotlin', 'C++', 'Rust', 'Ruby', 'Elixir']">
      <Tab value="Go">
        ```go title="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()
        }
        ```
      </Tab>

      <Tab value="Python">
        ```python title="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()
        ```
      </Tab>

      <Tab value="Node.js">
        ```javascript title="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'`);
        ```
      </Tab>

      <Tab value="Java">
        ```java title="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();
        ```
      </Tab>

      <Tab value="C#">
        ```csharp title="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)}");
        }
        ```
      </Tab>

      <Tab value="Kotlin">
        ```kotlin title="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()
        ```
      </Tab>

      <Tab value="C++">
        ```cpp title="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));
        ```
      </Tab>

      <Tab value="Rust">
        ```rust title="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(())
        }
        ```
      </Tab>

      <Tab value="Ruby">
        ```ruby title="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
        ```
      </Tab>

      <Tab value="Elixir">
        ```elixir title="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)
        ```
      </Tab>
    </Tabs>

    Run three instances with different `WORKER_ID` values:

    ```bash
    WORKER_ID=worker-A ./grouped_worker &
    WORKER_ID=worker-B ./grouped_worker &
    WORKER_ID=worker-C ./grouped_worker &
    ```
  </Step>

  <Step>
    ### Verify Load Distribution [#verify-load-distribution]

    Publish 9 events and observe each worker receives approximately 3 events:

    **Worker A** (receives \~3 events):

    ```text
    [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):

    ```text
    [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):

    ```text
    [worker-C] Processing: {"orderId":"ORD-003","status":"created"}
    [worker-C] Processing: {"orderId":"ORD-006","status":"created"}
    [worker-C] Processing: {"orderId":"ORD-009","status":"created"}
    ```
  </Step>
</Steps>

## Fan-Out vs Consumer Groups [#fan-out-vs-consumer-groups]

<Mermaid
  chart="graph TB
  subgraph FanOut[&#x22;Fan-Out (no group)&#x22;]
    P1[&#x22;Publisher&#x22;]
    EV1{{&#x22;Events channel&#x22;}}
    S1A[&#x22;Subscriber A<br/>receives ALL&#x22;]
    S1B[&#x22;Subscriber B<br/>receives ALL&#x22;]
    S1C[&#x22;Subscriber C<br/>receives ALL&#x22;]
    P1 -- publish --> EV1
    EV1 -- &#x22;every event&#x22; --> S1A
    EV1 -- &#x22;every event&#x22; --> S1B
    EV1 -- &#x22;every event&#x22; --> S1C
  end
  subgraph Group[&#x22;Consumer Group&#x22;]
    P2[&#x22;Publisher&#x22;]
    EV2{{&#x22;Events channel&#x22;}}
    S2A[&#x22;Worker A&#x22;]
    S2B[&#x22;Worker B&#x22;]
    S2C[&#x22;Worker C&#x22;]
    P2 -- publish --> EV2
    EV2 -- &#x22;1 of 3&#x22; --> S2A
    EV2 -- &#x22;1 of 3&#x22; --> S2B
    EV2 -- &#x22;1 of 3&#x22; --> S2C
  end

  class EV1 events
  class P1,S1A,S1B,S1C client
  class EV2 events
  class P2,S2A,S2B,S2C queue"
/>

*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            |

<Callout type="info">
  You can combine both patterns: grouped workers for load-balanced processing and an ungrouped monitor that sees all events. See [Scale Subscribers](/learn/events/how-to/scale-subscribers) for this pattern.
</Callout>

## Next Steps [#next-steps]

<Cards>
  <Card title="Scale Subscribers" href="/learn/events/how-to/scale-subscribers" description="Horizontal scaling strategies with groups." />

  <Card title="Multicast Events" href="/learn/events/tutorials/multicast" description="Publish to multiple channels simultaneously." />
</Cards>
