-
Notifications
You must be signed in to change notification settings - Fork 3
/
kafka.js
90 lines (78 loc) · 2.74 KB
/
kafka.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
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
const kafka = require('kafka-node');
class Subscription {
constructor(subscription) {
this.subscription = subscription;
this.subscription.options.requestTimeout = subscription.options.requestTimeout || 10000;
this.subscription.options.connectTimeout = subscription.options.connectTimeout || 10000;
this.client = new kafka.KafkaClient(this.subscription.client);
this.offset = new kafka.Offset(this.client);
}
async subscribe(context) {
try {
this.client.on('error', (err) => {
const message = `Error subscribing to kafka ${JSON.stringify(err)}`;
context.logger.error(message);
throw message;
});
return await this.fetchOffset();
} catch (exc) {
const message = `Error connecting kafka ${JSON.stringify(exc)}`;
context.logger.error(message);
throw message;
}
};
fetchOffset() {
return new Promise((resolve, reject) => {
try {
this.offset.fetchLatestOffsets([this.subscription.options.topic], async (error, offsets) => {
if (error) {
reject(error);
} else {
this.latestOffset = offsets[this.subscription.options.topic][0];
this.ableToUnsubscribe = true;
resolve();
}
});
} catch (err) {
reject(err);
}
});
}
async unsubscribe() {
if (this.ableToUnsubscribe) {
this.client.close();
}
}
createConsumer() {
return new kafka.Consumer(
this.client,
[{
topic: this.subscription.options.topic,
offset: this.latestOffset
}],
{
fromOffset: true
}
);
}
receiveMessage(context) {
return new Promise((resolve, reject) => {
const consumer = this.createConsumer();
consumer.on('message', (message) => {
context.logger.trace('Kafka message data: ' + JSON.stringify(message));
resolve(message);
consumer.close(() => {
context.logger.trace('Kafka consumer is closed');
});
});
consumer.on('error', (error) => {
context.logger.error('Kafka error message data: ' + JSON.stringify(error));
reject(error);
consumer.close(() => {
context.logger.trace('Kafka consumer is closed');
});
});
});
};
}
module.exports = {Subscription};