kakka parts
This commit is contained in:
+1
-1
@@ -16,6 +16,6 @@ module.exports = function(app) {
|
|||||||
.get(controller.getUsersEscrows);
|
.get(controller.getUsersEscrows);
|
||||||
|
|
||||||
app.route('/flutterOkHook')
|
app.route('/flutterOkHook')
|
||||||
.get(controller.flutterOkHook);
|
.post(controller.flutterOkHook);
|
||||||
|
|
||||||
};
|
};
|
||||||
@@ -0,0 +1,46 @@
|
|||||||
|
// import { Kafka } from "kafkajs";
|
||||||
|
const {Kafka} = require('kafkajs');
|
||||||
|
|
||||||
|
const kafka = new Kafka({
|
||||||
|
clientId: "nodejs-kafka",
|
||||||
|
brokers: ["10.10.10.120:9092"],
|
||||||
|
});
|
||||||
|
class KafkaConfig {
|
||||||
|
constructor() {
|
||||||
|
|
||||||
|
this.producer = kafka.producer();
|
||||||
|
this.consumer = kafka.consumer({ groupId: "test-group" });
|
||||||
|
}
|
||||||
|
|
||||||
|
async produce(topic, messages) {
|
||||||
|
try {
|
||||||
|
await this.producer.connect();
|
||||||
|
await this.producer.send({
|
||||||
|
topic: topic,
|
||||||
|
messages: messages,
|
||||||
|
});
|
||||||
|
} catch (error) {
|
||||||
|
console.error(error);
|
||||||
|
} finally {
|
||||||
|
await this.producer.disconnect();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
async consume(topic, callback) {
|
||||||
|
try {
|
||||||
|
await this.consumer.connect();
|
||||||
|
await this.consumer.subscribe({ topic: topic, fromBeginning: true });
|
||||||
|
await this.consumer.run({
|
||||||
|
eachMessage: async ({ topic, partition, message }) => {
|
||||||
|
const value = message.value.toString();
|
||||||
|
callback(value);
|
||||||
|
},
|
||||||
|
});
|
||||||
|
} catch (error) {
|
||||||
|
console.error(error);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
module.exports = KafkaConfig;
|
||||||
|
//export default KafkaConfig;
|
||||||
@@ -4,6 +4,7 @@ const cookieParser = require('cookie-parser');
|
|||||||
const bodyParser = require('body-parser');
|
const bodyParser = require('body-parser');
|
||||||
const logger = require('./app/logger');
|
const logger = require('./app/logger');
|
||||||
const port = process.env.PORT || 3000;
|
const port = process.env.PORT || 3000;
|
||||||
|
const KafkaConfig = require("./app/kconfig");
|
||||||
|
|
||||||
const app = express();
|
const app = express();
|
||||||
|
|
||||||
|
|||||||
+22
-8
@@ -3,20 +3,34 @@
|
|||||||
const request = require('request');
|
const request = require('request');
|
||||||
const db = require('../app/db')
|
const db = require('../app/db')
|
||||||
const logger = require('../app/logger');
|
const logger = require('../app/logger');
|
||||||
|
const KafkaConfig = require("../app/kconfig");
|
||||||
|
|
||||||
var ebroker = {
|
var ebroker = {
|
||||||
eventPublish: function (req, res, next) {
|
eventPublish: function (req, res, next) {
|
||||||
|
|
||||||
//console.log("REQ---->",req.body.uid);
|
try {
|
||||||
var data = {
|
const { message } = req.body;
|
||||||
"uid": req.body.uid,
|
console.log('THIS-> req.body -> ', req.body);
|
||||||
"member_id": req.body.member_id,
|
const kafkaConfig = new KafkaConfig();
|
||||||
"limit": (req.body.limit != null && req.body.limit !== "") ? req.body.limit : 20,
|
const messages = [{ key: "key1", value: message }];
|
||||||
"sessionid": req.body.sessionid,
|
// const messages = [{ key: "key1", value: "ameye olusesan" }];
|
||||||
"page": req.body.page
|
kafkaConfig.produce("FLUTTER_PAYMENT_RECEIVED", messages).then(r =>{
|
||||||
};
|
console.log('THIS->RET-> ',r);
|
||||||
|
} );
|
||||||
|
|
||||||
|
res.status(200).json({
|
||||||
|
status: "Ok!",
|
||||||
|
message: "Message successfully send!",
|
||||||
|
});
|
||||||
|
} catch (error) {
|
||||||
|
console.log(error);
|
||||||
|
}
|
||||||
|
|
||||||
|
let resultItem ={
|
||||||
|
"result": [],
|
||||||
|
"total_record": 0
|
||||||
|
}
|
||||||
|
next(null, resultItem ); // pass control to the next handler
|
||||||
},
|
},
|
||||||
|
|
||||||
eventConsumer: function (req, res, next) {
|
eventConsumer: function (req, res, next) {
|
||||||
|
|||||||
Reference in New Issue
Block a user