Repository navigation
Expand file tree
/
Copy pathconsumer.js
More file actions
143 lines (129 loc) · 4.75 KB
/
Copy pathconsumer.js
File metadata and controls
143 lines (129 loc) · 4.75 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
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
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
/**
* Kafka consumer
*/
'use strict';
const config = require('config');
const _ = require('lodash');
const Kafka = require('no-kafka');
const co = require('co');
global.Promise = require('bluebird');
const healthcheck = require('topcoder-healthcheck-dropin');
const logger = require('./src/common/logger');
const models = require('./src/models');
const processors = require('./src/processors');
/**
* Start Kafka consumer
*/
function startKafkaConsumer() {
const options = { groupId: config.KAFKA_GROUP_ID, connectionString: config.KAFKA_URL };
if (config.KAFKA_CLIENT_CERT && config.KAFKA_CLIENT_CERT_KEY) {
options.ssl = { cert: config.KAFKA_CLIENT_CERT, key: config.KAFKA_CLIENT_CERT_KEY };
}
const consumer = new Kafka.GroupConsumer(options);
// data handler
const messageHandler = (messageSet, topic, partition) => Promise.each(messageSet, (m) => {
const message = m.message.value.toString('utf8');
logger.info(`Handle Kafka event message; Topic: ${topic}; Partition: ${partition}; Offset: ${
m.offset}; Message: ${message}.`);
let messageJSON;
try {
messageJSON = JSON.parse(message);
} catch (e) {
logger.error('Invalid message JSON.');
logger.logFullError(e);
// commit the message and ignore it
consumer.commitOffset({ topic, partition, offset: m.offset });
return;
}
if (messageJSON.topic !== topic) {
logger.error(`The message topic ${messageJSON.topic} doesn't match the Kafka topic ${topic}.`);
// commit the message and ignore it
consumer.commitOffset({ topic, partition, offset: m.offset });
return;
}
// get rule sets for the topic
const ruleSets = config.KAFKA_CONSUMER_RULESETS[topic];
// TODO for NULL handler
if (!ruleSets || ruleSets.length === 0) {
logger.error(`No handler configured for Kafka topic ${topic}.`);
// commit the message and ignore it
consumer.commitOffset({ topic, partition, offset: m.offset });
return;
}
return co(function* () {
// run each handler
for (let i = 0; i < ruleSets.length; i += 1) {
const rule = ruleSets[i];
const handlerFuncArr = _.keys(rule);
const handlerFuncName = _.get(handlerFuncArr, '0');
try {
const handler = processors[handlerFuncName];
const handlerRuleSets = rule[handlerFuncName];
if (!handler) {
logger.error(`Handler ${handlerFuncName} is not defined`);
continue;
}
logger.info(`Run handler ${handlerFuncName}`);
// run handler to get notifications
const notifications = yield handler(messageJSON, handlerRuleSets);
if (notifications && notifications.length > 0) {
// save notifications in bulk to improve performance
logger.info(`Going to insert ${notifications.length} notifications in database.`);
yield models.Notification.bulkCreate(_.map(notifications, (n) => ({
userId: n.userId,
type: n.type || topic,
contents: n.contents || n.notification || messageJSON.payload || {},
read: false,
seen: false,
version: n.version || null,
})));
// logging
logger.info(`Saved ${notifications.length} notifications`);
/* logger.info(` for users: ${
_.map(notifications, (n) => n.userId).join(', ')
}`); */
}
logger.info(`Handler ${handlerFuncName} executed successfully`);
} catch (e) {
// log and ignore error, so that it won't block rest handlers
logger.error(`Handler ${handlerFuncName} failed`);
logger.logFullError(e);
}
}
})
// commit offset
.then(() => consumer.commitOffset({ topic, partition, offset: m.offset }))
.catch((err) => {
logger.error('Kafka handler failed');
logger.logFullError(err);
});
});
const check = function () {
if (!consumer.client.initialBrokers && !consumer.client.initialBrokers.length) {
return false;
}
let connected = true;
consumer.client.initialBrokers.forEach(conn => {
logger.debug(`url ${conn.server()} - connected=${conn.connected}`);
connected = conn.connected & connected;
});
return connected;
};
// Start kafka consumer
logger.info('Starting kafka consumer');
consumer
.init([{
// subscribe topics
subscriptions: _.keys(config.KAFKA_CONSUMER_RULESETS),
handler: messageHandler,
}])
.then(() => {
logger.info('Kafka consumer initialized successfully');
healthcheck.init([check]);
})
.catch((err) => {
logger.error('Kafka consumer failed');
logger.logFullError(err);
});
}
startKafkaConsumer();