-
Notifications
You must be signed in to change notification settings - Fork 2
Expand file tree
/
Copy pathMQProducer.js
More file actions
99 lines (81 loc) · 3.12 KB
/
Copy pathMQProducer.js
File metadata and controls
99 lines (81 loc) · 3.12 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
/**
* Created by Gang Lu on 6/12/16.
* E-mail: gang.lu.ict@gmail.com
*
* Copyright (c) 2016 bafst.com, All rights reserved.
*/
"use strict";
var settings = require("./settings_MQ"); //配置信息
var logger = settings.logger;
var moment = require('moment'); //时间
var java = require("java");
var DefaultMQProducer = java.import('com.alibaba.rocketmq.client.producer.DefaultMQProducer');
//需要注意的是node是支持同步/异步的,对于调用java内的函数,要根据情况选择是同步还是异步调用方式!
/**
* MQProducer
* @param {String} groupName group name
* @param {String} namesrvAddr addresses of the name servers
* @constructor
*/
var MQProducer = function (groupName, namesrvAddr){
this.producer = undefined; //初始化放在了init函数中
this.groupName = groupName;
this.namesrvAddr = namesrvAddr;
this.instanceName = moment().format("x"); //毫秒值作为instance name,默认返回string
this.compressMsgBodyOverHowmuch = 4096; //消息压缩阈值
};
/**
* init
* 批量设置一些基本项(为了尽可能少实现这些API接口,如以后有需要,可以逐个移出init)
* @param {Function} callback the callback function
*/
MQProducer.prototype.init = function(callback) {
var self = this;
logger.info('Initializing producer ' + self.instanceName + ' ...');
self.producer = new DefaultMQProducer(self.groupName); //创建实例
self.producer.setNamesrvAddrSync(self.namesrvAddr);
self.producer.setInstanceNameSync(self.instanceName);
self.producer.setCompressMsgBodyOverHowmuchSync(parseInt(self.compressMsgBodyOverHowmuch));
callback && callback();
};
/**
* start
* 批量设置一些基本项(为了尽可能少实现这些API接口,如以后有需要,可以逐个移出init)
* @param {Function} callback the callback function
*/
MQProducer.prototype.start = function () {
var self = this;
logger.info('Starting producer ' + self.instanceName + ' ...');
// 同步调用start
self.producer.startSync(); // start returns void
};
/**
* shutdown
* @param {Function} callback the callback function
*/
MQProducer.prototype.shutdown = function() {
var self = this;
logger.info('Shutting down producer ' + self.instanceName + ' ...');
// 同步调用stop
self.producer.shutdownSync() // stop returns void
};
/**
* send
* @param {Object} MQMsg the message to send
* @param {Function} callback the callback function
*/
MQProducer.prototype.send = function(MQMsg, callback) {
var self = this;
logger.debug('Producer ' + self.instanceName + ' sending message: ' + MQMsg.tostr());
// 异步调用模式下,callback的第一个参数是异常,第二个参数是调用的返回值
self.producer.send(MQMsg.msg, function(err, result){ // send returns SendResult
if(err) {
logger.error('Sending message failed! Please look up the exception reported!');
logger.error(err);
} else {
logger.debug('Sending message successfully from instance ' + self.instanceName + ' .');
}
callback && callback(result);
});
};
module.exports = MQProducer;