-
Notifications
You must be signed in to change notification settings - Fork 0
/
Copy pathindex.js
92 lines (78 loc) · 2.19 KB
/
index.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
91
92
import pkg from "kafkajs";
const { Kafka } = pkg;
import dotenv from "dotenv";
import axios from "axios";
dotenv.config();
// Create the client with the broker list, minimum 1 broker(bootstrap) is needed
// The client will auto-fetch the metadata of others
const kafka = new Kafka({
clientId: "push-data-service-" + Date.now(), // Append Current Epoch milliseconds for Random Id
brokers: [
process.env.KAFKA_BOOTSTRAP_SERVER_URL ||
"my-cluster-kafka-bootstrap.kafka:9092",
],
});
// Consumer
const consumerData = kafka.consumer({
groupId: "push-data-service-group-kinesis",
});
// Producer
const producer = kafka.producer();
const run = async () => {
await consumerData.connect();
await producer.connect();
console.info("Connected to Kafka Broker.");
await consumerData.subscribe({
topic: process.env.SUBSCRIBE_TOPIC,
fromBeginning: false,
});
consumerData.run({
eachMessage: async ({ topic, partition, message }) => {
try {
let payLoadParsed = JSON.parse(message.value.toString());
if (payLoadParsed) {
let payloadArr = [];
let obj = {
key: "TelematicsEvent",
value: JSON.stringify(payLoadParsed),
};
payloadArr.push(obj);
await producer.send({
topic: process.env.PUBLISH_TRACK_TOPIC,
messages: [
{
key: payLoadParsed.imeiNo,
value: JSON.stringify(payloadArr),
},
],
});
}
} catch (error) {
console.log("Eror: ",error);
}
},
});
};
run().catch("run error: ", console.error);
// Consumer Crash Events
consumerData.on("consumer.crash", function () {
console.log("Crash detected");
process.exit(0);
});
consumerData.on("consumer.disconnect", function () {
console.log("Disconnect detected");
process.exit(0);
});
consumerData.on("consumer.stop", function () {
console.log("Stop detected");
process.exit(0);
});
const errorTypes = ["unhandledRejection"];
errorTypes.map((type) => {
process.on(type, async (e) => {
console.log(`process.on ${type}`);
console.error(e);
// await consumer.disconnect()
process.exit(0);
});
});