|
1 | 1 | 'use strict'; |
2 | 2 |
|
3 | | -const kafka = require('no-kafka'); |
| 3 | +const { Kafka } = require('kafkajs'); |
4 | 4 |
|
5 | 5 | const KAFKA_HOSTS = process.env.KAFKA_HOSTS || 'localhost:9092'; |
6 | 6 | const KAFKA_SOURCE_TOPICS = process.env.KAFKA_SOURCE_TOPICS || 'example'; |
7 | 7 |
|
8 | | -const kafkaProducer = (message) => { |
9 | | - const producer = new kafka.Producer({ |
10 | | - connectionString: KAFKA_HOSTS, |
11 | | - }); |
12 | | - |
13 | | - const stringMessageValue = JSON.stringify(message); |
14 | | - |
15 | | - return producer.init() |
16 | | - .then(() => { |
17 | | - const data = { |
18 | | - topic: KAFKA_SOURCE_TOPICS, |
19 | | - message: { |
20 | | - value: stringMessageValue, |
21 | | - }, |
22 | | - }; |
23 | | - |
24 | | - return producer.send(data); |
25 | | - }) |
26 | | - .then((result) => { |
27 | | - producer.end(); |
28 | | - |
29 | | - const { error } = result[0]; |
30 | | - if (error) { |
31 | | - return error; |
32 | | - } |
33 | | - |
34 | | - console.log(`Message successfully produced to Kafka ${JSON.stringify(result)}`); |
35 | | - return result; |
36 | | - }) |
37 | | - .catch((e) => { |
38 | | - console.log(`Could not produce message to topic ${KAFKA_SOURCE_TOPICS}`); |
39 | | - console.log(e); |
40 | | - producer.end(); |
41 | | - return e; |
42 | | - }); |
43 | | -}; |
| 8 | +const kafka = new Kafka({ |
| 9 | + clientId: 'channel-service', |
| 10 | + brokers: [KAFKA_HOSTS], |
| 11 | +}); |
| 12 | +const producer = kafka.producer(); |
44 | 13 |
|
45 | 14 | exports.produce = { |
46 | | - handler: (request) => { |
47 | | - const msg = request.payload; |
48 | | - |
49 | | - return kafkaProducer(msg) |
50 | | - .then(message => ({ status: 'ok', message })) |
51 | | - .catch(message => ({ status: 'error', message })); |
52 | | - }, |
| 15 | + handler: async (request) => { |
| 16 | + const msg = request.payload; |
| 17 | + |
| 18 | + try { |
| 19 | + await producer.connect(); |
| 20 | + await producer.send({ |
| 21 | + topic: KAFKA_SOURCE_TOPICS, |
| 22 | + messages: [{ value: JSON.stringify(msg) }], |
| 23 | + }); |
| 24 | + await producer.disconnect(); |
| 25 | + console.log('Message successfully produced to Kafka'); |
| 26 | + |
| 27 | + return { status: 'ok' }; |
| 28 | + } catch (error) { |
| 29 | + console.error('Failed to produce message to Kafka', error); |
| 30 | + return { status: 'error', message: error }; |
| 31 | + } |
| 32 | + }, |
53 | 33 | }; |
0 commit comments