Adicionado REadme
This commit is contained in:
parent
82fb88a4e7
commit
be85ebef21
15
consumer.js
15
consumer.js
@ -5,19 +5,9 @@ const kafka = new Kafka({
|
|||||||
brokers: ['localhost:9092']
|
brokers: ['localhost:9092']
|
||||||
})
|
})
|
||||||
|
|
||||||
//const producer = kafka.producer()
|
|
||||||
const consumer = kafka.consumer({ groupId: 'test-group' })
|
const consumer = kafka.consumer({ groupId: 'test-group' })
|
||||||
|
|
||||||
const run = async () => {
|
const run = async () => {
|
||||||
// // Producing
|
|
||||||
// await producer.connect()
|
|
||||||
// await producer.send({
|
|
||||||
// topic: 'test-topic',
|
|
||||||
// messages: [
|
|
||||||
// { value: 'Hello KafkaJS user!' },
|
|
||||||
// ],
|
|
||||||
// })
|
|
||||||
|
|
||||||
// Consuming
|
// Consuming
|
||||||
await consumer.connect()
|
await consumer.connect()
|
||||||
await consumer.subscribe({ topic: process.env.TOPIC, fromBeginning: true })
|
await consumer.subscribe({ topic: process.env.TOPIC, fromBeginning: true })
|
||||||
@ -27,11 +17,6 @@ const run = async () => {
|
|||||||
const obj = JSON.parse(message.value)
|
const obj = JSON.parse(message.value)
|
||||||
console.log('Message consumer successfully!');
|
console.log('Message consumer successfully!');
|
||||||
console.log(obj.name)
|
console.log(obj.name)
|
||||||
// console.log({
|
|
||||||
// partition,
|
|
||||||
// offset: message.offset,
|
|
||||||
// value: Json.parse(message.value),
|
|
||||||
// })
|
|
||||||
},
|
},
|
||||||
|
|
||||||
})
|
})
|
||||||
|
2
package-lock.json
generated
2
package-lock.json
generated
@ -5,7 +5,7 @@
|
|||||||
"packages": {
|
"packages": {
|
||||||
"": {
|
"": {
|
||||||
"dependencies": {
|
"dependencies": {
|
||||||
"kafkajs": "^2.2.2"
|
"kafkajs": "^2.2.4"
|
||||||
}
|
}
|
||||||
},
|
},
|
||||||
"node_modules/kafkajs": {
|
"node_modules/kafkajs": {
|
||||||
|
@ -1,5 +1,5 @@
|
|||||||
{
|
{
|
||||||
"dependencies": {
|
"dependencies": {
|
||||||
"kafkajs": "^2.2.2"
|
"kafkajs": "^2.2.4"
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
15
producer.js
15
producer.js
@ -30,21 +30,6 @@ const run = async () => {
|
|||||||
} finally {
|
} finally {
|
||||||
await producer.disconnect();
|
await producer.disconnect();
|
||||||
}
|
}
|
||||||
|
|
||||||
// Consuming
|
|
||||||
// await consumer.connect()
|
|
||||||
// await consumer.subscribe({ topic: 'test-topic', fromBeginning: true })
|
|
||||||
|
|
||||||
// await consumer.run({
|
|
||||||
// eachMessage: async ({ topic, partition, message }) => {
|
|
||||||
// console.log({
|
|
||||||
// partition,
|
|
||||||
// offset: message.offset,
|
|
||||||
// value: message.value.toString(),
|
|
||||||
// })
|
|
||||||
// },
|
|
||||||
|
|
||||||
// })
|
|
||||||
}
|
}
|
||||||
|
|
||||||
run()
|
run()
|
@ -2,7 +2,7 @@
|
|||||||
# yarn lockfile v1
|
# yarn lockfile v1
|
||||||
|
|
||||||
|
|
||||||
"kafkajs@^2.2.2":
|
"kafkajs@^2.2.4":
|
||||||
"integrity" "sha512-j/YeapB1vfPT2iOIUn/vxdyKEuhuY2PxMBvf5JWux6iSaukAccrMtXEY/Lb7OvavDhOWME589bpLrEdnVHjfjA=="
|
"integrity" "sha512-j/YeapB1vfPT2iOIUn/vxdyKEuhuY2PxMBvf5JWux6iSaukAccrMtXEY/Lb7OvavDhOWME589bpLrEdnVHjfjA=="
|
||||||
"resolved" "https://registry.npmjs.org/kafkajs/-/kafkajs-2.2.4.tgz"
|
"resolved" "https://registry.npmjs.org/kafkajs/-/kafkajs-2.2.4.tgz"
|
||||||
"version" "2.2.4"
|
"version" "2.2.4"
|
||||||
|
Loading…
Reference in New Issue
Block a user