Node.js
Dependencies
npm install @grpc/grpc-js @grpc/proto-loader
No code-generation step required — proto-loader loads proto files at runtime.
Connect
const grpc = require('@grpc/grpc-js');
const protoLoader = require('@grpc/proto-loader');
const path = require('path');
const PROTO_DIR = path.join(__dirname, 'proto');
const packageDef = protoLoader.loadSync(
path.join(PROTO_DIR, 'tinybroker/v1/service.proto'),
{
keepCase: true,
longs: String,
defaults: true,
oneofs: true,
includeDirs: [PROTO_DIR],
}
);
const { tinybroker: { v1: { Broker } } } =
grpc.loadPackageDefinition(packageDef);
const client = new Broker(
'tinybroker:50051',
grpc.credentials.createInsecure()
);
Publish
client.publish(
{ topic: 'events.user.signup', payload: Buffer.from('{"user_id":"abc"}') },
(err, response) => {
if (err) console.error(err);
}
);
// Or with async/await via util.promisify:
const { promisify } = require('util');
const publish = promisify(client.publish.bind(client));
await publish({ topic: 'events.user.signup', payload: Buffer.from('{}') });
Subscribe (fan-out)
const stream = client.subscribe();
// Open subscription
stream.write({
open: {
topic_patterns: ['events.user.*', 'events.order.*'],
mode: 1, // DELIVERY_MODE_MULTI
},
});
let subscriptionId;
stream.on('data', (event) => {
if (event.subscription_id) {
subscriptionId = event.subscription_id;
console.log('subscription confirmed:', subscriptionId);
} else if (event.message) {
const { topic, payload, system_id } = event.message;
console.log('message:', topic, payload.toString());
} else if (event.patterns_updated) {
console.log('patterns now:', event.patterns_updated.topic_patterns);
}
});
stream.on('error', (err) => {
console.error('stream error:', err);
// implement reconnect
});
stream.on('end', () => {
console.log('stream ended');
});
Subscribe (consumer group)
const stream = client.subscribe();
stream.write({
open: {
subscription_id: 'job-workers',
topic_patterns: ['jobs.*'],
mode: 2, // DELIVERY_MODE_BALANCED
consumer_group: 'job-workers',
},
});
stream.on('data', async (event) => {
if (event.message) {
const { topic, payload, system_id } = event.message;
await processJob(payload);
// Acknowledge so the next message is delivered
await promisify(client.ack.bind(client))({
subscription_id: 'job-workers',
system_id,
});
}
});
Dynamic pattern management
// Add patterns on a live subscription
stream.write({
add_patterns: { topic_patterns: ['alerts.*'] },
});
// Remove patterns
stream.write({
remove_patterns: { topic_patterns: ['events.user.*'] },
});
TypeScript
Install types:
npm install --save-dev @types/node grpc-tools
Generate typed stubs with protoc + the grpc TypeScript plugin, or use ts-proto:
npm install ts-proto
protoc \
--plugin=./node_modules/.bin/protoc-gen-ts_proto \
--ts_proto_out=gen \
--ts_proto_opt=outputServices=grpc-js \
-I proto \
tinybroker/v1/service.proto tinybroker/v1/client.proto \
tinybroker/v1/server.proto tinybroker/v1/shared.proto
Reconnect pattern
function connectWithReconnect() {
const stream = client.subscribe();
stream.write({ open: { topic_patterns: ['events.*'], mode: 1 } });
stream.on('data', handleEvent);
stream.on('error', (err) => {
console.error('reconnecting in 2s:', err.message);
setTimeout(connectWithReconnect, 2000);
});
stream.on('end', connectWithReconnect);
}