基于Redis的高性能的延迟队列。不会每次去循环Redis即可判断是否有数据进入,极大减轻了redis的压力。 用法: 当有消息达到 delay time的时候会触发peekCallback。
publicclassTestRedisDelayQueue {
JedisClusterjedisCluster = null;
RedisDelayQueuequeue = null;
@Beforepublicvoidinit() {
Stringip = "192.168.2.160";
Set<HostAndPort> nodes = newHashSet<>();
nodes.add(newHostAndPort(ip, 7701));
nodes.add(newHostAndPort(ip, 7702));
nodes.add(newHostAndPort(ip, 7703));
nodes.add(newHostAndPort(ip, 7704));
nodes.add(newHostAndPort(ip, 7705));
nodes.add(newHostAndPort(ip, 7706));
JedisPoolConfigpool = newJedisPoolConfig();
pool.setMaxTotal(100);
pool.setFairness(false);
pool.setNumTestsPerEvictionRun(100);
pool.setMaxWaitMillis(5000);
pool.setTestOnBorrow(true);
jedisCluster = newJedisCluster(nodes, 1000, 1000, 100, null, pool); // maxAttempt必须调大jedisCluster.set("test", "test");
queue = newRedisDelayQueue("com.meipian", "delayqueue", jedisCluster, 60 * 1000,
newDelayQueueProcessListener() {
@OverridepublicvoidpushCallback(Messagemessage) {
}
@OverridepublicvoidpeekCallback(Messagemessage) {
System.out.println("message----->" + message);
queue.ack(message.getId());//确认操作。将会删除消息
}
@OverridepublicvoidackCallback(Messagemessage) {
}
});
}
@TestpublicvoidtestCreate() throwsInterruptedException {
Messagemessage = newMessage();
for (inti = 0; i < 10; i++) {
message.setId(i + "");
message.setPayload("test");
message.setPriority(0);
message.setTimeout(3000);
queue.push(message);
}
// message = queue.peek();// queue.ack("1234");queue.listen();
}
}