-
Notifications
You must be signed in to change notification settings - Fork 2
Expand file tree
/
Copy pathMQMessage.js
More file actions
131 lines (116 loc) · 5.7 KB
/
Copy pathMQMessage.js
File metadata and controls
131 lines (116 loc) · 5.7 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
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
/**
* 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 iconv = require('iconv-lite'); //char encoding/decoding
var settings = require("./settings_MQ"); //配置信息
var logger = settings.logger;
var java = require("java");
var Message = java.import('com.alibaba.rocketmq.common.message.Message');
// enum classes:
var SENDSTATUS = java.import('com.alibaba.rocketmq.client.producer.SendStatus');
var PULLSTATUS = java.import('com.alibaba.rocketmq.client.consumer.PullStatus');
var CONSUMEFROMWHERE = java.import('com.alibaba.rocketmq.common.consumer.ConsumeFromWhere');
var CONSUMECONCURRENTLYSTATUS = java.import('com.alibaba.rocketmq.client.consumer.listener.ConsumeConcurrentlyStatus');
var CONSUMEORDERLYSTATUS = java.import('com.alibaba.rocketmq.client.consumer.listener.ConsumeOrderlyStatus');
var MESSAGEMODEL = java.import('com.alibaba.rocketmq.common.protocol.heartbeat.MessageModel');
// Class MQMessage
var MQMessage = function(topic, tags, keys, body){
this.topic = topic;
this.tags = tags;
this.keys = keys;
this.body = body;
var buffer = iconv.encode(this.body, settings.MsgBodyEncoding); //string to encoded Buffer
//not buffer.toJSON().map(.....), but buffer.toJSON().data.map(......)
var byteArray = java.newArray("byte", buffer.toJSON().data.map(function(c) { return java.newByte(c); })); //从Buffer生成byte array的固定方式! 要先转换为Byte数组,也就是toJSON
//var byteArray = java.newArray("byte", this.body.split('').map(function(c) { return java.newByte(String.prototype.charCodeAt(c)); })); //从string生成byte array的固定方式!
//sync methods:
this.msg = new Message(this.topic, this.tags, this.keys, byteArray); //string to bytes
//this.msg = java.newInstanceSync("com.alibaba.rocketmq.common.message.Message", this.topic, this.tags, this.keys, byteArray); //string to bytes
//async method: java.newInstance(className, [args...], callback);
};
MQMessage.prototype.tostr = function(){
return this.topic + "::" + this.tags + "::" + this.keys + "::" + this.body;
};
exports.MQMessage = MQMessage;
// enum classes
// PullResult的返回结果
var PullStatus = {
//'FOUND': 0, // Founded
'FOUND': PULLSTATUS.FOUND,
//'NO_NEW_MSG': 1, // No new message can be pull
'NO_NEW_MSG': PULLSTATUS.NO_NEW_MSG,
//'NO_MATCHED_MSG': 2, // Filtering results can not match
'NO_MATCHED_MSG': PULLSTATUS.NO_MATCHED_MSG,
//'OFFSET_ILLEGAL': 3 // Illegal offset,may be too big or too small
'OFFSET_ILLEGAL': PULLSTATUS.OFFSET_ILLEGAL
};
exports.PullStatus = PullStatus;
// PullResult的返回结果
var SendStatus = {
//'SEND_OK': 0,
'SEND_OK': SENDSTATUS.SEND_OK,
//'FLUSH_DISK_TIMEOUT': 1,
'FLUSH_DISK_TIMEOUT': SENDSTATUS.FLUSH_DISK_TIMEOUT,
//'FLUSH_SLAVE_TIMEOUT': 2,
'FLUSH_SLAVE_TIMEOUT': SENDSTATUS.FLUSH_SLAVE_TIMEOUT,
//'SLAVE_NOT_AVAILABLE': 3
'SLAVE_NOT_AVAILABLE': SENDSTATUS.SLAVE_NOT_AVAILABLE
};
exports.SendStatus = SendStatus;
// PushConsumer消费时选择第一次订阅时的消费位置
var ConsumeFromWhere = {
// 一个新的订阅组第一次启动从队列的最后位置开始消费
// 后续再启动接着上次消费的进度开始消费
//'CONSUME_FROM_LAST_OFFSET': 0,
'CONSUME_FROM_LAST_OFFSET': CONSUMEFROMWHERE.CONSUME_FROM_LAST_OFFSET,
//@Deprecated
//'CONSUME_FROM_LAST_OFFSET_AND_FROM_MIN_WHEN_BOOT_FIRST': 1,
'CONSUME_FROM_LAST_OFFSET_AND_FROM_MIN_WHEN_BOOT_FIRST': CONSUMEFROMWHERE.CONSUME_FROM_LAST_OFFSET_AND_FROM_MIN_WHEN_BOOT_FIRST,
//@Deprecated
//'CONSUME_FROM_MIN_OFFSET': 2,
'CONSUME_FROM_MIN_OFFSET': CONSUMEFROMWHERE.CONSUME_FROM_MIN_OFFSET,
//@Deprecated
//'CONSUME_FROM_MAX_OFFSET': 3,
'CONSUME_FROM_MAX_OFFSET': CONSUMEFROMWHERE.CONSUME_FROM_MAX_OFFSET,
// 一个新的订阅组第一次启动从队列的最前位置开始消费<br>
// 后续再启动接着上次消费的进度开始消费
//'CONSUME_FROM_FIRST_OFFSET': 4,
'CONSUME_FROM_FIRST_OFFSET': CONSUMEFROMWHERE.CONSUME_FROM_FIRST_OFFSET,
// 一个新的订阅组第一次启动从指定时间点开始消费,时间点设置参见DefaultMQPushConsumer.consumeTimestamp参数
// 后续再启动接着上次消费的进度开始消费
//'CONSUME_FROM_TIMESTAMP': 5,
'CONSUME_FROM_TIMESTAMP': CONSUMEFROMWHERE.CONSUME_FROM_TIMESTAMP
};
exports.ConsumeFromWhere = ConsumeFromWhere;
// PushConsumer消费后的返回值(并发消费时)
var ConsumeConcurrentlyStatus = {
//'CONSUME_SUCCESS': 0, // Success consumption
'CONSUME_SUCCESS': CONSUMECONCURRENTLYSTATUS.CONSUME_SUCCESS,
//'RECONSUME_LATER': 1, // Failure consumption,later try to consume
'RECONSUME_LATER': CONSUMECONCURRENTLYSTATUS.RECONSUME_LATER
};
exports.ConsumeConcurrentlyStatus = ConsumeConcurrentlyStatus;
// PushConsumer消费后的返回值(顺序消费时)
var ConsumeOrderlyStatus ={
//'SUCCESS': 0, // Success consumption
'SUCCESS': CONSUMEORDERLYSTATUS.SUCCESS,
//'ROLLBACK': 1, // Rollback consumption(only for binlog consumption)
'ROLLBACK': CONSUMEORDERLYSTATUS.ROLLBACK,
//'COMMIT': 2, // Commit offset(only for binlog consumption)
'COMMIT': CONSUMEORDERLYSTATUS.COMMIT,
//'SUSPEND_CURRENT_QUEUE_A_MOMENT': 3 // Suspend current queue a moment
'SUSPEND_CURRENT_QUEUE_A_MOMENT': CONSUMEORDERLYSTATUS.SUSPEND_CURRENT_QUEUE_A_MOMENT
};
exports.ConsumeOrderlyStatus = ConsumeOrderlyStatus;
// PushConsumer的消息model
var MessageModel = {
//'BROADCASTING': 0, // broadcast
'BROADCASTING': MESSAGEMODEL.BROADCASTING,
//'CLUSTERING': 1 // clustering
'CLUSTERING': MESSAGEMODEL.CLUSTERING
};
exports.MessageModel = MessageModel;