Command Group
Load-balanced command handling via consumer groups
Overview
A command consumer group turns a single command handler into a scalable worker pool: run multiple identical instances subscribed with the same group name, and the broker load-balances each incoming command to exactly one member instead of broadcasting it to all of them. This is how you add capacity to handle a growing command volume — start more workers in the same group — without changing anything on the caller's side.
Every subscriber passes the same group alongside channel to subscribe_to_commands; the broker tracks membership and picks one live member per command. send_command on the caller side is unaware groups exist — it just awaits a response, which comes back from whichever handler happened to process it via send_command_response.
Gotchas: group membership is scoped per channel — subscribers on the same channel with different group names each get their own full copy of every command (fan-out), which looks like a bug when you expected load-balancing. A slow handler still holds up the caller's timeout, since only one worker is ever picked. And if every member of the group is offline when a command arrives, the send simply fails or times out — commands aren't queued or replayed for a group that has no active listener.
Prerequisites
- KubeMQ server running on
localhost:50000 - Rust SDK installed (
cargo add kubemq)
Code
use kubemq::prelude::*;
use kubemq::{CommandBuilder, CommandReplyBuilder};
use std::time::Duration;
#[tokio::main]
async fn main() -> kubemq::Result<()> {
let client = KubemqClient::builder()
.host("localhost")
.port(50000)
.build()
.await?;
let channel = "rust-rpc.command-group";
let group = "cmd-handler-group";
let rc1 = client.clone();
let sub1 = client.subscribe_to_commands(channel, group,
move |cmd| { let c = rc1.clone(); Box::pin(async move {
println!("Handler-1: id={}", cmd.id);
let reply = CommandReplyBuilder::new().request_id(&cmd.id).response_to(&cmd.response_to).build();
let _ = c.send_command_response(reply).await;
})}, None,
).await?;
let rc2 = client.clone();
let sub2 = client.subscribe_to_commands(channel, group,
move |cmd| { let c = rc2.clone(); Box::pin(async move {
println!("Handler-2: id={}", cmd.id);
let reply = CommandReplyBuilder::new().request_id(&cmd.id).response_to(&cmd.response_to).build();
let _ = c.send_command_response(reply).await;
})}, None,
).await?;
tokio::time::sleep(Duration::from_millis(500)).await;
for i in 0..5 {
let cmd = CommandBuilder::new()
.channel(channel)
.body(format!("cmd-{}", i).into_bytes())
.timeout(Duration::from_secs(10))
.build();
let resp = client.send_command(cmd).await?;
println!("Response {}: executed={}", i, resp.executed);
}
sub1.unsubscribe().await;
sub2.unsubscribe().await;
client.close().await?;
Ok(())
}How It Works
- Both handlers join the same
group, so each command is delivered to exactly one handler. - The server distributes commands across group members for load balancing.
- Review timeouts, channel names, and client IDs before running against shared environments.
- Run the program while the server from the prerequisites is available.
Related
Was this page helpful?