-
Notifications
You must be signed in to change notification settings - Fork 0
/
consumer.js
66 lines (55 loc) · 1.74 KB
/
consumer.js
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
66
const Kafka = require('node-rdkafka');
const { configFromPath } = require('./util');
function createConfigMap(config) {
if (config.hasOwnProperty('security.protocol')) {
return {
'bootstrap.servers': config['bootstrap.servers'],
'sasl.username': config['sasl.username'],
'sasl.password': config['sasl.password'],
'security.protocol': config['security.protocol'],
'sasl.mechanisms': config['sasl.mechanisms'],
'group.id': 'kafka-nodejs-getting-started'
}
} else {
return {
'bootstrap.servers': config['bootstrap.servers'],
'group.id': 'kafka-nodejs-getting-started'
}
}
}
function createConsumer(config, onData) {
const consumer = new Kafka.KafkaConsumer(
createConfigMap(config),
{'auto.offset.reset': 'earliest'});
return new Promise((resolve, reject) => {
consumer
.on('ready', () => resolve(consumer))
.on('data', onData);
consumer.connect();
});
};
async function consumerKafka() {
if (process.argv.length < 3) {
console.log("Please provide the configuration file path as the command line argument");
process.exit(1);
}
let configPath = process.argv.slice(2)[0];
const config = await configFromPath(configPath);
//let seen = 0;
let topic = "kafka-topic";
const consumer = await createConsumer(config, ({key, value}) => {
let k = key.toString().padEnd(10, ' ');
console.log(`Consumed event from topic ${topic}: key = ${k} value = ${value}`);
});
consumer.subscribe([topic]);
consumer.consume();
process.on('SIGINT', () => {
console.log('\nDisconnecting consumer ...');
consumer.disconnect();
});
}
consumerKafka()
.catch((err) => {
console.error(`Something went wrong:\n${err}`);
process.exit(1);
});