New consumer
This commit is contained in:
@@ -19,6 +19,7 @@
|
|||||||
"dotenv": "^16.4.5",
|
"dotenv": "^16.4.5",
|
||||||
"express": "^4.17.1",
|
"express": "^4.17.1",
|
||||||
"kafkajs": "^2.2.4",
|
"kafkajs": "^2.2.4",
|
||||||
|
"nodemailer": "^6.9.13",
|
||||||
"openapi-types": "^10.0.0",
|
"openapi-types": "^10.0.0",
|
||||||
"pg": "8.7.1",
|
"pg": "8.7.1",
|
||||||
"pg-pool": "^3.5.1",
|
"pg-pool": "^3.5.1",
|
||||||
|
|||||||
@@ -32,18 +32,45 @@ routes(app);
|
|||||||
|
|
||||||
//eventInterest
|
//eventInterest
|
||||||
|
|
||||||
|
kafka.consume("NEW_OFFER_110022", (value) => {
|
||||||
|
try {
|
||||||
|
console.log("Receive NEW_OFFER_110022 message xxxx: ", value);
|
||||||
|
var obj = phpUnserialize(value);
|
||||||
|
console.log(obj);
|
||||||
|
if (obj.depend_uid !== undefined && obj.depend_uid.length > 10){
|
||||||
|
console.log("We need contact people that did this job: ", obj.depend_uid);
|
||||||
|
|
||||||
|
}
|
||||||
|
} catch (exceptionVar) {
|
||||||
|
console.log(" Error ", exceptionVar.message);
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
|
kafka.consume("INTEREST_MSG", (value) => {
|
||||||
|
try {
|
||||||
|
console.log("Receive INTEREST_MSG message xxxx: ", value);
|
||||||
|
var obj = phpUnserialize(value);
|
||||||
|
console.log(obj);
|
||||||
|
|
||||||
|
} catch (exceptionVar) {
|
||||||
|
console.log(" Error ", exceptionVar.message);
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
kafka.consume("FLUTTER_PAYMENT_RECEIVED", (value) => {
|
kafka.consume("FLUTTER_PAYMENT_RECEIVED", (value) => {
|
||||||
try {
|
try {
|
||||||
console.log("📨 Receive message xxxx: ", value);
|
console.log("Receive FLUTTER_PAYMENT_RECEIVED message xxxx: ", value);
|
||||||
var obj = phpUnserialize(value);
|
var obj = phpUnserialize(value);
|
||||||
console.log(obj);
|
console.log(obj);
|
||||||
console.log("📨 Receive message yyyy: ", obj.txRef);
|
console.log("Receive message yyyy: ", obj.txRef);
|
||||||
} catch (exceptionVar) {
|
} catch (exceptionVar) {
|
||||||
console.log("📨 Error ", exceptionVar.message);
|
console.log("Error ", exceptionVar.message);
|
||||||
}
|
}
|
||||||
});
|
});
|
||||||
|
|
||||||
|
|
||||||
|
|
||||||
|
|
||||||
app.listen(port, "0.0.0.0", function() {
|
app.listen(port, "0.0.0.0", function() {
|
||||||
logger.info('***** Server started on port: ' + port + ' *****');
|
logger.info('***** Server started on port: ' + port + ' *****');
|
||||||
});
|
});
|
||||||
|
|||||||
Reference in New Issue
Block a user