@fastify/kafka
v4.0.0
Published
Fastify plugin to interact with Apache Kafka.
Downloads
2,002
Readme
@fastify/kafka
Fastify plugin to interact with Apache Kafka, supporting Kafka producers and consumers.
To achieve the best performance, the plugin uses @platformatic/kafka.
Install
npm i @fastify/kafkaCompatibility
| Plugin version | Fastify version |
| ---------------|-----------------|
| >=3.x | ^5.x |
| >=0.x <3.x | ^4.x |
| >=0.x <3.x | ^3.x |
| >=0.x <3.x | ^2.x |
| >=0.x <3.x | ^1.x |
Please note that if a Fastify version is out of support, then so are the corresponding versions of this plugin in the table above. See Fastify's LTS policy for more details.
Usage
const crypto = require('node:crypto')
const fastify = require('fastify')()
const group = crypto.randomBytes(20).toString('hex')
fastify
.register(require('@fastify/kafka'), {
producer: {
bootstrapBrokers: ['127.0.0.1:9092'],
clientId: 'my-producer',
allowAutoTopicCreation: true
},
consumer: {
bootstrapBrokers: ['127.0.0.1:9092'],
groupId: group,
clientId: 'my-consumer',
allowAutoTopicCreation: true
}
})
fastify.post('/data', (req, reply) => {
fastify.kafka.push({
topic: 'updates',
payload: req.body,
key: 'dataKey'
})
reply.send({ ok: true })
})
fastify.kafka.subscribe('updates')
fastify.kafka.on('updates', (msg, commit) => {
console.log(msg.value.toString())
commit()
})
fastify.listen({ port: 3000 }, err => {
if (err) throw err
console.log(`Server listening on ${fastify.server.address().port}`)
fastify.kafka.consume()
})For more examples on how to use this plugin, you can take a look at the examples directory.
Migration from node-rdkafka
Starting from version 3.x, this plugin uses @platformatic/kafka instead of node-rdkafka.
The main differences in the configuration options are:
| node-rdkafka (old) | @platformatic/kafka (new) |
|---------------------------------|-----------------------------------|
| metadata.broker.list | bootstrapBrokers: ['host:port'] |
| group.id | groupId |
| fetch.wait.max.ms | maxWaitTime |
| enable.auto.commit | autocommit |
| auto.offset.reset | passed via consumerTopicConf |
| dr_cb | not needed |
API
This module exposes the following APIs:
Plugin options
{
producer?: Record<string, unknown> // @platformatic/kafka Producer options
consumer?: Record<string, unknown> // @platformatic/kafka Consumer options
producerTopicConf?: Record<string, unknown>
consumerTopicConf?: Record<string, unknown>
metadataOptions?: Record<string, unknown>
}Producer
fastify.kafka.producer— the producer instancefastify.kafka.push(message)— send a message to a topic
fastify.kafka.push({
topic: 'my-topic', // required
payload: 'hello', // required - string or Buffer
key: 'myKey', // optional
partition: 0 // optional
})Consumer
fastify.kafka.consumer— the consumer instancefastify.kafka.subscribe(topic | topic[])— subscribe to one or more topicsfastify.kafka.consume()— start streaming messages (flow mode)fastify.kafka.consume(n, callback)— fetch exactlynmessages (batch mode)fastify.kafka.on(topic, handler)— listen for messages on a specific topic
// Flow mode
fastify.kafka.subscribe(['topic-a', 'topic-b'])
fastify.kafka.on('topic-a', (msg, commit) => {
console.log(msg.value.toString())
commit() // manually commit the offset
})
fastify.kafka.consume()
// Batch mode
fastify.kafka.subscribe('topic-a')
fastify.kafka.consume(10, (err, messages) => {
if (err) throw err
messages.forEach(msg => console.log(msg.value.toString()))
})Acknowledgments
This project is kindly sponsored by:
Past sponsors:
- LetzDoIt
License
Licensed under MIT.
