-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathkafka.test.js
More file actions
65 lines (57 loc) · 1.63 KB
/
Copy pathkafka.test.js
File metadata and controls
65 lines (57 loc) · 1.63 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
const expect = require('chai').expect;
const { KafkaClient, Producer, Consumer } = require('kafka-node');
describe('Kafka example', () => {
const topicName = 'kafka-node-testing';
const kafkaClientConfig = { kafkaHost: 'localhost:9092' };
before(done => {
createTopic(done);
});
it('sends and receives a message', (done) => {
const testMessage = 'a test message';
produceMessage(testMessage);
consumeMessage(message => {
expect(message.topic).to.equal(topicName);
expect(message.value).to.equal(testMessage);
done();
});
});
function createTopic(done) {
const client = new KafkaClient(kafkaClientConfig);
client.createTopics([{
topic: topicName,
partitions: 1,
replicationFactor: 1
}], (err, result) => {
client.close(done);
});
}
function produceMessage(message) {
const producer = new Producer(new KafkaClient(kafkaClientConfig));
producer.on('error', (err) => {
console.log(err);
console.log('connection errored');
throw err;
});
const payloads = [
{
topic: topicName,
messages: message
}
];
producer.on('ready', () => {
producer.send(payloads, (err, data) => {
if (err) {
console.log('broker update fail');
}
producer.close();
});
});
}
function consumeMessage(messageHandler) {
const consumerPayloads = [{ topic: topicName }];
const consumer = new Consumer(new KafkaClient(kafkaClientConfig), consumerPayloads, {});
consumer.on('message', (message) => {
consumer.close(() => messageHandler(message));
});
}
});