-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathindex.js
More file actions
52 lines (41 loc) · 1.4 KB
/
Copy pathindex.js
File metadata and controls
52 lines (41 loc) · 1.4 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
const express = require("express");
const cors = require('cors')
const grpcServer = require('./grpc/server');
const kafka = require('kafka-node');
const Queue = require('bull');
require("@opentelemetry/api");
const app = express();
app.use(cors());
app.all("/", (req, res) => {
res.json({ method: req.method, message: "You are accessing Micro App", ...req.body });
});
app.all("/command", async (req, res) => {
try {
res.json({ success: "true", message: "Success execute command" });
} catch (error) {
console.error('Error:', error);
}
});
// run grpc server
grpcServer();
// kafka
const kafkaClient = new kafka.KafkaClient({kafkaHost: process.env.KAFKA_HOST || 'localhost:9092'});
const kafkaTopics = [{ topic: process.env.KAFKA_TOPIC }];
const kafkaConsumer = new kafka.Consumer(kafkaClient, kafkaTopics, { autoCommit: true });
// Set up the event handlers
kafkaConsumer.on('message', function (message) {
console.log('Received message:', message);
});
kafkaConsumer.on('error', function (error) {
console.error('Error in consumer:', error);
});
// bull
const bullQueue = new Queue(process.env.REDIS_TOPIC, process.env.REDIS_URL || 'redis://127.0.0.1:6379');
bullQueue.process(async (job, done) => {
console.log('Received payment:', job.data.payment);
done();
});
const port = process.env.APP_PORT || "3002";
app.listen(port, function() {
console.log('server running on port ' + port + '.');
});