-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathredis-mq.js
More file actions
448 lines (408 loc) · 15.8 KB
/
redis-mq.js
File metadata and controls
448 lines (408 loc) · 15.8 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
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
'use strict';
const redis = require('./redis');
const CronJob = require('cron').CronJob;
const async = require('async');
const pullMessgae = require('./redis-pop');
const MQ_NAME = 'test-redis-mq';
const MQ_HASH_NAME = 'test-redis-hash';
const MQ_HASH_RETRY_TIMES = 'test-redis-retry-hash';
let MQ_FETCH_STATUS = false;
const maxAckTime = 8 * 1000; // 8s 最长消费时间
const maxAckTimeout = 60 * 1000; // 消费超时时间
// const retryTime = 20; // 重试时间 20s
const failureRate = 0.5; // 失败率
// const successRate = 0.4; // 成功率
const timeoutRate = 0.1; // 超时率
const eachMessageCount = 500; // 每次 Message 获取数量
const maxRetryTimes = 5; // 最大重试次数
// const maxRedisMessageCount = 500;
const minRedisMessageCount = 400;
// const eachPullCount = 200;
function messageCount(callback) {
callback = callback || function(err, count) {
if(err) {
console.error(`[获取 Redis 长度] ERROR: ${err.message}`);
} else {
console.log(`[获取 Redis 长度] 当前队列长度为: ${count}`);
}
}
return redis.llen(MQ_NAME, function(err, count) {
if(err) {
return callback(err);
}
callback(undefined, count);
});
}
function returnRate() {
return 'success';
const tmp = Math.random();
console.log(`随机 tmp: ${tmp}`);
if(tmp < timeoutRate) {
return 'timeout';
} else if(tmp - timeoutRate < failureRate) {
return 'failure';
} else {
return 'success';
}
}
function pullMessageFromKafuka(callback) {
if(MQ_FETCH_STATUS) {
console.log(`[fetch message] 但是已经 fetch 被堵塞`);
setTimeout(function() {
callback();
});
}
console.log(`[fetch message] 开始 fetch message`);
MQ_FETCH_STATUS = true;
pullMessgae(function(err) {
MQ_FETCH_STATUS = false;
if(err) {
console.log(`[fetch message] fetch message 失败 error: ${err.message}`);
return callback(err);
}
console.log(`[fetch message] fetch message 结束`);
messageCount(function(err, count) {
console.log(`[fetch message] 最后数量为: ${count}`);
callback();
});
});
}
function fetch(count, callback) {
let mqCount = 0;
if(count < 0) {
count = 1;
}
async.series([
function(callback) {
// 检查 redis 的队列长度
messageCount(function(err, count) {
if(err) {
return callback(err);
}
console.log(`[获取 Redis 长度] 当前队列长度为: ${count}`);
if(count > 0) {
mqCount = count;
}
return callback();
});
},
function(callback) {
console.log(`[检查 Redis 数量] 当前数量为: ${mqCount}, min 数量: ${minRedisMessageCount}`);
if(mqCount > minRedisMessageCount) {
console.log(`[检查 Redis 数量] 发现数量已经足够`);
return callback();
}
// 执行获取消息
console.log(`[检查 Redis 数量] 发现数量不足,需要重新获取`);
pullMessageFromKafuka(callback);
},
function(callback) {
console.log(`[fetch message key] 开始获取消息数量为: ${count}`);
const array = new Array(count);
const result = [];
async.eachLimit(array, 10, function(_a, callback) {
redis.lpop(MQ_NAME, function(err, key) {
if(err) {
console.log(`[fetch message key] 获取 key 失败 ${err.message}`);
return callback();
}
console.log(`[fetch message key] 成功 key: ${key}`) ;
redis.hset(MQ_HASH_NAME, key, Date.now(), function(err) {
console.log(`[redis hash set] 设置 hash id: ${key} 相关时间`);
if(err) {
console.log(`[redis hash set] 设置 hash id: ${key} 相关时间失败`);
console.log(`[redis hash set] 重新 push 消息`);
redis.lpush(MQ_NAME, key);
return callback(err);
}
console.log(`[fetch message]开始 fetch message`)
redis.get(key, function(err, _result) {
if(err) {
console.log(`[fetch message] fetch message 失败, key: ${key} error: ${err.messahe}`);
redis.rpush(MQ_NAME, key);
redis.hset(MQ_HASH_NAME, key, null);
return callback(err);
}
console.log(`[fetch message] fetch message id: ${key} 成功,消息体为: ${_result}`);
result.push({
key: key,
data: _result
});
callback();
});
});
});
}, function(err) {
if(err) {
return callback(err);
}
callback(undefined, result);
});
}
], function(err, result) {
if(err) {
return callback(err);
}
callback(undefined, result[result.length - 1]);
});
}
// 检查超时的消息
// 重置已经过
// fetch(20, function(err, result) {
// if(err) {
// console.error(err);
// }
// console.log('获取消息成功')
// console.log(result);
// });
function incrFailureTimes(key, callback) {
console.log(`[retry times] 添加 & 查询 key: ${key}`);
redis.hincrby(MQ_HASH_RETRY_TIMES, key, 1, function(err, count) {
console.log(`[retry times] 查询到 key: ${key} 重试之后的次数为: ${count}`);
if(err) {
return callback(err);
}
if(count > maxRetryTimes) {
console.log(`[retry times] key: ${key} 超时次数为 ${count} 超过预定次数 ${maxRetryTimes}`);
redis.hdel(MQ_HASH_NAME, key);
return callback();
// todo 保存到 MY_SQL
} else {
console.log(`[retry times] 将 key: ${key} 重新加入到队列,并清除原来的时间计数 `);
redis.rpush(MQ_NAME, key, function(err) {
if(err) {
return callback();
}
console.log(`[retry times] key: ${key} 重新加入队列成功 `);
redis.hset(MQ_HASH_NAME, key, null, function(err) {
console.log(`[retry times] key: ${key} 重新初始化时间成功 `);
callback();
});
});
}
});
}
let checkStatus = true;
function checkExpireMessage() {
// 1. list 有 hash null => 未消费
// 2. list 有 hash now => 重试消费
// 3. list 无 hash null => 数据丢失 => 重新 push 数据到 list
// 4. list 无 hash now => 消费中 or 消费超时
// 先 list 后 hash
// 如果 3:
// 检查 hash 是否存在,如果存在,那就是铁定的数据丢失,重新导入数据
// 如果 4:
// 先检查 hash 是否存在,如果存在,就计算时间是否超时,如果超时就记一次超时,重新加入到队列
let mgCount;
let mqMessages;
let hashMap;
if(checkStatus) {
checkStatus = false;
} else {
return;
}
console.log(`[message check] 开始检查消息过期问题`);
async.series([
function(callback) {
console.log(`[message check] 检查队列长度`);
messageCount(function(err, count) {
if(err) {
return callback(err);
}
mgCount = count;
console.log(`[message check] 队列长度为: ${count}`);
callback();
});
},
function(callback) {
const count = mgCount + 200
console.log(`[message check] 加上默认 pull 数量, 获取总数为: ${count}`);
redis.lrange(MQ_NAME, 0, count, function(err, data) {
if(err) {
return callback(err);
}
mqMessages = data;
console.log(`[message check] 获取消息的数量为: ${mqMessages.length}`);
callback();
});
},
function(callback) {
console.log(`[message check] 准备获取 HASH MAP`);
redis.hgetall(MQ_HASH_NAME, function(err, data) {
if(err) {
return callback(err);
}
hashMap = data;
callback();
});
}
], function(err) {
if(err) {
console.error(`[message check] 检查消息失败: ${err.message}`);
return;
}
// console.log(mqMessages);
// console.log(hashMap);
const missingList = [];
const timeoutList = [];
const keys = Object.keys(hashMap);
console.log(`[message check] 消息 HASH 内的消息总数量为 ${keys.length}`);
for(const key of keys) {
const index = mqMessages.indexOf(key)
if(index < 0 && !hashMap[key]) {
console.log(index, hashMap[key]);
// 数据缺失
console.log(`[message check] 检查发现 ${key} 这个数据缺失,需要重新 push 到 mq`)
// redis.rpush(key);
missingList.push(key);
} else if(index < 0 && hashMap[key] && (Date.now() - hashMap[key]) > maxAckTimeout) {
// 数据 ack 超时
console.log(`[message check] 检查发现 ${key} 超时, 需要重新 push 到 mq,并及加次数`);
timeoutList.push(key);
}
}
console.log(`[message check] 最后数据缺失的 list 长度为: ${missingList.length}`);
console.log(`[message check] 最后数据超时的 list 长度为: ${timeoutList.length}`);
async.series([
function(callback) {
console.log(`[message check] 修复 数据缺失`);
async.eachLimit(missingList, 10, function(key, callback) {
redis.hexists(MQ_HASH_NAME, key, function(err, data) {
if(err) {
console.log(err);
return callback();
}
if(!data) {
console.log(`[message check] 确定数据是因为延迟问题`);
return callback();
}
redis.rpush(MQ_NAME, key, function(err) {
if(err) {
console.log(err);
}
callback();
});
});
}, callback);
},
function(callback) {
console.log(`[messsage check] 修复数据超时问题`);
async.eachLimit(timeoutList, 10, function(key, callback) {
// 重新确认
redis.hget(MQ_HASH_NAME, key, function(err, data) {
if(err) {
console.log(err);
return callback();
}
if(!data) {
console.log(`[message check] 确定数据是因为延迟问题`);
return callback();
}
incrFailureTimes(key, function(err) {
if(err) {
console.log(err);
}
callback();
});
});
}, callback);
}
], function(err) {
console.log(`[message check] 修复数据成功`);
checkStatus = true;
})
});
}
function ackMessage(key, success, callback) {
console.log(`[ack message] key 为 ${key} 且状态为: ${success ? 'success' : 'false'}`);
async.series([
function(callback) {
redis.hget(MQ_HASH_NAME, key, function(err, time) {
if(!time) {
console.log(`[ack message] key: ${key} 已经过期了`);
return callback(new Error('key 已被过期消费'));
}
if(success) {
console.log(`[ack message] 消费成功删除 key: ${key}`);
redis.hdel(MQ_HASH_NAME, key, function(err) {
if(err) {
return callback(err);
}
redis.hdel(MQ_HASH_RETRY_TIMES, key);
redis.del(key);
callback();
});
} else {
console.log(`[ack message] 消费失败重试 key: ${key}`);
incrFailureTimes(key, function(err) {
if(err) {
console.log(err);
}
callback();
});
}
});
}
], function(err) {
if(err) {
return callback(err);
}
callback();
});
};
function clientDone() {
console.log(`[client] 我是一个快乐的客户端`);
let datas;
async.series([
function(callback) {
console.log(`[client] fetch ${eachMessageCount} 条数据`);
fetch(eachMessageCount, function(err, result) {
if(err) {
return callback(err);
}
datas = result;
console.log(`[client] 获得数据的数量为: ${datas.length}`);
return callback();
});
},
function(callback) {
async.eachSeries(datas, function(_data, callback) {
const { key, data } = _data;
console.log(`[client] 开始处理数据 key: ${key}` );
const res = returnRate();
if(res === 'failure') {
console.log(`[client] 抽到了失败 key: ${key}` );
setTimeout(() => {
ackMessage(key, false, callback);
}, 100);
} else if (res === 'success') {
console.log(`[client] 抽到了成功 key: ${key}` );
setTimeout(() => {
ackMessage(key, true, callback);
}, 100);
} else if(res === 'timeout') {
console.log(`[client] 抽到了超时 key: ${key}` );
setTimeout(() => {
callback();
}, 100);
} else {
throw new Error(`抽到了奇怪的东西: ${res}`);
}
}, function() {
console.log(`[client] 所有消息都处理完了, 棒棒哒💯`);
callback();
});
},
], function() {
console.log(`[client] 快乐的客户端完成了一次使命`);
console.log(`[client] 重新开始新的使命`);
setTimeout(() => {
clientDone();
}, 100);
})
}
clientDone();
// new CronJob('*/6 * * * * *', function() {
// console.log('=================== 快乐的修复数据 ====================')
// checkExpireMessage();
// }, null, true);