Consumer Group
Load-balance KubeMQ events across a consumer group of subscribers in Node.js so each event is handled by exactly one member.
Overview
A consumer group turns Events pub/sub from a broadcast into a work queue. By default every subscriber on a channel gets every event — fine for notifications, but wasteful when you want a pool of workers to split a stream of tasks so each one is handled exactly once. Reach for a consumer group whenever you're scaling out event processing and duplicate work isn't just wasteful but actively wrong (double-charging a customer, double-sending an alert).
It works by naming a group when you subscribe: every subscriber that passes the same group string to subscribeToEvents joins that group, and the broker round-robins each event to exactly one member instead of fanning it out to all of them. Passing an empty group string reverts to normal fan-out, so the same subscription call can flip between the two delivery models with one argument.
Gotchas: consumer groups are scoped per channel — subscribing to the same group on a different channel does not share load balancing across channels. A group with zero active subscribers behaves like no subscribers at all; events aren't queued for a group that's temporarily empty the way they are for durable queue messages. And because delivery is round-robin rather than content-aware, you can't route specific events to specific workers within a group — if you need that, partition by channel instead.
Prerequisites
- KubeMQ server running on
localhost:50000 - Node.js SDK installed (
npm install kubemq-js)
Code
/**
* Example: Events Subscribe with Consumer Group
*
* Demonstrates load-balanced event delivery using a consumer group.
* Two subscribers join the same group — each event is delivered to
* exactly one subscriber in the group instead of being fanned out
* to all of them.
*
* Prerequisites:
* - KubeMQ server running on localhost:50000
*
* Run: npx tsx examples/events/consumer-group.ts
*/
import { KubeMQClient, createEventMessage } from 'kubemq-js';
async function main(): Promise<void> {
const client = await KubeMQClient.create({
address: 'localhost:50000',
clientId: 'js-events-consumer-group-client',
});
try {
const sub1 = client.subscribeToEvents({
channel: 'js-events.consumer-group',
group: 'workers',
onEvent: (event) => {
console.log('[Worker A]', new TextDecoder().decode(event.body));
},
onError: (err) => {
console.error('Worker A error:', err.message);
},
});
const sub2 = client.subscribeToEvents({
channel: 'js-events.consumer-group',
group: 'workers',
onEvent: (event) => {
console.log('[Worker B]', new TextDecoder().decode(event.body));
},
onError: (err) => {
console.error('Worker B error:', err.message);
},
});
for (let i = 1; i <= 6; i++) {
await client.sendEvent(
createEventMessage({ channel: 'js-events.consumer-group', body: `task-${i}` }),
);
}
await new Promise((resolve) => setTimeout(resolve, 1000));
sub1.cancel();
sub2.cancel();
} finally {
await client.close();
}
}
main().catch(console.error);
How It Works
- Both
sub1andsub2passgroup: 'workers'tosubscribeToEvents()— this enrolls them in the same competing-consumers group. - KubeMQ load-balances delivery: each of the 6 published events goes to exactly one worker, never both.
- Without the
groupoption each subscriber would receive all events (fan-out); the group option switches to round-robin dispatch. - Both handles are cancelled at the end; the
close()call on the client then drains any pending callbacks before disconnecting.
Related
Was this page helpful?