From 72a9d9597f36e4b7fa937cb164bf6158d7ee2147 Mon Sep 17 00:00:00 2001 From: Sonu Kumar Date: Sat, 13 May 2023 15:21:27 +0530 Subject: [PATCH 1/9] Fixing #193. --- .../rqueue/config/RqueueSchedulerConfig.java | 3 + .../sonus21/rqueue/core/MessageScheduler.java | 215 +++++++++++------- .../rqueue/core/ScheduledTaskDetail.java | 57 ----- .../rqueue/example/MessageListener.java | 6 +- .../src/main/resources/logback.xml | 14 +- 5 files changed, 148 insertions(+), 147 deletions(-) delete mode 100644 rqueue-core/src/main/java/com/github/sonus21/rqueue/core/ScheduledTaskDetail.java diff --git a/rqueue-core/src/main/java/com/github/sonus21/rqueue/config/RqueueSchedulerConfig.java b/rqueue-core/src/main/java/com/github/sonus21/rqueue/config/RqueueSchedulerConfig.java index d062e65e8..db9b2723f 100644 --- a/rqueue-core/src/main/java/com/github/sonus21/rqueue/config/RqueueSchedulerConfig.java +++ b/rqueue-core/src/main/java/com/github/sonus21/rqueue/config/RqueueSchedulerConfig.java @@ -72,4 +72,7 @@ public class RqueueSchedulerConfig { // Maximum delay for message mover task due to failure @Value("${rqueue.scheduler.max.message.mover.delay:60000}") private long maxMessageMoverDelay; + + @Value("${rqueue.scheduler.max.message.count:100}") + private long maxMessageCount; } diff --git a/rqueue-core/src/main/java/com/github/sonus21/rqueue/core/MessageScheduler.java b/rqueue-core/src/main/java/com/github/sonus21/rqueue/core/MessageScheduler.java index e6fe94b73..b27b25b75 100644 --- a/rqueue-core/src/main/java/com/github/sonus21/rqueue/core/MessageScheduler.java +++ b/rqueue-core/src/main/java/com/github/sonus21/rqueue/core/MessageScheduler.java @@ -16,7 +16,6 @@ package com.github.sonus21.rqueue.core; -import static com.github.sonus21.rqueue.utils.Constants.MAX_MESSAGES; import static com.github.sonus21.rqueue.utils.Constants.MIN_DELAY; import static java.lang.Math.min; @@ -31,10 +30,12 @@ import java.util.Arrays; import java.util.List; import java.util.Map; +import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.Future; -import lombok.AllArgsConstructor; -import lombok.ToString; +import java.util.concurrent.ScheduledFuture; +import java.util.concurrent.TimeUnit; +import lombok.Getter; import org.slf4j.Logger; import org.springframework.beans.factory.DisposableBean; import org.springframework.beans.factory.annotation.Autowired; @@ -52,8 +53,8 @@ import org.springframework.util.Assert; import org.springframework.util.CollectionUtils; -public abstract class MessageScheduler - implements DisposableBean, ApplicationListener { +public abstract class MessageScheduler implements DisposableBean, + ApplicationListener { private final Object monitor = new Object(); @Autowired @@ -102,8 +103,8 @@ private void subscribeToRedisTopic(String queueName) { if (isRedisEnabled()) { String channelName = getChannelName(queueName); getLogger().debug("Queue {} subscribe to channel {}", queueName, channelName); - this.rqueueRedisListenerContainerFactory.addMessageListener( - messageSchedulerListener, new ChannelTopic(channelName)); + this.rqueueRedisListenerContainerFactory.addMessageListener(messageSchedulerListener, + new ChannelTopic(channelName)); channelNameToQueueName.put(channelName, queueName); } } @@ -114,8 +115,7 @@ private void startQueue(String queueName) { } queueRunningState.put(queueName, true); if (scheduleTaskAtStartup() || !isRedisEnabled()) { - long scheduleAt = getQueueStartTime(); - schedule(queueName, scheduleAt, false); + schedule(queueName, getQueueStartTime(), false); } subscribeToRedisTopic(queueName); } @@ -152,8 +152,7 @@ private void waitForRunningQueuesToStop() { } private void stopQueue(String queueName) { - Assert.isTrue( - queueRunningState.containsKey(queueName), + Assert.isTrue(queueRunningState.containsKey(queueName), "Queue with name '" + queueName + "' does not exist"); queueRunningState.put(queueName, false); } @@ -262,42 +261,25 @@ protected long getMinDelay() { private class QueueScheduler { - private void updateLastScheduleTime(String queueName, long time) { - queueNameToLastMessageScheduleTime.put(queueName, time); + private void scheduleNewTask(QueueDetail queueDetail, String zsetName, long startTime) { + MessageMoverTask task = new MessageMoverTask(queueDetail, zsetName); + Future future = scheduler.schedule(task, Instant.ofEpochMilli(startTime)); + addTask(task, new ScheduledTaskDetail(task.id, future, startTime)); } - private void scheduleNewTask( - QueueDetail queueDetail, String queueName, String zsetName, long startTime) { - MessageMoverTask timerTask = - new MessageMoverTask( - queueDetail.getName(), - queueDetail.getQueueName(), - zsetName, - isProcessingQueue(zsetName)); - Future future = - scheduler.schedule( - timerTask, Instant.ofEpochMilli(getNextScheduleTime(queueName, startTime))); - addTask(timerTask, new ScheduledTaskDetail(startTime, future)); - } - - private void scheduleTask( - long startTime, long currentTime, QueueDetail queueDetail, String zsetName) { + private void scheduleTask(QueueDetail queueDetail, String zsetName, long startTime, + long currentTime) { long requiredDelay = Math.max(1, startTime - currentTime); - long taskStartTime = startTime; - MessageMoverTask timerTask = - new MessageMoverTask( - queueDetail.getName(), - queueDetail.getQueueName(), - zsetName, - isProcessingQueue(queueDetail.getName())); + long taskStartTime = currentTime; + MessageMoverTask task = new MessageMoverTask(queueDetail, zsetName); Future future; if (requiredDelay < MIN_DELAY) { - future = scheduler.submit(timerTask); - taskStartTime = currentTime; + future = scheduler.submit(task); } else { - future = scheduler.schedule(timerTask, Instant.ofEpochMilli(currentTime + requiredDelay)); + taskStartTime = currentTime + requiredDelay; + future = scheduler.schedule(task, Instant.ofEpochMilli(taskStartTime)); } - addTask(timerTask, new ScheduledTaskDetail(taskStartTime, future)); + addTask(task, new ScheduledTaskDetail(task.id, future, taskStartTime)); } private boolean shouldNotSchedule(String queueName, boolean forceSchedule) { @@ -311,62 +293,84 @@ private boolean shouldNotSchedule(String queueName, boolean forceSchedule) { return !forceSchedule && currentTime - lastSeenTime < getMinDelay(); } + private void handleTaskOverride(ScheduledTaskDetail scheduledTaskDetail, + QueueDetail queueDetail, String zsetName, long startTime) { + // we should not schedule too frequent calls + long difference = startTime - scheduledTaskDetail.getStartTime(); + if (difference < getMinDelay()) { + return; + } + long currentTime = System.currentTimeMillis(); + if (cancelExistingTask(scheduledTaskDetail, currentTime, queueDetail, zsetName)) { + scheduleNewTask(queueDetail, zsetName, startTime); + } + } + protected synchronized void schedule(String queueName, Long startTime, boolean forceSchedule) { + getLogger().debug("Schedule Task queue={}, force={}", queueName, forceSchedule); if (shouldNotSchedule(queueName, forceSchedule)) { return; } long currentTime = System.currentTimeMillis(); - updateLastScheduleTime(queueName, currentTime); ScheduledTaskDetail scheduledTaskDetail = getScheduledTask(queueName); QueueDetail queueDetail = EndpointRegistry.get(queueName); String zsetName = getZsetName(queueName); + // no task was scheduled or call came from existing MessageMoverTask if (scheduledTaskDetail == null || forceSchedule) { - scheduleTask(startTime, currentTime, queueDetail, zsetName); - return; + scheduleTask(queueDetail, zsetName, startTime, currentTime); + } else { + handleTaskOverride(scheduledTaskDetail, queueDetail, zsetName, startTime); } - checkExistingTask(scheduledTaskDetail, currentTime, queueDetail, zsetName); - scheduleNewTask(queueDetail, queueName, zsetName, startTime); } - private void addTask(MessageMoverTask timerTask, ScheduledTaskDetail scheduledTaskDetail) { - getLogger().debug("Timer: {}, Task: {}", timerTask, scheduledTaskDetail); - queueNameToScheduledTask.put(timerTask.getName(), scheduledTaskDetail); + private void addTask(MessageMoverTask task, ScheduledTaskDetail scheduledTaskDetail) { + getLogger().debug("Adding Task task={}, startTime={}", task, + scheduledTaskDetail.getStartTime()); + queueNameToLastMessageScheduleTime.put(task.getName(), System.currentTimeMillis()); + queueNameToScheduledTask.put(task.getName(), scheduledTaskDetail); } - private void checkExistingTask( - ScheduledTaskDetail scheduledTaskDetail, - long currentTime, - QueueDetail queueDetail, - String zsetName) { - // run existing tasks continue - long existingDelay = scheduledTaskDetail.getStartTime() - currentTime; + private boolean cancelExistingTask(ScheduledTaskDetail scheduledTaskDetail, long currentTime, + QueueDetail queueDetail, String zsetName) { Future submittedTask = scheduledTaskDetail.getFuture(); boolean completedOrCancelled = submittedTask.isDone() || submittedTask.isCancelled(); - // tasks older than TASK_ALIVE_TIME are considered dead - if (!completedOrCancelled - && existingDelay < MIN_DELAY - && existingDelay > Constants.TASK_ALIVE_TIME) { - ThreadUtils.waitForTermination( - getLogger(), - submittedTask, - Constants.DEFAULT_SCRIPT_EXECUTION_TIME, - "LIST: {} ZSET: {}, Task: {} failed", - queueDetail.getQueueName(), - zsetName, - scheduledTaskDetail); + // this can happen when task was not scheduled due to some exception in scheduling + if (completedOrCancelled) { + return false; + } + // run existing tasks continue + long existingDelay = scheduledTaskDetail.getStartTime() - currentTime; + // task is about to run or was scheduled to run in last alive time, but could not so run it + if (existingDelay < MIN_DELAY && existingDelay > Constants.TASK_ALIVE_TIME) { + ThreadUtils.waitForTermination(getLogger(), submittedTask, + Constants.DEFAULT_SCRIPT_EXECUTION_TIME, "LIST: {} ZSET: {}, Task: {} failed", + queueDetail.getQueueName(), zsetName, scheduledTaskDetail); + return false; + } + boolean cancelled = submittedTask.cancel(false); + if (cancelled) { + getLogger().debug("Task {} cancelled", scheduledTaskDetail.getId()); } + return cancelled; } } - @ToString - @AllArgsConstructor private class MessageMoverTask implements Runnable { + private final String id; private final String name; private final String queueName; private final String zsetName; private final boolean processingQueue; + MessageMoverTask(QueueDetail queueDetail, String zsetName) { + this.id = UUID.randomUUID().toString(); + this.name = queueDetail.getName(); + this.queueName = queueDetail.getQueueName(); + this.zsetName = zsetName; + this.processingQueue = isProcessingQueue(zsetName); + } + private long getNextScheduleTimeInternal(Long value, Exception e) { int errCount = 0; long nextTime; @@ -375,7 +379,7 @@ private long getNextScheduleTimeInternal(Long value, Exception e) { if (errCount % 3 == 0) { getLogger().error("Message mover task is failing continuously queue: {}", name, e); } - long delay = (long) (100 * Math.pow(1.5, errCount)); + long delay = (long) (MIN_DELAY * Math.pow(1.5, errCount)); delay = Math.min(delay, rqueueSchedulerConfig.getMaxMessageMoverDelay()); nextTime = System.currentTimeMillis() + delay; } else { @@ -385,6 +389,24 @@ private long getNextScheduleTimeInternal(Long value, Exception e) { return nextTime; } + @Override + public String toString() { + return String.format("MessageMoverTask(id=%s, queue=%s)", id, name); + } + + private long getMessageCount() { + return rqueueSchedulerConfig.getMaxMessageCount(); + } + + private List scriptKeys() { + return Arrays.asList(queueName, zsetName); + } + + private Object[] scriptArgs() { + long currentTime = System.currentTimeMillis(); + return new Object[]{currentTime, getMessageCount(), processingQueue ? 1 : 0}; + } + @Override public void run() { getLogger().debug("Running {}", this); @@ -392,14 +414,7 @@ public void run() { Exception e = null; try { if (isQueueActive(name)) { - long currentTime = System.currentTimeMillis(); - value = - defaultScriptExecutor.execute( - redisScript, - Arrays.asList(queueName, zsetName), - currentTime, - MAX_MESSAGES, - processingQueue ? 1 : 0); + value = defaultScriptExecutor.execute(redisScript, scriptKeys(), scriptArgs()); } } catch (RedisSystemException ex) { e = ex; @@ -407,10 +422,8 @@ public void run() { e = ex; getLogger().warn("Task execution failed for the queue: {}", getName(), e); } finally { - if (isQueueActive(name)) { - long nextExecutionTime = getNextScheduleTimeInternal(value, e); - schedule(name, nextExecutionTime, true); - } + long nextExecutionTime = getNextScheduleTimeInternal(value, e); + schedule(name, nextExecutionTime, true); } } @@ -427,7 +440,8 @@ private void handleMessage(String queueName, Long startTime) { if (currentTime - lastSeenTime < getMinDelay()) { return; } - schedule(queueName, startTime, false); + long jobStartTime = getNextScheduleTime(queueName, startTime); + schedule(queueName, jobStartTime, false); } @Override @@ -451,4 +465,39 @@ public void onMessage(Message message, byte[] pattern) { } } } + + @Getter + private static class ScheduledTaskDetail { + + private final Future future; + private final long startTime; + private final String id; + + ScheduledTaskDetail(String id, Future future, long startTime) { + this.startTime = startTime; + this.future = future; + this.id = id; + } + + @Override + public String toString() { + StringBuilder sb = new StringBuilder(); + if (future instanceof ScheduledFuture) { + sb.append("ScheduledFuture(delay="); + sb.append(((ScheduledFuture) future).getDelay(TimeUnit.MILLISECONDS)); + sb.append("Ms, "); + } else { + sb.append("Future("); + } + sb.append("id="); + sb.append(id); + sb.append(", startTime="); + sb.append(startTime); + sb.append(", currentTime="); + sb.append(System.currentTimeMillis()); + sb.append(")"); + return sb.toString(); + } + } + } diff --git a/rqueue-core/src/main/java/com/github/sonus21/rqueue/core/ScheduledTaskDetail.java b/rqueue-core/src/main/java/com/github/sonus21/rqueue/core/ScheduledTaskDetail.java deleted file mode 100644 index 32e7057f7..000000000 --- a/rqueue-core/src/main/java/com/github/sonus21/rqueue/core/ScheduledTaskDetail.java +++ /dev/null @@ -1,57 +0,0 @@ -/* - * Copyright (c) 2019-2023 Sonu Kumar - * - * Licensed under the Apache License, Version 2.0 (the "License"); - * You may not use this file except in compliance with the License. - * You may obtain a copy of the License at - * - * https://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and limitations under the License. - * - */ - -package com.github.sonus21.rqueue.core; - -import java.util.UUID; -import java.util.concurrent.Future; -import java.util.concurrent.ScheduledFuture; -import java.util.concurrent.TimeUnit; -import lombok.Getter; - -@Getter -class ScheduledTaskDetail { - - private final Future future; - private final long startTime; - private final String id; - - ScheduledTaskDetail(long startTime, Future future) { - this.startTime = startTime; - this.future = future; - this.id = UUID.randomUUID().toString(); - } - - @Override - public String toString() { - StringBuilder sb = new StringBuilder(); - if (future instanceof ScheduledFuture) { - sb.append("ScheduledFuture(delay="); - sb.append(((ScheduledFuture) future).getDelay(TimeUnit.MILLISECONDS)); - sb.append("Ms, "); - } else { - sb.append("Future("); - } - sb.append("id="); - sb.append(id); - sb.append(", startTime="); - sb.append(startTime); - sb.append(", currentTime="); - sb.append(System.currentTimeMillis()); - sb.append(")"); - return sb.toString(); - } -} diff --git a/rqueue-spring-boot-example/src/main/java/com/github/sonus21/rqueue/example/MessageListener.java b/rqueue-spring-boot-example/src/main/java/com/github/sonus21/rqueue/example/MessageListener.java index acedfcacf..569563bd4 100644 --- a/rqueue-spring-boot-example/src/main/java/com/github/sonus21/rqueue/example/MessageListener.java +++ b/rqueue-spring-boot-example/src/main/java/com/github/sonus21/rqueue/example/MessageListener.java @@ -66,6 +66,7 @@ public void onSimpleMessage(String message) { } @RqueueListener( + active = "false", value = {"${rqueue.delay.queue}", "${rqueue.delay2.queue}"}, numRetries = "${rqueue.delay.queue.retries}", visibilityTimeout = "60*60*1000") @@ -74,6 +75,7 @@ public void onMessage(String message) { } @RqueueListener( + active = "false", value = "job-queue", deadLetterQueue = "job-morgue", numRetries = "2", @@ -85,6 +87,7 @@ public void onJobMessage(Job job) { @RqueueListener( + active = "false", value = "sch-job-queue", deadLetterQueue = "job-morgue", numRetries = "2", @@ -100,7 +103,8 @@ public void onSchJobMessage(Job job, @Header(RqueueMessageHeaders.ID) String mes } - @RqueueListener(value = "job-morgue", numRetries = "1", concurrency = "1-3") + @RqueueListener(value = "job-morgue", active = "false", + numRetries = "1", concurrency = "1-3") public void onJobDlqMessage(Job job) { execute("job-morgue: {}", job, true); } diff --git a/rqueue-spring-boot-example/src/main/resources/logback.xml b/rqueue-spring-boot-example/src/main/resources/logback.xml index e77d8d73b..d6a612743 100644 --- a/rqueue-spring-boot-example/src/main/resources/logback.xml +++ b/rqueue-spring-boot-example/src/main/resources/logback.xml @@ -17,16 +17,12 @@ - [%date{dd-MM-yyyy HH:mm:ss.SSS}] [%thread] %-5level %X{traceId:-} %X{spanId:-} - ${PID:-} %logger{36} - %msg%n - + [%date{dd-MM-yyyy HH:mm:ss.SSS}] [%thread] %-5level %X{traceId:-} %X{spanId:-} ${PID:-} %logger{36} - %msg%n - [%date{dd-MM-yyyy HH:mm:ss.SSS}] [%thread] %-5level %X{traceId:-} %X{spanId:-} - ${PID:-} %logger{36} - %msg%n - + [%date{dd-MM-yyyy HH:mm:ss.SSS}] [%thread] %-5level %X{traceId:-} %X{spanId:-} ${PID:-} %logger{36} - %msg%n log/app.log @@ -39,4 +35,10 @@ + + + + + + \ No newline at end of file From 6059d8f0fb5e4cf8414dce37f57031b20e550dd8 Mon Sep 17 00:00:00 2001 From: Sonu Kumar Date: Sat, 13 May 2023 15:39:38 +0530 Subject: [PATCH 2/9] added doc --- docs/CHANGELOG.md | 36 +++++++++++-------- docs/message-handling/producer-consumer.md | 13 ++++--- .../spring-configuration-metadata.json | 7 ++++ .../rqueue/example/MessageListener.java | 26 +++----------- 4 files changed, 42 insertions(+), 40 deletions(-) diff --git a/docs/CHANGELOG.md b/docs/CHANGELOG.md index 77cf61875..15b175291 100644 --- a/docs/CHANGELOG.md +++ b/docs/CHANGELOG.md @@ -8,22 +8,30 @@ layout: default All notable user-facing changes to this project are documented in this file. -## Release [3.0.1] 17-Jan-2022 +## Release [3.0.2] TBD +{: .highlight } +Migrate to this version to reduce resource utilization + +This will fix an important bug happening due to task multiplications. This is causing more Redis +resource usage Please check #[193] + -We're so excited to release Rqueue `3.0.1`. This release supports the Java 17, Spring Boot 3.x and Spring Framework 6.x +## Release [3.0.1] 17-Jan-2022 +We're so excited to release Rqueue `3.0.1`. This release supports the Java 17, Spring Boot 3.x and +Spring Framework 6.x ### [2.13.0] - 25-Dec-2022 ### Fixes -{: .highlight} + +{: .highlight} Migrate to this version as soon as possible to avoid duplicate message consumption post deletion. * Important fix for parallel message deletion or delete the message from message listener * No threads are available, improvement on message poller * Use System Zone ID for UI bottom screen - ### [2.12.0] - 14-Dec-2022 ### Fixes @@ -33,15 +41,14 @@ Migrate to this version as soon as possible to avoid duplicate message consumpti ### [2.11.1] - 18-Nov-2022 -{: .highlight} -Migrate to this version as soon as possible to avoid message build up. Messages in scheduled queue -can grow if poller is failing. Workaround is to restart the application. +{: .highlight} +Migrate to this version as soon as possible to avoid message build up. Messages in scheduled queue +can grow if poller is failing. Workaround is to restart the application. -* Message mover unreliability, scheduled message were not getting consumed once redis connection error occurs +* Message mover unreliability, scheduled message were not getting consumed once redis connection + error occurs * Upgraded Jquery version - - ### [2.10.2] - 16-Jul-2022 ### Fixes @@ -60,10 +67,10 @@ can grow if poller is failing. Workaround is to restart the application. ### [2.10.0] - 10-Oct-2021 -{: .warning } +{: .warning } Breaking change, if you're controlling any internal settings of Rqueue using application environment -or configuration variable than application can break. We've renamed some config keys, [see](./migration#290-to-210) - +or configuration variable than application can break. We've renamed some config +keys, [see](./migration#290-to-210) ### Fixes @@ -235,7 +242,6 @@ Breaking change, for migration [see](./migration#1x-to-2x) REDIS with prefix `rqueue-` then it will consider version 2. - Renamed annotation field `maxJobExecutionTime` to `visibilityTimeout` - ### Added - Web interface to visualize queue @@ -294,7 +300,6 @@ Breaking change, for migration [see](./migration#1x-to-2x) * The basic version of Asynchronous task execution using Redis for Spring and Spring Boot - [1.0]: https://repo1.maven.org/maven2/com/github/sonus21/rqueue/1.0-RELEASE [1.1]: https://repo1.maven.org/maven2/com/github/sonus21/rqueue/1.1-RELEASE @@ -356,3 +361,4 @@ Breaking change, for migration [see](./migration#1x-to-2x) [3.0.1]: https://repo1.maven.org/maven2/com/github/sonus21/rqueue-core/3.0.0-RELEASE [122]: https://github.com/sonus21/rqueue/issues/122 +[193]: https://github.com/sonus21/rqueue/issues/193 diff --git a/docs/message-handling/producer-consumer.md b/docs/message-handling/producer-consumer.md index b6c10b70e..3bac670f1 100644 --- a/docs/message-handling/producer-consumer.md +++ b/docs/message-handling/producer-consumer.md @@ -163,10 +163,10 @@ supported configurations. setting. * `rqueue.scheduler.auto.start=true` Rqueue scheduler also have thread pools, that handles the message. If you would like to use only event based message movement than set auto start as false. -* `rqueue.scheduler.scheduled.message.thread.pool.size=5` There could be many delayed queues, in that - case Rqueue has to move more messages from ZSET to LIST. In such cases, you can increase thread - pool size, the number of threads used for message movement is minimum of queue count and pool - size. +* `rqueue.scheduler.scheduled.message.thread.pool.size=5` There could be many delayed queues, in + that case Rqueue has to move more messages from ZSET to LIST. In such cases, you can increase + thread pool size, the number of threads used for message movement is minimum of queue count and + pool size. * `rqueue.scheduler.processing.message.thread.pool.size=1` there could be some dead message in processing queue as well, if you're seeing large number of dead messages in processing queue, then you should increase the thread pool size. Processing queue is used for at least once message @@ -174,6 +174,11 @@ supported configurations. * `rqueue.scheduler.scheduled.message.time.interval=5000` At what interval message should be moved from scheduled queue to normal queue. The default value is 5 seconds, that means, you can observe minimum delay of 5 seconds in delayed message consumption. +* `rqueue.scheduler.max.message.count=100` Rqueue continuously move scheduled messages from + processing/scheduled queue to normal queue so that we can process them asap. There are many + instances when large number of messages are scheduled to be run in next 5 minutes. In such cases + Rqueue can pull message from scheduled queue to normal queue at higher rate. By default, it copies + 100 messages from scheduled/processing queue to normal queue. ### Dead Letter Queue Consumer/Listener diff --git a/rqueue-core/src/main/resources/META-INF/spring-configuration-metadata.json b/rqueue-core/src/main/resources/META-INF/spring-configuration-metadata.json index 3cc46a5eb..93a28d51f 100644 --- a/rqueue-core/src/main/resources/META-INF/spring-configuration-metadata.json +++ b/rqueue-core/src/main/resources/META-INF/spring-configuration-metadata.json @@ -281,6 +281,13 @@ "type" : "java.lang.Integer", "defaultValue" : 1 }, + { + "sourceType" : "com.github.sonus21.rqueue.config.RqueueSchedulerConfig", + "name" : "rqueue.scheduler.max.message.count", + "description" : "Maximum number of messages that should be moved to normal queue from processing/schedule queue", + "type" : "java.lang.Integer", + "defaultValue" : 100 + }, { "sourceType" : "com.github.sonus21.rqueue.config.RqueueSchedulerConfig", "name" : "rqueue.scheduler.scheduled.message.time.interval", diff --git a/rqueue-spring-boot-example/src/main/java/com/github/sonus21/rqueue/example/MessageListener.java b/rqueue-spring-boot-example/src/main/java/com/github/sonus21/rqueue/example/MessageListener.java index 569563bd4..467b6f112 100644 --- a/rqueue-spring-boot-example/src/main/java/com/github/sonus21/rqueue/example/MessageListener.java +++ b/rqueue-spring-boot-example/src/main/java/com/github/sonus21/rqueue/example/MessageListener.java @@ -65,34 +65,19 @@ public void onSimpleMessage(String message) { execute("simple: {}", message, false); } - @RqueueListener( - active = "false", - value = {"${rqueue.delay.queue}", "${rqueue.delay2.queue}"}, - numRetries = "${rqueue.delay.queue.retries}", - visibilityTimeout = "60*60*1000") + @RqueueListener(value = {"${rqueue.delay.queue}", + "${rqueue.delay2.queue}"}, numRetries = "${rqueue.delay.queue.retries}", visibilityTimeout = "60*60*1000") public void onMessage(String message) { execute("delay: {}", message, true); } - @RqueueListener( - active = "false", - value = "job-queue", - deadLetterQueue = "job-morgue", - numRetries = "2", - deadLetterQueueListenerEnabled = "false", - concurrency = "10-20") + @RqueueListener(value = "job-queue", deadLetterQueue = "job-morgue", numRetries = "2", deadLetterQueueListenerEnabled = "false", concurrency = "10-20") public void onJobMessage(Job job) { execute("job-queue: {}", job, true); } - @RqueueListener( - active = "false", - value = "sch-job-queue", - deadLetterQueue = "job-morgue", - numRetries = "2", - deadLetterQueueListenerEnabled = "false", - concurrency = "1-3") + @RqueueListener(value = "sch-job-queue", deadLetterQueue = "job-morgue", numRetries = "2", deadLetterQueueListenerEnabled = "false", concurrency = "1-3") public void onSchJobMessage(Job job, @Header(RqueueMessageHeaders.ID) String messageId) { execute("sch-job-queue: {}", job, false); count += 1; @@ -103,8 +88,7 @@ public void onSchJobMessage(Job job, @Header(RqueueMessageHeaders.ID) String mes } - @RqueueListener(value = "job-morgue", active = "false", - numRetries = "1", concurrency = "1-3") + @RqueueListener(value = "job-morgue", numRetries = "1", concurrency = "1-3") public void onJobDlqMessage(Job job) { execute("job-morgue: {}", job, true); } From cd6464def02e67557f08469492dccab0ad7922a8 Mon Sep 17 00:00:00 2001 From: Sonu Kumar Date: Sun, 14 May 2023 10:36:54 +0530 Subject: [PATCH 3/9] add a warning log for not running tasks --- .../com/github/sonus21/rqueue/core/MessageScheduler.java | 5 +++++ .../main/java/com/github/sonus21/rqueue/utils/Constants.java | 2 +- 2 files changed, 6 insertions(+), 1 deletion(-) diff --git a/rqueue-core/src/main/java/com/github/sonus21/rqueue/core/MessageScheduler.java b/rqueue-core/src/main/java/com/github/sonus21/rqueue/core/MessageScheduler.java index b27b25b75..e997b30e6 100644 --- a/rqueue-core/src/main/java/com/github/sonus21/rqueue/core/MessageScheduler.java +++ b/rqueue-core/src/main/java/com/github/sonus21/rqueue/core/MessageScheduler.java @@ -347,6 +347,11 @@ private boolean cancelExistingTask(ScheduledTaskDetail scheduledTaskDetail, long queueDetail.getQueueName(), zsetName, scheduledTaskDetail); return false; } + if (existingDelay < Constants.TASK_ALIVE_TIME) { + getLogger().warn( + "MessageMoverTask {} has not run, you should consider increasing scheduled thread pool size", + scheduledTaskDetail); + } boolean cancelled = submittedTask.cancel(false); if (cancelled) { getLogger().debug("Task {} cancelled", scheduledTaskDetail.getId()); diff --git a/rqueue-core/src/main/java/com/github/sonus21/rqueue/utils/Constants.java b/rqueue-core/src/main/java/com/github/sonus21/rqueue/utils/Constants.java index 65e8449e9..47314b2ab 100644 --- a/rqueue-core/src/main/java/com/github/sonus21/rqueue/utils/Constants.java +++ b/rqueue-core/src/main/java/com/github/sonus21/rqueue/utils/Constants.java @@ -37,7 +37,7 @@ public final class Constants { public static final long MIN_DELAY = 100L; public static final long MIN_EXECUTION_TIME = MIN_DELAY; public static final long DELTA_BETWEEN_RE_ENQUEUE_TIME = ONE_MILLI; - public static final long TASK_ALIVE_TIME = -30 * ONE_MILLI; + public static final long TASK_ALIVE_TIME = -10 * ONE_MILLI; public static final int DEFAULT_RETRY_DEAD_LETTER_QUEUE = 3; public static final int MAX_MESSAGES = 100; public static final int DEFAULT_WORKER_COUNT_PER_QUEUE = 2; From 05bfcac965ee16509f152033a5a2bbca0d7922f9 Mon Sep 17 00:00:00 2001 From: Sonu Kumar Date: Sun, 14 May 2023 12:58:34 +0530 Subject: [PATCH 4/9] wip --- .../sonus21/rqueue/core/MessageScheduler.java | 20 +- .../rqueue/core/MessageSchedulerTest.java | 4 +- ...leTest.java => MessageSchedulingTest.java} | 53 ++--- .../ProcessingQueueMessageSchedulerTest.java | 17 +- .../ScheduledQueueMessageSchedulerTest.java | 214 ++++++++++-------- .../sonus21/test/TestTaskScheduler.java | 6 +- .../src/main/resources/logback.xml | 22 +- 7 files changed, 182 insertions(+), 154 deletions(-) rename rqueue-core/src/test/java/com/github/sonus21/rqueue/core/{MessageScheduleTest.java => MessageSchedulingTest.java} (85%) diff --git a/rqueue-core/src/main/java/com/github/sonus21/rqueue/core/MessageScheduler.java b/rqueue-core/src/main/java/com/github/sonus21/rqueue/core/MessageScheduler.java index e997b30e6..4c815cbe1 100644 --- a/rqueue-core/src/main/java/com/github/sonus21/rqueue/core/MessageScheduler.java +++ b/rqueue-core/src/main/java/com/github/sonus21/rqueue/core/MessageScheduler.java @@ -62,7 +62,7 @@ public abstract class MessageScheduler implements DisposableBean, @Autowired protected RqueueConfig rqueueConfig; private RedisScript redisScript; - private MessageSchedulerListener messageSchedulerListener; + protected MessageListener messageSchedulerListener; private DefaultScriptExecutor defaultScriptExecutor; private Map queueRunningState; private Map queueNameToScheduledTask; @@ -296,7 +296,7 @@ private boolean shouldNotSchedule(String queueName, boolean forceSchedule) { private void handleTaskOverride(ScheduledTaskDetail scheduledTaskDetail, QueueDetail queueDetail, String zsetName, long startTime) { // we should not schedule too frequent calls - long difference = startTime - scheduledTaskDetail.getStartTime(); + long difference = startTime - scheduledTaskDetail.getScheduleTime(); if (difference < getMinDelay()) { return; } @@ -324,8 +324,8 @@ protected synchronized void schedule(String queueName, Long startTime, boolean f } private void addTask(MessageMoverTask task, ScheduledTaskDetail scheduledTaskDetail) { - getLogger().debug("Adding Task task={}, startTime={}", task, - scheduledTaskDetail.getStartTime()); + getLogger().debug("Adding Task task={}, scheduleTime={}", task, + scheduledTaskDetail.getScheduleTime()); queueNameToLastMessageScheduleTime.put(task.getName(), System.currentTimeMillis()); queueNameToScheduledTask.put(task.getName(), scheduledTaskDetail); } @@ -339,7 +339,7 @@ private boolean cancelExistingTask(ScheduledTaskDetail scheduledTaskDetail, long return false; } // run existing tasks continue - long existingDelay = scheduledTaskDetail.getStartTime() - currentTime; + long existingDelay = scheduledTaskDetail.getScheduleTime() - currentTime; // task is about to run or was scheduled to run in last alive time, but could not so run it if (existingDelay < MIN_DELAY && existingDelay > Constants.TASK_ALIVE_TIME) { ThreadUtils.waitForTermination(getLogger(), submittedTask, @@ -475,11 +475,11 @@ public void onMessage(Message message, byte[] pattern) { private static class ScheduledTaskDetail { private final Future future; - private final long startTime; + private final long scheduleTime; private final String id; - ScheduledTaskDetail(String id, Future future, long startTime) { - this.startTime = startTime; + ScheduledTaskDetail(String id, Future future, long scheduleTime) { + this.scheduleTime = scheduleTime; this.future = future; this.id = id; } @@ -496,8 +496,8 @@ public String toString() { } sb.append("id="); sb.append(id); - sb.append(", startTime="); - sb.append(startTime); + sb.append(", scheduleTime="); + sb.append(scheduleTime); sb.append(", currentTime="); sb.append(System.currentTimeMillis()); sb.append(")"); diff --git a/rqueue-core/src/test/java/com/github/sonus21/rqueue/core/MessageSchedulerTest.java b/rqueue-core/src/test/java/com/github/sonus21/rqueue/core/MessageSchedulerTest.java index 8b52ce9a9..cea20b045 100644 --- a/rqueue-core/src/test/java/com/github/sonus21/rqueue/core/MessageSchedulerTest.java +++ b/rqueue-core/src/test/java/com/github/sonus21/rqueue/core/MessageSchedulerTest.java @@ -23,7 +23,7 @@ import com.github.sonus21.rqueue.CoreUnitTest; import com.github.sonus21.rqueue.config.RqueueConfig; import com.github.sonus21.rqueue.config.RqueueSchedulerConfig; -import com.github.sonus21.rqueue.core.ScheduledQueueMessageSchedulerTest.TestMessageScheduler; +import com.github.sonus21.rqueue.core.ScheduledQueueMessageSchedulerTest.TestScheduledQueueMessageScheduler; import com.github.sonus21.rqueue.listener.QueueDetail; import com.github.sonus21.rqueue.models.event.RqueueBootstrapEvent; import com.github.sonus21.rqueue.utils.TestUtils; @@ -56,7 +56,7 @@ class MessageSchedulerTest extends TestBase { @Mock private RedisTemplate redisTemplate; @InjectMocks - private TestMessageScheduler messageScheduler; + private TestScheduledQueueMessageScheduler messageScheduler; @BeforeEach public void init() { diff --git a/rqueue-core/src/test/java/com/github/sonus21/rqueue/core/MessageScheduleTest.java b/rqueue-core/src/test/java/com/github/sonus21/rqueue/core/MessageSchedulingTest.java similarity index 85% rename from rqueue-core/src/test/java/com/github/sonus21/rqueue/core/MessageScheduleTest.java rename to rqueue-core/src/test/java/com/github/sonus21/rqueue/core/MessageSchedulingTest.java index f26a75952..95523df44 100644 --- a/rqueue-core/src/test/java/com/github/sonus21/rqueue/core/MessageScheduleTest.java +++ b/rqueue-core/src/test/java/com/github/sonus21/rqueue/core/MessageSchedulingTest.java @@ -31,9 +31,12 @@ import com.github.sonus21.rqueue.core.ProcessingQueueMessageSchedulerTest.ProcessingQTestMessageScheduler; import com.github.sonus21.rqueue.listener.QueueDetail; import com.github.sonus21.rqueue.models.event.RqueueBootstrapEvent; +import com.github.sonus21.rqueue.utils.Constants; import com.github.sonus21.rqueue.utils.TestUtils; import com.github.sonus21.rqueue.utils.ThreadUtils; +import com.github.sonus21.rqueue.utils.TimeoutUtils; import com.github.sonus21.test.TestTaskScheduler; +import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicInteger; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; @@ -46,12 +49,14 @@ import org.springframework.data.redis.RedisConnectionFailureException; import org.springframework.data.redis.RedisSystemException; import org.springframework.data.redis.TooManyClusterRedirectionsException; +import org.springframework.data.redis.connection.DefaultMessage; +import org.springframework.data.redis.connection.Message; import org.springframework.data.redis.core.RedisCallback; import org.springframework.data.redis.core.RedisTemplate; @CoreUnitTest @SuppressWarnings("unchecked") -class MessageScheduleTest extends TestBase { +class MessageSchedulingTest extends TestBase { @InjectMocks private final ProcessingQTestMessageScheduler messageScheduler = new ProcessingQTestMessageScheduler(); @@ -81,16 +86,12 @@ void onCompletionOfExistingTaskNewTaskShouldBeSubmitted() throws Exception { doReturn(true).when(rqueueSchedulerConfig).isEnabled(); doReturn(true).when(rqueueSchedulerConfig).isRedisEnabled(); AtomicInteger counter = new AtomicInteger(0); - doAnswer( - invocation -> { - counter.incrementAndGet(); - return null; - }) - .when(redisTemplate) - .execute(any(RedisCallback.class)); + doAnswer(invocation -> { + counter.incrementAndGet(); + return null; + }).when(redisTemplate).execute(any(RedisCallback.class)); TestTaskScheduler scheduler = new TestTaskScheduler(); - threadUtils - .when(() -> ThreadUtils.createTaskScheduler(1, "processingQueueMsgScheduler-", 60)) + threadUtils.when(() -> ThreadUtils.createTaskScheduler(1, "processingQueueMsgScheduler-", 60)) .thenReturn(scheduler); messageScheduler.onApplicationEvent(new RqueueBootstrapEvent("Test", true)); waitFor(() -> counter.get() >= 1, "scripts are getting executed"); @@ -109,16 +110,12 @@ void multipleTasksAreRunningForTheSameQueue() throws Exception { doReturn(true).when(rqueueSchedulerConfig).isEnabled(); doReturn(true).when(rqueueSchedulerConfig).isRedisEnabled(); AtomicInteger counter = new AtomicInteger(0); - doAnswer( - invocation -> { - counter.incrementAndGet(); - return System.currentTimeMillis(); - }) - .when(redisTemplate) - .execute(any(RedisCallback.class)); + doAnswer(invocation -> { + counter.incrementAndGet(); + return System.currentTimeMillis(); + }).when(redisTemplate).execute(any(RedisCallback.class)); TestTaskScheduler scheduler = new TestTaskScheduler(); - threadUtils - .when(() -> ThreadUtils.createTaskScheduler(1, "processingQueueMsgScheduler-", 60)) + threadUtils.when(() -> ThreadUtils.createTaskScheduler(1, "processingQueueMsgScheduler-", 60)) .thenReturn(scheduler); messageScheduler.onApplicationEvent(new RqueueBootstrapEvent("Test", true)); waitFor(() -> counter.get() >= 2, "scripts are getting executed"); @@ -137,17 +134,13 @@ void taskShouldBeScheduledOnFailure() throws Exception { doReturn(true).when(rqueueSchedulerConfig).isRedisEnabled(); doReturn(10000L).when(rqueueSchedulerConfig).getMaxMessageMoverDelay(); AtomicInteger counter = new AtomicInteger(0); - doAnswer( - invocation -> { - counter.incrementAndGet(); - throw new RedisSystemException("Something is not correct", - new NullPointerException("oops!")); - }) - .when(redisTemplate) - .execute(any(RedisCallback.class)); + doAnswer(invocation -> { + counter.incrementAndGet(); + throw new RedisSystemException("Something is not correct", + new NullPointerException("oops!")); + }).when(redisTemplate).execute(any(RedisCallback.class)); TestTaskScheduler scheduler = new TestTaskScheduler(); - threadUtils - .when(() -> ThreadUtils.createTaskScheduler(1, "processingQueueMsgScheduler-", 60)) + threadUtils.when(() -> ThreadUtils.createTaskScheduler(1, "processingQueueMsgScheduler-", 60)) .thenReturn(scheduler); messageScheduler.onApplicationEvent(new RqueueBootstrapEvent("Test", true)); waitFor(() -> counter.get() >= 2, "scripts are getting executed"); @@ -158,7 +151,7 @@ void taskShouldBeScheduledOnFailure() throws Exception { } @Test - void continuousTaskFailTask() throws Exception { + void continuousTaskFailure() throws Exception { try (MockedStatic threadUtils = Mockito.mockStatic(ThreadUtils.class)) { doReturn(1).when(rqueueSchedulerConfig).getProcessingMessageThreadPoolSize(); doReturn(true).when(rqueueSchedulerConfig).isAutoStart(); diff --git a/rqueue-core/src/test/java/com/github/sonus21/rqueue/core/ProcessingQueueMessageSchedulerTest.java b/rqueue-core/src/test/java/com/github/sonus21/rqueue/core/ProcessingQueueMessageSchedulerTest.java index fbf816f6c..3d1eed848 100644 --- a/rqueue-core/src/test/java/com/github/sonus21/rqueue/core/ProcessingQueueMessageSchedulerTest.java +++ b/rqueue-core/src/test/java/com/github/sonus21/rqueue/core/ProcessingQueueMessageSchedulerTest.java @@ -64,8 +64,7 @@ public void init() { @Test void getChannelName() { - assertEquals( - slowQueueDetail.getProcessingQueueChannelName(), + assertEquals(slowQueueDetail.getProcessingQueueChannelName(), messageScheduler.getChannelName(slowQueue)); } @@ -77,21 +76,19 @@ void getZsetName() { @Test void getNextScheduleTimeSlowQueue() { long currentTime = System.currentTimeMillis(); - assertThat( - messageScheduler.getNextScheduleTime(slowQueue, null), + assertThat(messageScheduler.getNextScheduleTime(slowQueue, null), greaterThanOrEqualTo(currentTime + 100000)); - assertEquals( - currentTime + 1000L, messageScheduler.getNextScheduleTime(slowQueue, currentTime + 1000L)); + assertEquals(currentTime + 1000L, + messageScheduler.getNextScheduleTime(slowQueue, currentTime + 1000L)); } @Test void getNextScheduleTimeFastQueue() { long currentTime = System.currentTimeMillis(); - assertThat( - messageScheduler.getNextScheduleTime(fastQueue, null), + assertThat(messageScheduler.getNextScheduleTime(fastQueue, null), greaterThanOrEqualTo(currentTime + 200000)); - assertEquals( - currentTime + 1000L, messageScheduler.getNextScheduleTime(fastQueue, currentTime + 1000L)); + assertEquals(currentTime + 1000L, + messageScheduler.getNextScheduleTime(fastQueue, currentTime + 1000L)); } diff --git a/rqueue-core/src/test/java/com/github/sonus21/rqueue/core/ScheduledQueueMessageSchedulerTest.java b/rqueue-core/src/test/java/com/github/sonus21/rqueue/core/ScheduledQueueMessageSchedulerTest.java index 1a7da35fa..791c15032 100644 --- a/rqueue-core/src/test/java/com/github/sonus21/rqueue/core/ScheduledQueueMessageSchedulerTest.java +++ b/rqueue-core/src/test/java/com/github/sonus21/rqueue/core/ScheduledQueueMessageSchedulerTest.java @@ -36,6 +36,7 @@ import com.github.sonus21.rqueue.config.RqueueSchedulerConfig; import com.github.sonus21.rqueue.listener.QueueDetail; import com.github.sonus21.rqueue.models.event.RqueueBootstrapEvent; +import com.github.sonus21.rqueue.utils.Constants; import com.github.sonus21.rqueue.utils.TestUtils; import com.github.sonus21.rqueue.utils.ThreadUtils; import com.github.sonus21.rqueue.utils.TimeoutUtils; @@ -43,7 +44,9 @@ import java.util.List; import java.util.Map; import java.util.Vector; +import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicInteger; +import java.util.concurrent.atomic.AtomicLong; import org.apache.commons.lang3.reflect.FieldUtils; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; @@ -66,8 +69,6 @@ @SuppressWarnings("unchecked") class ScheduledQueueMessageSchedulerTest extends TestBase { - @InjectMocks - private final TestMessageScheduler messageScheduler = new TestMessageScheduler(); private final String slowQueue = "slow-queue"; private final String fastQueue = "fast-queue"; private final QueueDetail slowQueueDetail = TestUtils.createQueueDetail(slowQueue); @@ -80,6 +81,8 @@ class ScheduledQueueMessageSchedulerTest extends TestBase { private RedisTemplate redisTemplate; @Mock private RqueueRedisListenerContainerFactory rqueueRedisListenerContainerFactory; + @InjectMocks + private TestScheduledQueueMessageScheduler messageScheduler; @BeforeEach public void init() { @@ -91,8 +94,8 @@ public void init() { @Test void getChannelName() { - assertEquals( - slowQueueDetail.getScheduledQueueChannelName(), messageScheduler.getChannelName(slowQueue)); + assertEquals(slowQueueDetail.getScheduledQueueChannelName(), + messageScheduler.getChannelName(slowQueue)); } @Test @@ -104,11 +107,9 @@ void getZsetName() { void getNextScheduleTime() { long currentTime = System.currentTimeMillis(); doReturn(5000L).when(rqueueSchedulerConfig).getScheduledMessageTimeIntervalInMilli(); - assertThat( - messageScheduler.getNextScheduleTime(slowQueue, null), + assertThat(messageScheduler.getNextScheduleTime(slowQueue, null), greaterThanOrEqualTo(currentTime + 5000L)); - assertThat( - messageScheduler.getNextScheduleTime(fastQueue, currentTime + 1000L), + assertThat(messageScheduler.getNextScheduleTime(fastQueue, currentTime + 1000L), greaterThanOrEqualTo(currentTime + 5000L)); } @@ -133,14 +134,14 @@ void start() throws Exception { doReturn(true).when(rqueueSchedulerConfig).isEnabled(); doReturn(1000L).when(rqueueSchedulerConfig).getScheduledMessageTimeIntervalInMilli(); messageScheduler.onApplicationEvent(new RqueueBootstrapEvent("Test", true)); - Map queueRunningState = - (Map) FieldUtils.readField(messageScheduler, "queueRunningState", true); + Map queueRunningState = (Map) FieldUtils.readField( + messageScheduler, "queueRunningState", true); assertEquals(2, queueRunningState.size()); assertTrue(queueRunningState.get(slowQueue)); - assertEquals( - 2, ((Map) FieldUtils.readField(messageScheduler, "queueNameToScheduledTask", true)).size()); - assertEquals( - 2, ((Map) FieldUtils.readField(messageScheduler, "channelNameToQueueName", true)).size()); + assertEquals(2, + ((Map) FieldUtils.readField(messageScheduler, "queueNameToScheduledTask", true)).size()); + assertEquals(2, + ((Map) FieldUtils.readField(messageScheduler, "channelNameToQueueName", true)).size()); assertEquals(2, ((Map) FieldUtils.readField(messageScheduler, "queueSchedulers", true)).size()); TimeoutUtils.sleep(500L); messageScheduler.destroy(); @@ -153,14 +154,10 @@ void startAddsChannelToMessageListener() throws Exception { doReturn(true).when(rqueueSchedulerConfig).isAutoStart(); doReturn(true).when(rqueueSchedulerConfig).isEnabled(); doReturn(true).when(rqueueSchedulerConfig).isRedisEnabled(); - doNothing() - .when(rqueueRedisListenerContainerFactory) - .addMessageListener( - any(), eq(new ChannelTopic(slowQueueDetail.getScheduledQueueChannelName()))); - doNothing() - .when(rqueueRedisListenerContainerFactory) - .addMessageListener( - any(), eq(new ChannelTopic(fastQueueDetail.getScheduledQueueChannelName()))); + doNothing().when(rqueueRedisListenerContainerFactory).addMessageListener(any(), + eq(new ChannelTopic(slowQueueDetail.getScheduledQueueChannelName()))); + doNothing().when(rqueueRedisListenerContainerFactory).addMessageListener(any(), + eq(new ChannelTopic(fastQueueDetail.getScheduledQueueChannelName()))); messageScheduler.onApplicationEvent(new RqueueBootstrapEvent("Test", true)); TimeoutUtils.sleep(500L); messageScheduler.destroy(); @@ -176,12 +173,12 @@ void stop() throws Exception { messageScheduler.onApplicationEvent(new RqueueBootstrapEvent("Test", true)); TimeoutUtils.sleep(500L); messageScheduler.onApplicationEvent(new RqueueBootstrapEvent("Test", false)); - Map queueRunningState = - (Map) FieldUtils.readField(messageScheduler, "queueRunningState", true); + Map queueRunningState = (Map) FieldUtils.readField( + messageScheduler, "queueRunningState", true); assertEquals(2, queueRunningState.size()); assertFalse(queueRunningState.get(slowQueue)); - assertEquals( - 2, ((Map) FieldUtils.readField(messageScheduler, "channelNameToQueueName", true)).size()); + assertEquals(2, + ((Map) FieldUtils.readField(messageScheduler, "channelNameToQueueName", true)).size()); assertTrue( ((Map) FieldUtils.readField(messageScheduler, "queueNameToScheduledTask", true)).isEmpty()); messageScheduler.destroy(); @@ -196,19 +193,17 @@ void destroy() throws Exception { doReturn(1000L).when(rqueueSchedulerConfig).getScheduledMessageTimeIntervalInMilli(); TestTaskScheduler scheduler = new TestTaskScheduler(); try (MockedStatic threadUtils = Mockito.mockStatic(ThreadUtils.class)) { - threadUtils - .when(() -> ThreadUtils.createTaskScheduler(1, "scheduledQueueMsgScheduler-", 60)) + threadUtils.when(() -> ThreadUtils.createTaskScheduler(1, "scheduledQueueMsgScheduler-", 60)) .thenReturn(scheduler); messageScheduler.onApplicationEvent(new RqueueBootstrapEvent("Test", true)); TimeoutUtils.sleep(500L); messageScheduler.destroy(); - Map queueRunningState = - (Map) FieldUtils.readField(messageScheduler, "queueRunningState", true); + Map queueRunningState = (Map) FieldUtils.readField( + messageScheduler, "queueRunningState", true); assertEquals(2, queueRunningState.size()); assertFalse(queueRunningState.get(slowQueue)); - assertTrue( - ((Map) FieldUtils.readField(messageScheduler, "queueNameToScheduledTask", true)) - .isEmpty()); + assertTrue(((Map) FieldUtils.readField(messageScheduler, "queueNameToScheduledTask", + true)).isEmpty()); assertTrue(scheduler.shutdown); } } @@ -221,8 +216,7 @@ void startSubmitsTask() throws Exception { doReturn(true).when(rqueueSchedulerConfig).isRedisEnabled(); TestTaskScheduler scheduler = new TestTaskScheduler(); try (MockedStatic threadUtils = Mockito.mockStatic(ThreadUtils.class)) { - threadUtils - .when(() -> ThreadUtils.createTaskScheduler(1, "scheduledQueueMsgScheduler-", 60)) + threadUtils.when(() -> ThreadUtils.createTaskScheduler(1, "scheduledQueueMsgScheduler-", 60)) .thenReturn(scheduler); messageScheduler.onApplicationEvent(new RqueueBootstrapEvent("Test", true)); assertTrue(scheduler.submittedTasks() >= 1); @@ -238,13 +232,10 @@ void startSubmitsTaskAndThatGetsExecuted() throws Exception { doReturn(true).when(rqueueSchedulerConfig).isRedisEnabled(); doReturn(1000L).when(rqueueSchedulerConfig).getScheduledMessageTimeIntervalInMilli(); AtomicInteger counter = new AtomicInteger(0); - doAnswer( - invocation -> { - counter.incrementAndGet(); - return null; - }) - .when(redisTemplate) - .execute(any(RedisCallback.class)); + doAnswer(invocation -> { + counter.incrementAndGet(); + return null; + }).when(redisTemplate).execute(any(RedisCallback.class)); messageScheduler.onApplicationEvent(new RqueueBootstrapEvent("Test", true)); waitFor(() -> counter.get() >= 1, "scripts are getting executed"); messageScheduler.destroy(); @@ -259,16 +250,12 @@ void onCompletionOfExistingTaskNewTaskShouldBeSubmitted() throws Exception { doReturn(true).when(rqueueSchedulerConfig).isRedisEnabled(); doReturn(1000L).when(rqueueSchedulerConfig).getScheduledMessageTimeIntervalInMilli(); AtomicInteger counter = new AtomicInteger(0); - doAnswer( - invocation -> { - counter.incrementAndGet(); - return null; - }) - .when(redisTemplate) - .execute(any(RedisCallback.class)); + doAnswer(invocation -> { + counter.incrementAndGet(); + return null; + }).when(redisTemplate).execute(any(RedisCallback.class)); TestTaskScheduler scheduler = new TestTaskScheduler(); - threadUtils - .when(() -> ThreadUtils.createTaskScheduler(1, "scheduledQueueMsgScheduler-", 60)) + threadUtils.when(() -> ThreadUtils.createTaskScheduler(1, "scheduledQueueMsgScheduler-", 60)) .thenReturn(scheduler); messageScheduler.onApplicationEvent(new RqueueBootstrapEvent("Test", true)); waitFor(() -> counter.get() >= 1, "scripts are getting executed"); @@ -287,17 +274,13 @@ void taskShouldBeScheduledOnFailure() throws Exception { doReturn(true).when(rqueueSchedulerConfig).isRedisEnabled(); doReturn(10000L).when(rqueueSchedulerConfig).getMaxMessageMoverDelay(); AtomicInteger counter = new AtomicInteger(0); - doAnswer( - invocation -> { - counter.incrementAndGet(); - throw new RedisSystemException("Something is not correct", - new NullPointerException("oops!")); - }) - .when(redisTemplate) - .execute(any(RedisCallback.class)); + doAnswer(invocation -> { + counter.incrementAndGet(); + throw new RedisSystemException("Something is not correct", + new NullPointerException("oops!")); + }).when(redisTemplate).execute(any(RedisCallback.class)); TestTaskScheduler scheduler = new TestTaskScheduler(); - threadUtils - .when(() -> ThreadUtils.createTaskScheduler(1, "scheduledQueueMsgScheduler-", 60)) + threadUtils.when(() -> ThreadUtils.createTaskScheduler(1, "scheduledQueueMsgScheduler-", 60)) .thenReturn(scheduler); messageScheduler.onApplicationEvent(new RqueueBootstrapEvent("Test", true)); waitFor(() -> counter.get() >= 1, "scripts are getting executed"); @@ -316,24 +299,20 @@ void continuousTaskFailTask() throws Exception { doReturn(true).when(rqueueSchedulerConfig).isRedisEnabled(); doReturn(100L).when(rqueueSchedulerConfig).getMaxMessageMoverDelay(); AtomicInteger counter = new AtomicInteger(0); - doAnswer( - invocation -> { - int count = counter.incrementAndGet(); - if (count % 3 == 0) { - throw new RedisSystemException("Something is not correct", - new NullPointerException("oops!")); - } - if (count % 3 == 1) { - throw new RedisConnectionFailureException("Unknown host"); - } - throw new ClusterRedirectException(3, "localhost", 9004, - new TooManyClusterRedirectionsException("too many redirects")); - }) - .when(redisTemplate) - .execute(any(RedisCallback.class)); + doAnswer(invocation -> { + int count = counter.incrementAndGet(); + if (count % 3 == 0) { + throw new RedisSystemException("Something is not correct", + new NullPointerException("oops!")); + } + if (count % 3 == 1) { + throw new RedisConnectionFailureException("Unknown host"); + } + throw new ClusterRedirectException(3, "localhost", 9004, + new TooManyClusterRedirectionsException("too many redirects")); + }).when(redisTemplate).execute(any(RedisCallback.class)); TestTaskScheduler scheduler = new TestTaskScheduler(); - threadUtils - .when(() -> ThreadUtils.createTaskScheduler(1, "scheduledQueueMsgScheduler-", 60)) + threadUtils.when(() -> ThreadUtils.createTaskScheduler(1, "scheduledQueueMsgScheduler-", 60)) .thenReturn(scheduler); messageScheduler.onApplicationEvent(new RqueueBootstrapEvent("Test", true)); waitFor(() -> counter.get() >= 10, "scripts are getting executed"); @@ -352,8 +331,7 @@ void onMessageListenerTest() throws Exception { doReturn(true).when(rqueueSchedulerConfig).isRedisEnabled(); messageScheduler.onApplicationEvent(new RqueueBootstrapEvent("Test", true)); - MessageListener messageListener = - (MessageListener) FieldUtils.readField(messageScheduler, "messageSchedulerListener", true); + MessageListener messageListener = messageScheduler.messageSchedulerListener; // invalid channel messageListener.onMessage(new DefaultMessage(slowQueue.getBytes(), "312".getBytes()), null); TimeoutUtils.sleep(50); @@ -361,29 +339,87 @@ void onMessageListenerTest() throws Exception { // invalid body messageListener.onMessage( - new DefaultMessage( - slowQueueDetail.getScheduledQueueChannelName().getBytes(), "sss".getBytes()), - null); + new DefaultMessage(slowQueueDetail.getScheduledQueueChannelName().getBytes(), + "sss".getBytes()), null); TimeoutUtils.sleep(50); assertEquals(2, messageScheduler.scheduleList.stream().filter(e -> !e).count()); TimeoutUtils.sleep(110); // both are correct messageListener.onMessage( - new DefaultMessage( - slowQueueDetail.getScheduledQueueChannelName().getBytes(), - String.valueOf(System.currentTimeMillis()).getBytes()), - null); + new DefaultMessage(slowQueueDetail.getScheduledQueueChannelName().getBytes(), + String.valueOf(System.currentTimeMillis()).getBytes()), null); assertEquals(3, messageScheduler.scheduleList.stream().filter(e -> !e).count()); messageScheduler.destroy(); } - static class TestMessageScheduler extends ScheduledQueueMessageScheduler { + @Test + void taskMultiplier() throws Exception { + TestTaskScheduler scheduler = new TestTaskScheduler(2); + AtomicBoolean startGenerateMessage = new AtomicBoolean(false); + AtomicBoolean generateMessage = new AtomicBoolean(true); + String channelName = messageScheduler.getChannelName(slowQueue); + doReturn(1).when(rqueueSchedulerConfig).getScheduledMessageThreadPoolSize(); + doReturn(true).when(rqueueSchedulerConfig).isAutoStart(); + doReturn(true).when(rqueueSchedulerConfig).isEnabled(); + doReturn(true).when(rqueueSchedulerConfig).isRedisEnabled(); + doReturn(3000L).when(rqueueSchedulerConfig).getScheduledMessageTimeIntervalInMilli(); + doReturn(10000L).when(rqueueSchedulerConfig).getMaxMessageMoverDelay(); + doReturn(100L).when(rqueueSchedulerConfig).getMaxMessageCount(); + + Runnable messageGenerator = () -> { + int counter = 0; + while (generateMessage.get()) { + if (startGenerateMessage.get()) { + counter += 1; + byte[] currentTime = String.valueOf(System.currentTimeMillis()).getBytes(); + messageScheduler.messageSchedulerListener.onMessage( + new DefaultMessage(channelName.getBytes(), currentTime), null); + } + // each message would enqueue would lead to one event, so ~250 QPS + TimeoutUtils.sleep(4L); + } + System.out.println("Exiting sent " + counter + " messages "); + }; + scheduler.submit(messageGenerator); + + AtomicInteger counter = new AtomicInteger(0); + doAnswer(invocation -> { + counter.incrementAndGet(); + sleep(5); + return System.currentTimeMillis() - Constants.DEFAULT_SCRIPT_EXECUTION_TIME; + }).when(redisTemplate).execute(any(RedisCallback.class)); + + try (MockedStatic threadUtils = Mockito.mockStatic(ThreadUtils.class)) { + threadUtils.when(() -> ThreadUtils.createTaskScheduler(1, "scheduledQueueMsgScheduler-", 60)) + .thenReturn(scheduler); + messageScheduler.onApplicationEvent(new RqueueBootstrapEvent("Test", true)); + waitFor(() -> scheduler.submittedTasks() >= 2, "one start task to be submitted"); + startGenerateMessage.set(true); + // sleep for 2 seconds -> this should send around 500 events + TimeoutUtils.sleep(2000); + generateMessage.set(false); // disable Redis pub/sub event + TimeoutUtils.sleep(100); + int oldCount = counter.get(); + for (int i = 0; i < 20; i++) { + TimeoutUtils.sleep(100); + System.out.println(i + "=" + counter.get()); + } + int newCount = counter.get(); + messageScheduler.destroy(); + int jobCount = scheduler.submittedTasks(); + System.out.println(oldCount); + System.out.println(newCount); + System.out.println(jobCount); + } + } + + static class TestScheduledQueueMessageScheduler extends ScheduledQueueMessageScheduler { List scheduleList; - TestMessageScheduler() { + TestScheduledQueueMessageScheduler() { this.scheduleList = new Vector<>(); } diff --git a/rqueue-test-util/src/main/java/com/github/sonus21/test/TestTaskScheduler.java b/rqueue-test-util/src/main/java/com/github/sonus21/test/TestTaskScheduler.java index 7572bbdb2..27d033bbe 100644 --- a/rqueue-test-util/src/main/java/com/github/sonus21/test/TestTaskScheduler.java +++ b/rqueue-test-util/src/main/java/com/github/sonus21/test/TestTaskScheduler.java @@ -25,7 +25,6 @@ public class TestTaskScheduler extends ThreadPoolTaskScheduler { - private static final long serialVersionUID = 3617860362304703358L; public boolean shutdown = false; List> tasks = new Vector<>(); @@ -34,6 +33,11 @@ public TestTaskScheduler() { afterPropertiesSet(); } + public TestTaskScheduler(int poolSize) { + setPoolSize(poolSize); + afterPropertiesSet(); + } + @Override public Future submit(Runnable r) { Future f = super.submit(r); diff --git a/rqueue-test-util/src/main/resources/logback.xml b/rqueue-test-util/src/main/resources/logback.xml index ff31df820..9d0ffeb50 100644 --- a/rqueue-test-util/src/main/resources/logback.xml +++ b/rqueue-test-util/src/main/resources/logback.xml @@ -17,16 +17,12 @@ - %property{testClass}::%property{testName} [%date{dd-MM-yyyy HH:mm:ss.SSS}] [%thread] - %-5level %logger{36} - %msg%n - + %property{testClass}::%property{testName} [%date{dd-MM-yyyy HH:mm:ss.SSS}] [%thread] %-5level %logger{36} - %msg%n - %property{testClass}::%property{testName} [%date{dd-MM-yyyy HH:mm:ss.SSS}] [%thread] - %-5level %logger{36} - %msg%n - + %property{testClass}::%property{testName} [%date{dd-MM-yyyy HH:mm:ss.SSS}] [%thread] %-5level %logger{36} - %msg%n log/test.log @@ -38,9 +34,7 @@ - %property{testClass}::%property{testName} [%date{dd-MM-yyyy HH:mm:ss.SSS}] [%thread] - %-5level %logger{36} - %msg%n - + %property{testClass}::%property{testName} [%date{dd-MM-yyyy HH:mm:ss.SSS}] [%thread] %-5level %logger{36} - %msg%n log/monitor.log @@ -52,9 +46,7 @@ - %property{testClass}::%property{testName} [%date{dd-MM-yyyy HH:mm:ss.SSS}] [%thread] - %-5level %logger{36} - %msg%n - + %property{testClass}::%property{testName} [%date{dd-MM-yyyy HH:mm:ss.SSS}] [%thread] %-5level %logger{36} - %msg%n log/message.log @@ -99,5 +91,11 @@ + + + + + + \ No newline at end of file From 92c8c83370e891c78f991fc3d7189296d7c79e3f Mon Sep 17 00:00:00 2001 From: Sonu Kumar Date: Sat, 27 May 2023 13:28:03 +0530 Subject: [PATCH 5/9] Use hybrid combination of Redis and fixed rate scheduler to avoid job multiplications. --- .../rqueue/config/RqueueSchedulerConfig.java | 14 + .../sonus21/rqueue/core/MessageScheduler.java | 318 +++++------------- .../core/ProcessingQueueMessageScheduler.java | 7 +- .../core/RedisScheduleTriggerHandler.java | 161 +++++++++ .../RqueueRedisListenerContainerFactory.java | 7 +- .../core/ScheduledQueueMessageScheduler.java | 5 +- .../sonus21/rqueue/utils/Constants.java | 3 +- .../sonus21/rqueue/utils/ThreadUtils.java | 14 +- .../rqueue/core/MessageSchedulerTest.java | 2 +- .../ProcessingQueueMessageSchedulerTest.java | 29 +- .../core/RedisAndNormalSchedulingTest.java | 123 +++++++ .../core/RedisScheduleTriggerHandlerTest.java | 151 +++++++++ .../ScheduledQueueMessageSchedulerTest.java | 131 ++------ .../sonus21/test/TestTaskScheduler.java | 6 +- 14 files changed, 589 insertions(+), 382 deletions(-) create mode 100644 rqueue-core/src/main/java/com/github/sonus21/rqueue/core/RedisScheduleTriggerHandler.java create mode 100644 rqueue-core/src/test/java/com/github/sonus21/rqueue/core/RedisAndNormalSchedulingTest.java create mode 100644 rqueue-core/src/test/java/com/github/sonus21/rqueue/core/RedisScheduleTriggerHandlerTest.java diff --git a/rqueue-core/src/main/java/com/github/sonus21/rqueue/config/RqueueSchedulerConfig.java b/rqueue-core/src/main/java/com/github/sonus21/rqueue/config/RqueueSchedulerConfig.java index db9b2723f..b2406f38a 100644 --- a/rqueue-core/src/main/java/com/github/sonus21/rqueue/config/RqueueSchedulerConfig.java +++ b/rqueue-core/src/main/java/com/github/sonus21/rqueue/config/RqueueSchedulerConfig.java @@ -16,6 +16,7 @@ package com.github.sonus21.rqueue.config; +import com.github.sonus21.rqueue.utils.Constants; import lombok.Getter; import lombok.Setter; import org.springframework.beans.factory.annotation.Value; @@ -68,11 +69,24 @@ public class RqueueSchedulerConfig { @Value("${rqueue.scheduler.scheduled.message.time.interval:2000}") private long scheduledMessageTimeIntervalInMilli; + @Value("${rqueue.scheduler.termination.wait.time:1000}") + private long terminationWaitTime; // Maximum delay for message mover task due to failure @Value("${rqueue.scheduler.max.message.mover.delay:60000}") private long maxMessageMoverDelay; + // Minimum amount of time between two consecutive message move calls + @Value("${rqueue.scheduler.min.message.mover.delay:100}") + private long minMessageMoverDelay; + @Value("${rqueue.scheduler.max.message.count:100}") private long maxMessageCount; + + public long minMessageMoveDelay() { + if (minMessageMoverDelay <= 0) { + return Constants.MIN_SCHEDULE_INTERVAL; + } + return minMessageMoverDelay; + } } diff --git a/rqueue-core/src/main/java/com/github/sonus21/rqueue/core/MessageScheduler.java b/rqueue-core/src/main/java/com/github/sonus21/rqueue/core/MessageScheduler.java index 4c815cbe1..cd0d4cbc5 100644 --- a/rqueue-core/src/main/java/com/github/sonus21/rqueue/core/MessageScheduler.java +++ b/rqueue-core/src/main/java/com/github/sonus21/rqueue/core/MessageScheduler.java @@ -16,7 +16,6 @@ package com.github.sonus21.rqueue.core; -import static com.github.sonus21.rqueue.utils.Constants.MIN_DELAY; import static java.lang.Math.min; import com.github.sonus21.rqueue.config.RqueueConfig; @@ -26,7 +25,7 @@ import com.github.sonus21.rqueue.models.event.RqueueBootstrapEvent; import com.github.sonus21.rqueue.utils.Constants; import com.github.sonus21.rqueue.utils.ThreadUtils; -import java.time.Instant; +import java.time.Duration; import java.util.Arrays; import java.util.List; import java.util.Map; @@ -34,20 +33,16 @@ import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.Future; import java.util.concurrent.ScheduledFuture; -import java.util.concurrent.TimeUnit; -import lombok.Getter; +import com.google.common.annotations.VisibleForTesting; import org.slf4j.Logger; import org.springframework.beans.factory.DisposableBean; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.context.ApplicationListener; import org.springframework.data.redis.RedisSystemException; -import org.springframework.data.redis.connection.Message; -import org.springframework.data.redis.connection.MessageListener; import org.springframework.data.redis.core.RedisTemplate; import org.springframework.data.redis.core.script.DefaultScriptExecutor; import org.springframework.data.redis.core.script.RedisScript; -import org.springframework.data.redis.listener.ChannelTopic; import org.springframework.scheduling.annotation.Async; import org.springframework.scheduling.concurrent.ThreadPoolTaskScheduler; import org.springframework.util.Assert; @@ -62,12 +57,14 @@ public abstract class MessageScheduler implements DisposableBean, @Autowired protected RqueueConfig rqueueConfig; private RedisScript redisScript; - protected MessageListener messageSchedulerListener; private DefaultScriptExecutor defaultScriptExecutor; private Map queueRunningState; - private Map queueNameToScheduledTask; - private Map channelNameToQueueName; - private Map queueNameToLastMessageScheduleTime; + private Map> queueNameToScheduledTask; + private Map queueNameToNextRunTime; + + @VisibleForTesting + protected RedisScheduleTriggerHandler redisScheduleTriggerHandler; + private ThreadPoolTaskScheduler scheduler; @Autowired private RqueueRedisListenerContainerFactory rqueueRedisListenerContainerFactory; @@ -75,13 +72,11 @@ public abstract class MessageScheduler implements DisposableBean, @Autowired @Qualifier("rqueueRedisLongTemplate") private RedisTemplate redisTemplate; - - private Map queueSchedulers; private Map errorCount; protected abstract Logger getLogger(); - protected abstract long getNextScheduleTime(String queueName, Long value); + protected abstract long getNextScheduleTime(String queueName, long currentTime, Long value); protected abstract String getChannelName(String queueName); @@ -91,7 +86,15 @@ public abstract class MessageScheduler implements DisposableBean, protected abstract int getThreadPoolSize(); - protected abstract boolean isProcessingQueue(String queueName); + protected Duration getPeriod() { + long delay = rqueueSchedulerConfig.getScheduledMessageTimeIntervalInMilli(); + if (delay <= 0) { + delay = Constants.MIN_SCHEDULE_INTERVAL; + } + return Duration.ofMillis(delay); + } + + protected abstract boolean isProcessingQueue(); private void doStart() { for (String queueName : queueRunningState.keySet()) { @@ -99,14 +102,16 @@ private void doStart() { } } - private void subscribeToRedisTopic(String queueName) { - if (isRedisEnabled()) { - String channelName = getChannelName(queueName); - getLogger().debug("Queue {} subscribe to channel {}", queueName, channelName); - this.rqueueRedisListenerContainerFactory.addMessageListener(messageSchedulerListener, - new ChannelTopic(channelName)); - channelNameToQueueName.put(channelName, queueName); - } + private MessageMoverTask task(String queueName, boolean periodic) { + QueueDetail queueDetail = EndpointRegistry.get(queueName); + String zsetName = getZsetName(queueName); + return new MessageMoverTask(queueDetail, zsetName, periodic); + } + + protected void schedule(String queueName) { + Runnable task = task(queueName, true); + ScheduledFuture future = scheduler.scheduleAtFixedRate(task, getPeriod()); + queueNameToScheduledTask.put(queueName, future); } private void startQueue(String queueName) { @@ -114,14 +119,12 @@ private void startQueue(String queueName) { return; } queueRunningState.put(queueName, true); - if (scheduleTaskAtStartup() || !isRedisEnabled()) { - schedule(queueName, getQueueStartTime(), false); + if (scheduleTaskAtStartup()) { + schedule(queueName); + } + if (isRedisEnabled()) { + redisScheduleTriggerHandler.startQueue(queueName); } - subscribeToRedisTopic(queueName); - } - - protected long getQueueStartTime() { - return System.currentTimeMillis() + MIN_DELAY; } private void doStop() { @@ -135,19 +138,22 @@ private void doStop() { } waitForRunningQueuesToStop(); queueNameToScheduledTask.clear(); + if (isRedisEnabled()) { + redisScheduleTriggerHandler.stop(); + } } private void waitForRunningQueuesToStop() { for (Map.Entry runningState : queueRunningState.entrySet()) { String queueName = runningState.getKey(); - ScheduledTaskDetail scheduledTaskDetail = queueNameToScheduledTask.get(queueName); - if (scheduledTaskDetail != null) { - Future future = scheduledTaskDetail.getFuture(); - boolean completedOrCancelled = future.isCancelled() || future.isDone(); - if (!completedOrCancelled) { - future.cancel(true); - } - } + ScheduledFuture scheduledFuture = queueNameToScheduledTask.get(queueName); + ThreadUtils.waitForTermination( + getLogger(), + scheduledFuture, + rqueueSchedulerConfig.getTerminationWaitTime(), + "An exception occurred while stopping scheduler queue '{}'", + queueName + ); } } @@ -194,40 +200,32 @@ private boolean isQueueActive(String queueName) { return val; } - private long getLastScheduleTime(String queueName) { - return queueNameToLastMessageScheduleTime.getOrDefault(queueName, 0L); - } - - protected ScheduledTaskDetail getScheduledTask(String queueName) { - return queueNameToScheduledTask.get(queueName); - } - - protected void schedule(String queueName, Long startTime, boolean forceSchedule) { - this.queueSchedulers.get(queueName).schedule(queueName, startTime, forceSchedule); - } - protected void initialize() { List queueNames = EndpointRegistry.getActiveQueues(); defaultScriptExecutor = new DefaultScriptExecutor<>(redisTemplate); redisScript = RedisScriptFactory.getScript(ScriptType.MOVE_EXPIRED_MESSAGE); queueRunningState = new ConcurrentHashMap<>(queueNames.size()); queueNameToScheduledTask = new ConcurrentHashMap<>(queueNames.size()); - channelNameToQueueName = new ConcurrentHashMap<>(queueNames.size()); - queueNameToLastMessageScheduleTime = new ConcurrentHashMap<>(queueNames.size()); - queueSchedulers = new ConcurrentHashMap<>(queueNames.size()); + queueNameToNextRunTime = new ConcurrentHashMap<>(queueNames.size()); errorCount = new ConcurrentHashMap<>(queueNames.size()); createScheduler(queueNames.size()); - if (isRedisEnabled()) { - messageSchedulerListener = new MessageSchedulerListener(); - } for (String queueName : queueNames) { initQueue(queueName); } + if (isRedisEnabled()) { + redisScheduleTriggerHandler = new RedisScheduleTriggerHandler( + getLogger(), + rqueueRedisListenerContainerFactory, + rqueueSchedulerConfig, + queueNames, + this::addTask, + this::getChannelName); + redisScheduleTriggerHandler.initialize(); + } } private void initQueue(String queueName) { queueRunningState.put(queueName, false); - queueSchedulers.put(queueName, new QueueScheduler()); } @Override @@ -255,109 +253,8 @@ public void onApplicationEvent(RqueueBootstrapEvent event) { } } - protected long getMinDelay() { - return MIN_DELAY; - } - - private class QueueScheduler { - - private void scheduleNewTask(QueueDetail queueDetail, String zsetName, long startTime) { - MessageMoverTask task = new MessageMoverTask(queueDetail, zsetName); - Future future = scheduler.schedule(task, Instant.ofEpochMilli(startTime)); - addTask(task, new ScheduledTaskDetail(task.id, future, startTime)); - } - - private void scheduleTask(QueueDetail queueDetail, String zsetName, long startTime, - long currentTime) { - long requiredDelay = Math.max(1, startTime - currentTime); - long taskStartTime = currentTime; - MessageMoverTask task = new MessageMoverTask(queueDetail, zsetName); - Future future; - if (requiredDelay < MIN_DELAY) { - future = scheduler.submit(task); - } else { - taskStartTime = currentTime + requiredDelay; - future = scheduler.schedule(task, Instant.ofEpochMilli(taskStartTime)); - } - addTask(task, new ScheduledTaskDetail(task.id, future, taskStartTime)); - } - - private boolean shouldNotSchedule(String queueName, boolean forceSchedule) { - boolean isQueueActive = isQueueActive(queueName); - if (!isQueueActive || scheduler == null) { - return true; - } - long lastSeenTime = getLastScheduleTime(queueName); - long currentTime = System.currentTimeMillis(); - // ignore too frequents events - return !forceSchedule && currentTime - lastSeenTime < getMinDelay(); - } - - private void handleTaskOverride(ScheduledTaskDetail scheduledTaskDetail, - QueueDetail queueDetail, String zsetName, long startTime) { - // we should not schedule too frequent calls - long difference = startTime - scheduledTaskDetail.getScheduleTime(); - if (difference < getMinDelay()) { - return; - } - long currentTime = System.currentTimeMillis(); - if (cancelExistingTask(scheduledTaskDetail, currentTime, queueDetail, zsetName)) { - scheduleNewTask(queueDetail, zsetName, startTime); - } - } - - protected synchronized void schedule(String queueName, Long startTime, boolean forceSchedule) { - getLogger().debug("Schedule Task queue={}, force={}", queueName, forceSchedule); - if (shouldNotSchedule(queueName, forceSchedule)) { - return; - } - long currentTime = System.currentTimeMillis(); - ScheduledTaskDetail scheduledTaskDetail = getScheduledTask(queueName); - QueueDetail queueDetail = EndpointRegistry.get(queueName); - String zsetName = getZsetName(queueName); - // no task was scheduled or call came from existing MessageMoverTask - if (scheduledTaskDetail == null || forceSchedule) { - scheduleTask(queueDetail, zsetName, startTime, currentTime); - } else { - handleTaskOverride(scheduledTaskDetail, queueDetail, zsetName, startTime); - } - } - - private void addTask(MessageMoverTask task, ScheduledTaskDetail scheduledTaskDetail) { - getLogger().debug("Adding Task task={}, scheduleTime={}", task, - scheduledTaskDetail.getScheduleTime()); - queueNameToLastMessageScheduleTime.put(task.getName(), System.currentTimeMillis()); - queueNameToScheduledTask.put(task.getName(), scheduledTaskDetail); - } - - private boolean cancelExistingTask(ScheduledTaskDetail scheduledTaskDetail, long currentTime, - QueueDetail queueDetail, String zsetName) { - Future submittedTask = scheduledTaskDetail.getFuture(); - boolean completedOrCancelled = submittedTask.isDone() || submittedTask.isCancelled(); - // this can happen when task was not scheduled due to some exception in scheduling - if (completedOrCancelled) { - return false; - } - // run existing tasks continue - long existingDelay = scheduledTaskDetail.getScheduleTime() - currentTime; - // task is about to run or was scheduled to run in last alive time, but could not so run it - if (existingDelay < MIN_DELAY && existingDelay > Constants.TASK_ALIVE_TIME) { - ThreadUtils.waitForTermination(getLogger(), submittedTask, - Constants.DEFAULT_SCRIPT_EXECUTION_TIME, "LIST: {} ZSET: {}, Task: {} failed", - queueDetail.getQueueName(), zsetName, scheduledTaskDetail); - return false; - } - if (existingDelay < Constants.TASK_ALIVE_TIME) { - getLogger().warn( - "MessageMoverTask {} has not run, you should consider increasing scheduled thread pool size", - scheduledTaskDetail); - } - boolean cancelled = submittedTask.cancel(false); - if (cancelled) { - getLogger().debug("Task {} cancelled", scheduledTaskDetail.getId()); - } - return cancelled; - } + protected Future addTask(String queueName) { + return scheduler.submit(task(queueName, false)); } private class MessageMoverTask implements Runnable { @@ -367,16 +264,18 @@ private class MessageMoverTask implements Runnable { private final String queueName; private final String zsetName; private final boolean processingQueue; + private final boolean periodic; - MessageMoverTask(QueueDetail queueDetail, String zsetName) { + MessageMoverTask(QueueDetail queueDetail, String zsetName, boolean periodic) { this.id = UUID.randomUUID().toString(); this.name = queueDetail.getName(); this.queueName = queueDetail.getQueueName(); this.zsetName = zsetName; - this.processingQueue = isProcessingQueue(zsetName); + this.periodic = periodic; + this.processingQueue = isProcessingQueue(); } - private long getNextScheduleTimeInternal(Long value, Exception e) { + private long getNextScheduleTimeInternal(Long value, long currentTime, Exception e) { int errCount = 0; long nextTime; if (null != e) { @@ -384,11 +283,12 @@ private long getNextScheduleTimeInternal(Long value, Exception e) { if (errCount % 3 == 0) { getLogger().error("Message mover task is failing continuously queue: {}", name, e); } - long delay = (long) (MIN_DELAY * Math.pow(1.5, errCount)); - delay = Math.min(delay, rqueueSchedulerConfig.getMaxMessageMoverDelay()); - nextTime = System.currentTimeMillis() + delay; + // delay = x * 1.5^errorCount + double delay = rqueueSchedulerConfig.minMessageMoveDelay() * Math.pow(1.5, errCount); + long maxDelay = Math.min((long) delay, rqueueSchedulerConfig.getMaxMessageMoverDelay()); + nextTime = currentTime + maxDelay; } else { - nextTime = getNextScheduleTime(name, value); + nextTime = getNextScheduleTime(name, currentTime, value); } errorCount.put(name, errCount); return nextTime; @@ -396,7 +296,7 @@ private long getNextScheduleTimeInternal(Long value, Exception e) { @Override public String toString() { - return String.format("MessageMoverTask(id=%s, queue=%s)", id, name); + return String.format("MessageMoverTask(id=%s, queue=%s, periodic=%s)", id, name, periodic); } private long getMessageCount() { @@ -412,8 +312,18 @@ private Object[] scriptArgs() { return new Object[]{currentTime, getMessageCount(), processingQueue ? 1 : 0}; } + private boolean shouldSkip(long currentTime) { + Long nextRunTime = queueNameToNextRunTime.get(queueName); + return nextRunTime != null && nextRunTime > currentTime; + } + @Override public void run() { + long currentTime = System.currentTimeMillis(); + if (shouldSkip(currentTime)) { + getLogger().debug("Skipped {}", this); + return; + } getLogger().debug("Running {}", this); Long value = null; Exception e = null; @@ -427,8 +337,8 @@ public void run() { e = ex; getLogger().warn("Task execution failed for the queue: {}", getName(), e); } finally { - long nextExecutionTime = getNextScheduleTimeInternal(value, e); - schedule(name, nextExecutionTime, true); + long nextExecutionTime = getNextScheduleTimeInternal(value, currentTime, e); + queueNameToNextRunTime.put(queueName, nextExecutionTime); } } @@ -437,72 +347,4 @@ public String getName() { } } - private class MessageSchedulerListener implements MessageListener { - - private void handleMessage(String queueName, Long startTime) { - long lastSeenTime = getLastScheduleTime(queueName); - long currentTime = System.currentTimeMillis(); - if (currentTime - lastSeenTime < getMinDelay()) { - return; - } - long jobStartTime = getNextScheduleTime(queueName, startTime); - schedule(queueName, jobStartTime, false); - } - - @Override - public void onMessage(Message message, byte[] pattern) { - if (message.getBody().length == 0 || message.getChannel().length == 0) { - return; - } - String body = new String(message.getBody()); - String channel = new String(message.getChannel()); - getLogger().trace("Body: {} Channel: {}", body, channel); - try { - Long startTime = Long.parseLong(body); - String queueName = channelNameToQueueName.get(channel); - if (queueName == null) { - getLogger().warn("Unknown channel name {}", channel); - return; - } - handleMessage(queueName, startTime); - } catch (Exception e) { - getLogger().error("Error occurred on a channel {}, body: {}", channel, body, e); - } - } - } - - @Getter - private static class ScheduledTaskDetail { - - private final Future future; - private final long scheduleTime; - private final String id; - - ScheduledTaskDetail(String id, Future future, long scheduleTime) { - this.scheduleTime = scheduleTime; - this.future = future; - this.id = id; - } - - @Override - public String toString() { - StringBuilder sb = new StringBuilder(); - if (future instanceof ScheduledFuture) { - sb.append("ScheduledFuture(delay="); - sb.append(((ScheduledFuture) future).getDelay(TimeUnit.MILLISECONDS)); - sb.append("Ms, "); - } else { - sb.append("Future("); - } - sb.append("id="); - sb.append(id); - sb.append(", scheduleTime="); - sb.append(scheduleTime); - sb.append(", currentTime="); - sb.append(System.currentTimeMillis()); - sb.append(")"); - return sb.toString(); - } - } - } diff --git a/rqueue-core/src/main/java/com/github/sonus21/rqueue/core/ProcessingQueueMessageScheduler.java b/rqueue-core/src/main/java/com/github/sonus21/rqueue/core/ProcessingQueueMessageScheduler.java index da1420749..276a09032 100644 --- a/rqueue-core/src/main/java/com/github/sonus21/rqueue/core/ProcessingQueueMessageScheduler.java +++ b/rqueue-core/src/main/java/com/github/sonus21/rqueue/core/ProcessingQueueMessageScheduler.java @@ -19,9 +19,11 @@ import static java.lang.Long.max; import com.github.sonus21.rqueue.listener.QueueDetail; +import java.time.Duration; import java.util.List; import java.util.Map; import java.util.concurrent.ConcurrentHashMap; +import com.github.sonus21.rqueue.utils.Constants; import lombok.extern.slf4j.Slf4j; import org.slf4j.Logger; @@ -61,7 +63,7 @@ protected int getThreadPoolSize() { } @Override - protected boolean isProcessingQueue(String queueName) { + protected boolean isProcessingQueue() { return true; } @@ -71,8 +73,7 @@ protected String getThreadNamePrefix() { } @Override - protected long getNextScheduleTime(String queueName, Long value) { - long currentTime = System.currentTimeMillis(); + protected long getNextScheduleTime(String queueName, long currentTime, Long value) { if (value == null) { long delay = queueNameToDelay.get(queueName); return currentTime + delay; diff --git a/rqueue-core/src/main/java/com/github/sonus21/rqueue/core/RedisScheduleTriggerHandler.java b/rqueue-core/src/main/java/com/github/sonus21/rqueue/core/RedisScheduleTriggerHandler.java new file mode 100644 index 000000000..4fb41d1b8 --- /dev/null +++ b/rqueue-core/src/main/java/com/github/sonus21/rqueue/core/RedisScheduleTriggerHandler.java @@ -0,0 +1,161 @@ +/* + * Copyright (c) 2019-2023 Sonu Kumar + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * You may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and limitations under the License. + * + */ +package com.github.sonus21.rqueue.core; + +import com.github.sonus21.rqueue.config.RqueueSchedulerConfig; +import com.github.sonus21.rqueue.utils.ThreadUtils; +import com.google.common.annotations.VisibleForTesting; +import org.slf4j.Logger; +import org.springframework.data.redis.connection.Message; +import org.springframework.data.redis.connection.MessageListener; +import org.springframework.data.redis.listener.ChannelTopic; + +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.Future; +import java.util.function.Function; + +class RedisScheduleTriggerHandler { + + private final RqueueRedisListenerContainerFactory rqueueRedisListenerContainerFactory; + private final RqueueSchedulerConfig rqueueSchedulerConfig; + private final Logger logger; + private final Function> scheduler; + private final Function channelNameProducer; + private final List queueNames; + + @VisibleForTesting + Map queueNameToLastRunTime; + @VisibleForTesting + Map> queueNameToFuture; + @VisibleForTesting + Map channelNameToQueueName; + @VisibleForTesting + MessageListener messageListener; + + RedisScheduleTriggerHandler(Logger logger, + RqueueRedisListenerContainerFactory rqueueRedisListenerContainerFactory, + RqueueSchedulerConfig rqueueSchedulerConfig, List queueNames, + Function> scheduler, Function channelNameProducer) { + this.queueNames = queueNames; + this.rqueueSchedulerConfig = rqueueSchedulerConfig; + this.rqueueRedisListenerContainerFactory = rqueueRedisListenerContainerFactory; + this.logger = logger; + this.scheduler = scheduler; + this.channelNameProducer = channelNameProducer; + } + + void initialize() { + this.messageListener = new MessageSchedulerListener(); + this.channelNameToQueueName = new HashMap<>(queueNames.size()); + this.queueNameToFuture = new ConcurrentHashMap<>(queueNames.size()); + this.queueNameToLastRunTime = new ConcurrentHashMap<>(queueNames.size()); + } + + void stop() { + for (String queue : queueNames) { + stopQueue(queue); + } + } + + void startQueue(String queueName) { + queueNameToLastRunTime.put(queueName, 0L); + subscribeToRedisTopic(queueName); + } + + void stopQueue(String queueName) { + Future future = queueNameToFuture.get(queueName); + ThreadUtils.waitForTermination(logger, future, rqueueSchedulerConfig.getTerminationWaitTime(), + "An exception occurred while stopping scheduler queue '{}'", queueName); + queueNameToLastRunTime.put(queueName, 0L); + queueNameToFuture.remove(queueName); + unsubscribeFromRedis(queueName); + } + + private void unsubscribeFromRedis(String queueName) { + String channelName = channelNameProducer.apply(queueName); + logger.debug("Queue {} unsubscribe from channel {}", queueName, channelName); + rqueueRedisListenerContainerFactory.removeMessageListener(messageListener, + new ChannelTopic(channelName)); + channelNameToQueueName.put(channelName, queueName); + } + + private void subscribeToRedisTopic(String queueName) { + String channelName = channelNameProducer.apply(queueName); + channelNameToQueueName.put(channelName, queueName); + logger.debug("Queue {} subscribe to channel {}", queueName, channelName); + rqueueRedisListenerContainerFactory.addMessageListener(messageListener, + new ChannelTopic(channelName)); + } + + protected long getMinDelay() { + return rqueueSchedulerConfig.minMessageMoveDelay(); + } + + + /** + * This MessageListener listen the event from Redis, its expected that the event should be only + * raised when elements in the ZSET are lagging behind current time. + */ + private class MessageSchedulerListener implements MessageListener { + + private void schedule(String queueName, long currentTime) { + Future future = queueNameToFuture.get(queueName); + if (future == null || future.isCancelled() || future.isDone()) { + queueNameToLastRunTime.put(queueName, currentTime); + Future newFuture = scheduler.apply(queueName); + queueNameToFuture.put(queueName, newFuture); + } + } + + private void handleMessage(String queueName, long startTime) { + long currentTime = System.currentTimeMillis(); + if (startTime > currentTime) { + logger.warn("Received message body is not correct queue: {}, time: {}", queueName, + startTime); + return; + } + long lastRunTime = queueNameToLastRunTime.get(queueName); + if (currentTime - lastRunTime < getMinDelay()) { + return; + } + schedule(queueName, currentTime); + } + + @Override + public void onMessage(Message message, byte[] pattern) { + if (message.getBody().length == 0 || message.getChannel().length == 0) { + return; + } + String body = new String(message.getBody()); + String channel = new String(message.getChannel()); + logger.trace("Body: {} Channel: {}", body, channel); + try { + long startTime = Long.parseLong(body); + String queueName = channelNameToQueueName.get(channel); + if (queueName == null) { + logger.warn("Unknown channel name {}", channel); + return; + } + handleMessage(queueName, startTime); + } catch (Exception e) { + logger.error("Error occurred on a channel {}, body: {}", channel, body, e); + } + } + } +} diff --git a/rqueue-core/src/main/java/com/github/sonus21/rqueue/core/RqueueRedisListenerContainerFactory.java b/rqueue-core/src/main/java/com/github/sonus21/rqueue/core/RqueueRedisListenerContainerFactory.java index 95215d3d2..1b6cf3edb 100644 --- a/rqueue-core/src/main/java/com/github/sonus21/rqueue/core/RqueueRedisListenerContainerFactory.java +++ b/rqueue-core/src/main/java/com/github/sonus21/rqueue/core/RqueueRedisListenerContainerFactory.java @@ -24,6 +24,7 @@ import org.springframework.beans.factory.annotation.Autowired; import org.springframework.context.SmartLifecycle; import org.springframework.data.redis.connection.MessageListener; +import org.springframework.data.redis.listener.ChannelTopic; import org.springframework.data.redis.listener.RedisMessageListenerContainer; import org.springframework.data.redis.listener.Topic; @@ -96,7 +97,7 @@ public void afterPropertiesSet() throws Exception { createContainer(); return; } - if (rqueueConfig.isSharedConnection() || rqueueSchedulerConfig.isListenerShared()) { + if (rqueueSchedulerConfig.isListenerShared()) { if (systemContainer != null) { container = systemContainer; sharedContainer = true; @@ -121,4 +122,8 @@ public void stop(Runnable callback) { public int getPhase() { return Integer.MAX_VALUE; } + + public void removeMessageListener(MessageListener messageListener, ChannelTopic channelTopic) { + getContainer().removeMessageListener(messageListener, channelTopic); + } } diff --git a/rqueue-core/src/main/java/com/github/sonus21/rqueue/core/ScheduledQueueMessageScheduler.java b/rqueue-core/src/main/java/com/github/sonus21/rqueue/core/ScheduledQueueMessageScheduler.java index 91b815230..ac660b091 100644 --- a/rqueue-core/src/main/java/com/github/sonus21/rqueue/core/ScheduledQueueMessageScheduler.java +++ b/rqueue-core/src/main/java/com/github/sonus21/rqueue/core/ScheduledQueueMessageScheduler.java @@ -28,8 +28,7 @@ protected Logger getLogger() { } @Override - protected long getNextScheduleTime(String queueName, Long value) { - long currentTime = System.currentTimeMillis(); + protected long getNextScheduleTime(String queueName, long currentTime, Long value) { if (value == null) { return currentTime + rqueueSchedulerConfig.getScheduledMessageTimeIntervalInMilli(); } @@ -60,7 +59,7 @@ protected int getThreadPoolSize() { } @Override - protected boolean isProcessingQueue(String queueName) { + protected boolean isProcessingQueue() { return false; } } diff --git a/rqueue-core/src/main/java/com/github/sonus21/rqueue/utils/Constants.java b/rqueue-core/src/main/java/com/github/sonus21/rqueue/utils/Constants.java index 47314b2ab..12fee1798 100644 --- a/rqueue-core/src/main/java/com/github/sonus21/rqueue/utils/Constants.java +++ b/rqueue-core/src/main/java/com/github/sonus21/rqueue/utils/Constants.java @@ -32,12 +32,11 @@ public final class Constants { public static final int MINUTES_IN_A_DAY = HOURS_IN_A_DAY * MINUTES_IN_AN_HOUR; public static final int SECONDS_IN_A_DAY = MINUTES_IN_A_DAY * SECONDS_IN_A_MINUTE; public static final long MILLIS_IN_A_DAY = SECONDS_IN_A_DAY * ONE_MILLI; - public static final int SECONDS_IN_A_WEEK = DAYS_IN_A_WEEK * SECONDS_IN_A_DAY; public static final long DEFAULT_SCRIPT_EXECUTION_TIME = 5 * ONE_MILLI; public static final long MIN_DELAY = 100L; + public static final long MIN_SCHEDULE_INTERVAL = 100L; public static final long MIN_EXECUTION_TIME = MIN_DELAY; public static final long DELTA_BETWEEN_RE_ENQUEUE_TIME = ONE_MILLI; - public static final long TASK_ALIVE_TIME = -10 * ONE_MILLI; public static final int DEFAULT_RETRY_DEAD_LETTER_QUEUE = 3; public static final int MAX_MESSAGES = 100; public static final int DEFAULT_WORKER_COUNT_PER_QUEUE = 2; diff --git a/rqueue-core/src/main/java/com/github/sonus21/rqueue/utils/ThreadUtils.java b/rqueue-core/src/main/java/com/github/sonus21/rqueue/utils/ThreadUtils.java index a033dc69d..45cff7067 100644 --- a/rqueue-core/src/main/java/com/github/sonus21/rqueue/utils/ThreadUtils.java +++ b/rqueue-core/src/main/java/com/github/sonus21/rqueue/utils/ThreadUtils.java @@ -78,16 +78,14 @@ private static void waitForShutdown( public static void waitForTermination( Logger log, Future future, long waitTimeInMillis, String msg, Object... msgParams) { - if (future == null) { + if (future == null || future.isCancelled() || future.isDone()) { return; } - boolean completedOrCancelled = future.isCancelled() || future.isDone(); - if (!completedOrCancelled) { - if (future instanceof ScheduledFuture) { - ScheduledFuture f = (ScheduledFuture) future; - if (f.getDelay(TimeUnit.MILLISECONDS) > Constants.MIN_DELAY) { - return; - } + if (future instanceof ScheduledFuture) { + ScheduledFuture f = (ScheduledFuture) future; + if (f.getDelay(TimeUnit.MILLISECONDS) > Constants.MIN_DELAY) { + f.cancel(false); + return; } } waitForShutdown(log, future, waitTimeInMillis, msg, msgParams); diff --git a/rqueue-core/src/test/java/com/github/sonus21/rqueue/core/MessageSchedulerTest.java b/rqueue-core/src/test/java/com/github/sonus21/rqueue/core/MessageSchedulerTest.java index cea20b045..bafefd9a0 100644 --- a/rqueue-core/src/test/java/com/github/sonus21/rqueue/core/MessageSchedulerTest.java +++ b/rqueue-core/src/test/java/com/github/sonus21/rqueue/core/MessageSchedulerTest.java @@ -69,7 +69,7 @@ public void init() { void afterPropertiesSetWithEmptyQueSet() throws Exception { EndpointRegistry.delete(); messageScheduler.onApplicationEvent(new RqueueBootstrapEvent("Test", true)); - assertEquals(0, messageScheduler.scheduleList.size()); + assertEquals(0, messageScheduler.scheduleCounter.get()); messageScheduler.destroy(); } diff --git a/rqueue-core/src/test/java/com/github/sonus21/rqueue/core/ProcessingQueueMessageSchedulerTest.java b/rqueue-core/src/test/java/com/github/sonus21/rqueue/core/ProcessingQueueMessageSchedulerTest.java index 3d1eed848..da1d9035e 100644 --- a/rqueue-core/src/test/java/com/github/sonus21/rqueue/core/ProcessingQueueMessageSchedulerTest.java +++ b/rqueue-core/src/test/java/com/github/sonus21/rqueue/core/ProcessingQueueMessageSchedulerTest.java @@ -26,8 +26,8 @@ import com.github.sonus21.rqueue.config.RqueueSchedulerConfig; import com.github.sonus21.rqueue.listener.QueueDetail; import com.github.sonus21.rqueue.utils.TestUtils; -import java.util.List; -import java.util.Vector; +import java.util.concurrent.Future; +import java.util.concurrent.atomic.AtomicInteger; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; import org.mockito.InjectMocks; @@ -76,34 +76,41 @@ void getZsetName() { @Test void getNextScheduleTimeSlowQueue() { long currentTime = System.currentTimeMillis(); - assertThat(messageScheduler.getNextScheduleTime(slowQueue, null), + assertThat(messageScheduler.getNextScheduleTime(slowQueue, currentTime, null), greaterThanOrEqualTo(currentTime + 100000)); assertEquals(currentTime + 1000L, - messageScheduler.getNextScheduleTime(slowQueue, currentTime + 1000L)); + messageScheduler.getNextScheduleTime(slowQueue, currentTime, currentTime + 1000L)); } @Test void getNextScheduleTimeFastQueue() { long currentTime = System.currentTimeMillis(); - assertThat(messageScheduler.getNextScheduleTime(fastQueue, null), + assertThat(messageScheduler.getNextScheduleTime(fastQueue, currentTime, null), greaterThanOrEqualTo(currentTime + 200000)); assertEquals(currentTime + 1000L, - messageScheduler.getNextScheduleTime(fastQueue, currentTime + 1000L)); + messageScheduler.getNextScheduleTime(fastQueue, currentTime, currentTime + 1000L)); } static class ProcessingQTestMessageScheduler extends ProcessingQueueMessageScheduler { - List scheduleList; + private final AtomicInteger schedulesCalls; + private final AtomicInteger addTaskCalls; ProcessingQTestMessageScheduler() { - this.scheduleList = new Vector<>(); + this.schedulesCalls = new AtomicInteger(0); + this.addTaskCalls = new AtomicInteger(0); } @Override - protected synchronized void schedule(String queueName, Long startTime, boolean forceSchedule) { - super.schedule(queueName, startTime, forceSchedule); - this.scheduleList.add(forceSchedule); + protected synchronized void schedule(String queueName) { + schedulesCalls.incrementAndGet(); + super.schedule(queueName); + } + + protected Future addTask(String queueName) { + addTaskCalls.incrementAndGet(); + return super.addTask(queueName); } } } diff --git a/rqueue-core/src/test/java/com/github/sonus21/rqueue/core/RedisAndNormalSchedulingTest.java b/rqueue-core/src/test/java/com/github/sonus21/rqueue/core/RedisAndNormalSchedulingTest.java new file mode 100644 index 000000000..3cb395304 --- /dev/null +++ b/rqueue-core/src/test/java/com/github/sonus21/rqueue/core/RedisAndNormalSchedulingTest.java @@ -0,0 +1,123 @@ +package com.github.sonus21.rqueue.core; + +import com.github.sonus21.TestBase; +import com.github.sonus21.rqueue.CoreUnitTest; +import com.github.sonus21.rqueue.config.RqueueConfig; +import com.github.sonus21.rqueue.config.RqueueSchedulerConfig; +import com.github.sonus21.rqueue.core.ScheduledQueueMessageSchedulerTest.TestScheduledQueueMessageScheduler; +import com.github.sonus21.rqueue.listener.QueueDetail; +import com.github.sonus21.rqueue.models.event.RqueueBootstrapEvent; +import com.github.sonus21.rqueue.utils.Constants; +import com.github.sonus21.rqueue.utils.TestUtils; +import com.github.sonus21.rqueue.utils.ThreadUtils; +import com.github.sonus21.rqueue.utils.TimeoutUtils; +import com.github.sonus21.test.TestTaskScheduler; +import lombok.extern.slf4j.Slf4j; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.mockito.InjectMocks; +import org.mockito.Mock; +import org.mockito.MockedStatic; +import org.mockito.Mockito; +import org.mockito.MockitoAnnotations; +import org.springframework.data.redis.connection.DefaultMessage; +import org.springframework.data.redis.core.RedisCallback; +import org.springframework.data.redis.core.RedisTemplate; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicInteger; + +import static com.github.sonus21.rqueue.utils.TimeoutUtils.sleep; +import static com.github.sonus21.rqueue.utils.TimeoutUtils.waitFor; +import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.doAnswer; +import static org.mockito.Mockito.doReturn; + +@CoreUnitTest +@Slf4j +class RedisAndNormalSchedulingTest extends TestBase { + + private final String slowQueue = "slow-queue"; + private final String fastQueue = "fast-queue"; + private final QueueDetail slowQueueDetail = TestUtils.createQueueDetail(slowQueue); + private final QueueDetail fastQueueDetail = TestUtils.createQueueDetail(fastQueue); + @Mock + private RqueueSchedulerConfig rqueueSchedulerConfig; + @Mock + private RqueueConfig rqueueConfig; + @Mock + private RedisTemplate redisTemplate; + @Mock + private RqueueRedisListenerContainerFactory rqueueRedisListenerContainerFactory; + @InjectMocks + private TestScheduledQueueMessageScheduler messageScheduler; + + @BeforeEach + public void init() { + MockitoAnnotations.openMocks(this); + EndpointRegistry.delete(); + EndpointRegistry.register(fastQueueDetail); + EndpointRegistry.register(slowQueueDetail); + } + + @Test + void redisAndNormalScheduling() throws Exception { + TestTaskScheduler scheduler = new TestTaskScheduler(2); + AtomicBoolean startGenerateMessage = new AtomicBoolean(false); + AtomicBoolean generateMessage = new AtomicBoolean(true); + long totalTime = 2000L; + long minDelay = 10L; + //15% buffer due to short polling intervals + double buffer = 0.13; + String channelName = messageScheduler.getChannelName(slowQueue); + doReturn(1).when(rqueueSchedulerConfig).getScheduledMessageThreadPoolSize(); + doReturn(true).when(rqueueSchedulerConfig).isAutoStart(); + doReturn(true).when(rqueueSchedulerConfig).isEnabled(); + doReturn(true).when(rqueueSchedulerConfig).isRedisEnabled(); + doReturn(3000L).when(rqueueSchedulerConfig).getScheduledMessageTimeIntervalInMilli(); + doReturn(100L).when(rqueueSchedulerConfig).getMaxMessageCount(); + doReturn(minDelay).when(rqueueSchedulerConfig).minMessageMoveDelay(); + + Runnable messageGenerator = () -> { + int counter = 0; + while (generateMessage.get()) { + if (startGenerateMessage.get()) { + counter += 1; + byte[] currentTime = String.valueOf(System.currentTimeMillis()).getBytes(); + messageScheduler.redisScheduleTriggerHandler.messageListener.onMessage( + new DefaultMessage(channelName.getBytes(), currentTime), null); + } + // each message enqueue would lead to one event, so ~250 QPS + TimeoutUtils.sleep(4L); + } + System.out.println("Exiting sent " + counter + " messages "); + }; + scheduler.submit(messageGenerator); + AtomicInteger counter = new AtomicInteger(0); + doAnswer(invocation -> { + counter.incrementAndGet(); + sleep(5); + return System.currentTimeMillis() - Constants.DEFAULT_SCRIPT_EXECUTION_TIME; + }).when(redisTemplate).execute(any(RedisCallback.class)); + + try (MockedStatic threadUtils = Mockito.mockStatic(ThreadUtils.class)) { + threadUtils.when(() -> ThreadUtils.createTaskScheduler(1, "scheduledQueueMsgScheduler-", 60)) + .thenReturn(scheduler); + messageScheduler.onApplicationEvent(new RqueueBootstrapEvent("Test", true)); + waitFor(() -> scheduler.submittedTasks() >= 2, "one start task to be submitted"); + startGenerateMessage.set(true); + // sleep for 2 seconds -> this should send around 500 events + TimeoutUtils.sleep(totalTime); + generateMessage.set(false); // disable Redis pub/sub event + TimeoutUtils.sleep(100); + int expectedMessageMoveCalls = (int) ((totalTime / minDelay) * (1 - buffer)); + messageScheduler.destroy(); + int ranJobs = counter.get(); + int jobCount = scheduler.submittedTasks(); + log.info("Expected Job={}, Ran Jobs={}, Submitted Jobs={}", + expectedMessageMoveCalls, ranJobs, jobCount); + assertTrue(jobCount >= ranJobs); + assertTrue(jobCount >= expectedMessageMoveCalls); + } + } +} diff --git a/rqueue-core/src/test/java/com/github/sonus21/rqueue/core/RedisScheduleTriggerHandlerTest.java b/rqueue-core/src/test/java/com/github/sonus21/rqueue/core/RedisScheduleTriggerHandlerTest.java new file mode 100644 index 000000000..4c4eaa4ad --- /dev/null +++ b/rqueue-core/src/test/java/com/github/sonus21/rqueue/core/RedisScheduleTriggerHandlerTest.java @@ -0,0 +1,151 @@ +package com.github.sonus21.rqueue.core; + +import com.github.sonus21.TestBase; +import com.github.sonus21.rqueue.CoreUnitTest; +import com.github.sonus21.rqueue.config.RqueueSchedulerConfig; +import com.github.sonus21.rqueue.listener.QueueDetail; +import com.github.sonus21.rqueue.utils.TestUtils; +import com.github.sonus21.rqueue.utils.TimeoutUtils; +import lombok.extern.slf4j.Slf4j; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.mockito.Mock; +import org.mockito.MockitoAnnotations; +import org.springframework.data.redis.connection.DefaultMessage; +import org.springframework.data.redis.connection.MessageListener; + +import java.util.List; +import java.util.concurrent.Callable; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.Future; +import java.util.concurrent.atomic.AtomicInteger; +import java.util.function.Function; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.doReturn; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; + +@CoreUnitTest +@Slf4j +class RedisScheduleTriggerHandlerTest extends TestBase { + + private final ExecutorService executor = Executors.newFixedThreadPool(1); + private final String slowQueue = "slow-queue"; + private final QueueDetail slowQueueDetail = TestUtils.createQueueDetail(slowQueue); + private Scheduler scheduler; + @Mock + private RqueueSchedulerConfig rqueueSchedulerConfig; + @Mock + private RqueueRedisListenerContainerFactory rqueueRedisListenerContainerFactory; + private RedisScheduleTriggerHandler redisScheduleTriggerHandler; + private long runningTime = 0; + + + class Task implements Callable { + + @Override + public Void call() throws Exception { + log.info("Running"); + TimeoutUtils.sleep(runningTime); + return null; + } + } + + class Scheduler implements Function> { + + AtomicInteger counter = new AtomicInteger(0); + + @Override + public Future apply(String s) { + counter.incrementAndGet(); + return executor.submit(new Task()); + } + } + + @BeforeEach + public void init() { + MockitoAnnotations.openMocks(this); + EndpointRegistry.delete(); + EndpointRegistry.register(slowQueueDetail); + scheduler = new Scheduler(); + redisScheduleTriggerHandler = new RedisScheduleTriggerHandler(log, + rqueueRedisListenerContainerFactory, rqueueSchedulerConfig, + List.of(slowQueue), scheduler, + (e) -> { + return slowQueueDetail.getScheduledQueueChannelName(); + }); + redisScheduleTriggerHandler.initialize(); + redisScheduleTriggerHandler.startQueue(slowQueue); + } + + @Test + void onMessageListenerTest() throws Exception { + MessageListener messageListener = redisScheduleTriggerHandler.messageListener; + // invalid channel + messageListener.onMessage(new DefaultMessage(slowQueue.getBytes(), "312".getBytes()), null); + TimeoutUtils.sleep(50); + assertEquals(0, scheduler.counter.get()); + + // invalid body + messageListener.onMessage( + new DefaultMessage(slowQueueDetail.getScheduledQueueChannelName().getBytes(), + "sss".getBytes()), null); + TimeoutUtils.sleep(50); + assertEquals(0, scheduler.counter.get()); + + // future time + messageListener.onMessage( + new DefaultMessage(slowQueueDetail.getScheduledQueueChannelName().getBytes(), + String.valueOf(System.currentTimeMillis() + 500).getBytes()), null); + TimeoutUtils.sleep(50); + assertEquals(0, scheduler.counter.get()); + + // both are correct + messageListener.onMessage( + new DefaultMessage(slowQueueDetail.getScheduledQueueChannelName().getBytes(), + String.valueOf(System.currentTimeMillis()).getBytes()), null); + assertEquals(1, scheduler.counter.get()); + // let it run + TimeoutUtils.sleep(100); + + // send another message while one is running + runningTime = 200; + messageListener.onMessage( + new DefaultMessage(slowQueueDetail.getScheduledQueueChannelName().getBytes(), + String.valueOf(System.currentTimeMillis()).getBytes()), null); + assertEquals(2, scheduler.counter.get()); + TimeoutUtils.sleep(100); + + // this should be rejected as another task is already running + messageListener.onMessage( + new DefaultMessage(slowQueueDetail.getScheduledQueueChannelName().getBytes(), + String.valueOf(System.currentTimeMillis()).getBytes()), null); + TimeoutUtils.sleep(50); + assertEquals(2, scheduler.counter.get()); + TimeoutUtils.sleep(100); + + // this should success + long lastRunTime = System.currentTimeMillis(); + messageListener.onMessage( + new DefaultMessage(slowQueueDetail.getScheduledQueueChannelName().getBytes(), + String.valueOf(System.currentTimeMillis()).getBytes()), null); + assertEquals(3, scheduler.counter.get()); + verify(rqueueRedisListenerContainerFactory, times(1)).addMessageListener(any(), any()); + doReturn(400L).when(rqueueSchedulerConfig).getTerminationWaitTime(); + + assertEquals(1, redisScheduleTriggerHandler.queueNameToFuture.size()); + assertEquals(1, redisScheduleTriggerHandler.channelNameToQueueName.size()); + assertTrue(redisScheduleTriggerHandler.queueNameToLastRunTime.get(slowQueue) >= lastRunTime); + + redisScheduleTriggerHandler.stop(); + + assertTrue(redisScheduleTriggerHandler.queueNameToFuture.isEmpty()); + assertEquals(1, redisScheduleTriggerHandler.channelNameToQueueName.size()); + assertEquals(0L, redisScheduleTriggerHandler.queueNameToLastRunTime.get(slowQueue)); + verify(rqueueRedisListenerContainerFactory, times(1)).removeMessageListener(any(), any()); + } +} diff --git a/rqueue-core/src/test/java/com/github/sonus21/rqueue/core/ScheduledQueueMessageSchedulerTest.java b/rqueue-core/src/test/java/com/github/sonus21/rqueue/core/ScheduledQueueMessageSchedulerTest.java index 791c15032..25e91f1ed 100644 --- a/rqueue-core/src/test/java/com/github/sonus21/rqueue/core/ScheduledQueueMessageSchedulerTest.java +++ b/rqueue-core/src/test/java/com/github/sonus21/rqueue/core/ScheduledQueueMessageSchedulerTest.java @@ -41,12 +41,10 @@ import com.github.sonus21.rqueue.utils.ThreadUtils; import com.github.sonus21.rqueue.utils.TimeoutUtils; import com.github.sonus21.test.TestTaskScheduler; -import java.util.List; import java.util.Map; -import java.util.Vector; +import java.util.concurrent.Future; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicInteger; -import java.util.concurrent.atomic.AtomicLong; import org.apache.commons.lang3.reflect.FieldUtils; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; @@ -59,8 +57,6 @@ import org.springframework.data.redis.RedisConnectionFailureException; import org.springframework.data.redis.RedisSystemException; import org.springframework.data.redis.TooManyClusterRedirectionsException; -import org.springframework.data.redis.connection.DefaultMessage; -import org.springframework.data.redis.connection.MessageListener; import org.springframework.data.redis.core.RedisCallback; import org.springframework.data.redis.core.RedisTemplate; import org.springframework.data.redis.listener.ChannelTopic; @@ -107,9 +103,9 @@ void getZsetName() { void getNextScheduleTime() { long currentTime = System.currentTimeMillis(); doReturn(5000L).when(rqueueSchedulerConfig).getScheduledMessageTimeIntervalInMilli(); - assertThat(messageScheduler.getNextScheduleTime(slowQueue, null), + assertThat(messageScheduler.getNextScheduleTime(slowQueue, currentTime, null), greaterThanOrEqualTo(currentTime + 5000L)); - assertThat(messageScheduler.getNextScheduleTime(fastQueue, currentTime + 1000L), + assertThat(messageScheduler.getNextScheduleTime(fastQueue, currentTime, currentTime + 1000L), greaterThanOrEqualTo(currentTime + 5000L)); } @@ -121,9 +117,8 @@ void afterPropertiesSetWithEmptyQueSet() throws Exception { assertNull(FieldUtils.readField(messageScheduler, "scheduler", true)); assertNull(FieldUtils.readField(messageScheduler, "queueRunningState", true)); assertNull(FieldUtils.readField(messageScheduler, "queueNameToScheduledTask", true)); - assertNull(FieldUtils.readField(messageScheduler, "channelNameToQueueName", true)); - assertNull(FieldUtils.readField(messageScheduler, "queueNameToLastMessageScheduleTime", true)); - assertNull(FieldUtils.readField(messageScheduler, "queueSchedulers", true)); + assertNull(FieldUtils.readField(messageScheduler, "queueNameToNextRunTime", true)); + assertNull(FieldUtils.readField(messageScheduler, "redisScheduleTriggerHandler", true)); } @Test @@ -140,9 +135,6 @@ void start() throws Exception { assertTrue(queueRunningState.get(slowQueue)); assertEquals(2, ((Map) FieldUtils.readField(messageScheduler, "queueNameToScheduledTask", true)).size()); - assertEquals(2, - ((Map) FieldUtils.readField(messageScheduler, "channelNameToQueueName", true)).size()); - assertEquals(2, ((Map) FieldUtils.readField(messageScheduler, "queueSchedulers", true)).size()); TimeoutUtils.sleep(500L); messageScheduler.destroy(); } @@ -178,7 +170,7 @@ void stop() throws Exception { assertEquals(2, queueRunningState.size()); assertFalse(queueRunningState.get(slowQueue)); assertEquals(2, - ((Map) FieldUtils.readField(messageScheduler, "channelNameToQueueName", true)).size()); + ((Map) FieldUtils.readField(messageScheduler, "queueNameToNextRunTime", true)).size()); assertTrue( ((Map) FieldUtils.readField(messageScheduler, "queueNameToScheduledTask", true)).isEmpty()); messageScheduler.destroy(); @@ -318,115 +310,30 @@ void continuousTaskFailTask() throws Exception { waitFor(() -> counter.get() >= 10, "scripts are getting executed"); sleep(10); messageScheduler.destroy(); - assertTrue(scheduler.submittedTasks() >= 11); } } - @Test - void onMessageListenerTest() throws Exception { - doReturn(true).when(rqueueSchedulerConfig).isEnabled(); - doReturn(5000L).when(rqueueSchedulerConfig).getScheduledMessageTimeIntervalInMilli(); - doReturn(1).when(rqueueSchedulerConfig).getScheduledMessageThreadPoolSize(); - doReturn(true).when(rqueueSchedulerConfig).isAutoStart(); - doReturn(true).when(rqueueSchedulerConfig).isRedisEnabled(); - messageScheduler.onApplicationEvent(new RqueueBootstrapEvent("Test", true)); - - MessageListener messageListener = messageScheduler.messageSchedulerListener; - // invalid channel - messageListener.onMessage(new DefaultMessage(slowQueue.getBytes(), "312".getBytes()), null); - TimeoutUtils.sleep(50); - assertEquals(2, messageScheduler.scheduleList.stream().filter(e -> !e).count()); - - // invalid body - messageListener.onMessage( - new DefaultMessage(slowQueueDetail.getScheduledQueueChannelName().getBytes(), - "sss".getBytes()), null); - TimeoutUtils.sleep(50); - assertEquals(2, messageScheduler.scheduleList.stream().filter(e -> !e).count()); - - TimeoutUtils.sleep(110); - // both are correct - messageListener.onMessage( - new DefaultMessage(slowQueueDetail.getScheduledQueueChannelName().getBytes(), - String.valueOf(System.currentTimeMillis()).getBytes()), null); - - assertEquals(3, messageScheduler.scheduleList.stream().filter(e -> !e).count()); - messageScheduler.destroy(); - } - - @Test - void taskMultiplier() throws Exception { - TestTaskScheduler scheduler = new TestTaskScheduler(2); - AtomicBoolean startGenerateMessage = new AtomicBoolean(false); - AtomicBoolean generateMessage = new AtomicBoolean(true); - String channelName = messageScheduler.getChannelName(slowQueue); - doReturn(1).when(rqueueSchedulerConfig).getScheduledMessageThreadPoolSize(); - doReturn(true).when(rqueueSchedulerConfig).isAutoStart(); - doReturn(true).when(rqueueSchedulerConfig).isEnabled(); - doReturn(true).when(rqueueSchedulerConfig).isRedisEnabled(); - doReturn(3000L).when(rqueueSchedulerConfig).getScheduledMessageTimeIntervalInMilli(); - doReturn(10000L).when(rqueueSchedulerConfig).getMaxMessageMoverDelay(); - doReturn(100L).when(rqueueSchedulerConfig).getMaxMessageCount(); - - Runnable messageGenerator = () -> { - int counter = 0; - while (generateMessage.get()) { - if (startGenerateMessage.get()) { - counter += 1; - byte[] currentTime = String.valueOf(System.currentTimeMillis()).getBytes(); - messageScheduler.messageSchedulerListener.onMessage( - new DefaultMessage(channelName.getBytes(), currentTime), null); - } - // each message would enqueue would lead to one event, so ~250 QPS - TimeoutUtils.sleep(4L); - } - System.out.println("Exiting sent " + counter + " messages "); - }; - scheduler.submit(messageGenerator); - - AtomicInteger counter = new AtomicInteger(0); - doAnswer(invocation -> { - counter.incrementAndGet(); - sleep(5); - return System.currentTimeMillis() - Constants.DEFAULT_SCRIPT_EXECUTION_TIME; - }).when(redisTemplate).execute(any(RedisCallback.class)); - - try (MockedStatic threadUtils = Mockito.mockStatic(ThreadUtils.class)) { - threadUtils.when(() -> ThreadUtils.createTaskScheduler(1, "scheduledQueueMsgScheduler-", 60)) - .thenReturn(scheduler); - messageScheduler.onApplicationEvent(new RqueueBootstrapEvent("Test", true)); - waitFor(() -> scheduler.submittedTasks() >= 2, "one start task to be submitted"); - startGenerateMessage.set(true); - // sleep for 2 seconds -> this should send around 500 events - TimeoutUtils.sleep(2000); - generateMessage.set(false); // disable Redis pub/sub event - TimeoutUtils.sleep(100); - int oldCount = counter.get(); - for (int i = 0; i < 20; i++) { - TimeoutUtils.sleep(100); - System.out.println(i + "=" + counter.get()); - } - int newCount = counter.get(); - messageScheduler.destroy(); - int jobCount = scheduler.submittedTasks(); - System.out.println(oldCount); - System.out.println(newCount); - System.out.println(jobCount); - } - } static class TestScheduledQueueMessageScheduler extends ScheduledQueueMessageScheduler { - List scheduleList; + final AtomicInteger scheduleCounter; + final AtomicInteger addTaskCounter; TestScheduledQueueMessageScheduler() { - this.scheduleList = new Vector<>(); + this.scheduleCounter = new AtomicInteger(0); + this.addTaskCounter = new AtomicInteger(0); + } + + @Override + protected void schedule(String queueName) { + scheduleCounter.incrementAndGet(); + super.schedule(queueName); } @Override - protected synchronized void schedule(String queueName, Long startTime, boolean forceSchedule) { - super.schedule(queueName, startTime, forceSchedule); - this.scheduleList.add(forceSchedule); + protected Future addTask(String queueName) { + addTaskCounter.incrementAndGet(); + return super.addTask(queueName); } } } diff --git a/rqueue-test-util/src/main/java/com/github/sonus21/test/TestTaskScheduler.java b/rqueue-test-util/src/main/java/com/github/sonus21/test/TestTaskScheduler.java index 27d033bbe..106357b68 100644 --- a/rqueue-test-util/src/main/java/com/github/sonus21/test/TestTaskScheduler.java +++ b/rqueue-test-util/src/main/java/com/github/sonus21/test/TestTaskScheduler.java @@ -16,7 +16,7 @@ package com.github.sonus21.test; -import java.time.Instant; +import java.time.Duration; import java.util.List; import java.util.Vector; import java.util.concurrent.Future; @@ -46,8 +46,8 @@ public Future submit(Runnable r) { } @Override - public ScheduledFuture schedule(Runnable r, Instant instant) { - ScheduledFuture f = super.schedule(r, instant); + public ScheduledFuture scheduleAtFixedRate(Runnable task, Duration period) { + ScheduledFuture f = super.scheduleAtFixedRate(task, period); tasks.add(f); return f; } From 8e2baba0544c85e9f6e33617cbf7bdfe9084964b Mon Sep 17 00:00:00 2001 From: Sonu Kumar Date: Sat, 27 May 2023 13:55:47 +0530 Subject: [PATCH 6/9] test fixes --- .../sonus21/rqueue/core/MessageScheduler.java | 6 +-- .../core/MessageSchedulerDisabledTest.java | 9 ++-- .../rqueue/core/MessageSchedulingTest.java | 53 +++++-------------- .../core/RedisAndNormalSchedulingTest.java | 2 +- 4 files changed, 19 insertions(+), 51 deletions(-) diff --git a/rqueue-core/src/main/java/com/github/sonus21/rqueue/core/MessageScheduler.java b/rqueue-core/src/main/java/com/github/sonus21/rqueue/core/MessageScheduler.java index cd0d4cbc5..2998b6add 100644 --- a/rqueue-core/src/main/java/com/github/sonus21/rqueue/core/MessageScheduler.java +++ b/rqueue-core/src/main/java/com/github/sonus21/rqueue/core/MessageScheduler.java @@ -119,7 +119,7 @@ private void startQueue(String queueName) { return; } queueRunningState.put(queueName, true); - if (scheduleTaskAtStartup()) { + if (rqueueSchedulerConfig.isAutoStart()) { schedule(queueName); } if (isRedisEnabled()) { @@ -163,10 +163,6 @@ private void stopQueue(String queueName) { queueRunningState.put(queueName, false); } - private boolean scheduleTaskAtStartup() { - return rqueueSchedulerConfig.isAutoStart(); - } - private boolean isRedisEnabled() { return rqueueSchedulerConfig.isRedisEnabled(); } diff --git a/rqueue-core/src/test/java/com/github/sonus21/rqueue/core/MessageSchedulerDisabledTest.java b/rqueue-core/src/test/java/com/github/sonus21/rqueue/core/MessageSchedulerDisabledTest.java index 162fe17ed..6f29bb08c 100644 --- a/rqueue-core/src/test/java/com/github/sonus21/rqueue/core/MessageSchedulerDisabledTest.java +++ b/rqueue-core/src/test/java/com/github/sonus21/rqueue/core/MessageSchedulerDisabledTest.java @@ -66,6 +66,7 @@ public void init() { void startShouldSubmitsTaskWhenRedisIsDisabled() throws Exception { doReturn(1).when(rqueueSchedulerConfig).getScheduledMessageThreadPoolSize(); doReturn(true).when(rqueueSchedulerConfig).isEnabled(); + doReturn(true).when(rqueueSchedulerConfig).isAutoStart(); TestTaskScheduler scheduler = new TestTaskScheduler(); try (MockedStatic threadUtils = Mockito.mockStatic(ThreadUtils.class)) { threadUtils @@ -73,7 +74,7 @@ void startShouldSubmitsTaskWhenRedisIsDisabled() throws Exception { .thenReturn(scheduler); messageScheduler.onApplicationEvent(new RqueueBootstrapEvent("Test", true)); assertEquals(1, scheduler.submittedTasks()); - assertNull(FieldUtils.readField(messageScheduler, "messageSchedulerListener", true)); + assertNull(messageScheduler.redisScheduleTriggerHandler); messageScheduler.destroy(); } } @@ -86,8 +87,7 @@ void afterPropertiesSetProducerMode() throws Exception { assertNull(FieldUtils.readField(messageScheduler, "scheduler", true)); assertNull(FieldUtils.readField(messageScheduler, "queueRunningState", true)); assertNull(FieldUtils.readField(messageScheduler, "queueNameToScheduledTask", true)); - assertNull(FieldUtils.readField(messageScheduler, "channelNameToQueueName", true)); - assertNull(FieldUtils.readField(messageScheduler, "queueNameToLastMessageScheduleTime", true)); + assertNull(FieldUtils.readField(messageScheduler, "queueNameToNextRunTime", true)); } @Test @@ -112,7 +112,6 @@ void afterPropertiesSetDisabled() throws Exception { assertNull(FieldUtils.readField(messageScheduler, "scheduler", true)); assertNull(FieldUtils.readField(messageScheduler, "queueRunningState", true)); assertNull(FieldUtils.readField(messageScheduler, "queueNameToScheduledTask", true)); - assertNull(FieldUtils.readField(messageScheduler, "channelNameToQueueName", true)); - assertNull(FieldUtils.readField(messageScheduler, "queueNameToLastMessageScheduleTime", true)); + assertNull(FieldUtils.readField(messageScheduler, "queueNameToNextRunTime", true)); } } diff --git a/rqueue-core/src/test/java/com/github/sonus21/rqueue/core/MessageSchedulingTest.java b/rqueue-core/src/test/java/com/github/sonus21/rqueue/core/MessageSchedulingTest.java index 95523df44..e5bf012cb 100644 --- a/rqueue-core/src/test/java/com/github/sonus21/rqueue/core/MessageSchedulingTest.java +++ b/rqueue-core/src/test/java/com/github/sonus21/rqueue/core/MessageSchedulingTest.java @@ -19,6 +19,7 @@ import static com.github.sonus21.rqueue.utils.TimeoutUtils.sleep; import static com.github.sonus21.rqueue.utils.TimeoutUtils.waitFor; +import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertTrue; import static org.mockito.ArgumentMatchers.any; import static org.mockito.Mockito.doAnswer; @@ -76,39 +77,17 @@ public void init() { MockitoAnnotations.openMocks(this); EndpointRegistry.delete(); EndpointRegistry.register(queueDetail); - } - - @Test - void onCompletionOfExistingTaskNewTaskShouldBeSubmitted() throws Exception { - try (MockedStatic threadUtils = Mockito.mockStatic(ThreadUtils.class)) { - doReturn(1).when(rqueueSchedulerConfig).getProcessingMessageThreadPoolSize(); - doReturn(true).when(rqueueSchedulerConfig).isAutoStart(); - doReturn(true).when(rqueueSchedulerConfig).isEnabled(); - doReturn(true).when(rqueueSchedulerConfig).isRedisEnabled(); - AtomicInteger counter = new AtomicInteger(0); - doAnswer(invocation -> { - counter.incrementAndGet(); - return null; - }).when(redisTemplate).execute(any(RedisCallback.class)); - TestTaskScheduler scheduler = new TestTaskScheduler(); - threadUtils.when(() -> ThreadUtils.createTaskScheduler(1, "processingQueueMsgScheduler-", 60)) - .thenReturn(scheduler); - messageScheduler.onApplicationEvent(new RqueueBootstrapEvent("Test", true)); - waitFor(() -> counter.get() >= 1, "scripts are getting executed"); - sleep(10); - messageScheduler.destroy(); - assertTrue(scheduler.submittedTasks() >= 2); - } + doReturn(1).when(rqueueSchedulerConfig).getProcessingMessageThreadPoolSize(); + doReturn(200L).when(rqueueSchedulerConfig).getScheduledMessageTimeIntervalInMilli(); + doReturn(true).when(rqueueSchedulerConfig).isAutoStart(); + doReturn(true).when(rqueueSchedulerConfig).isEnabled(); + doReturn(true).when(rqueueSchedulerConfig).isRedisEnabled(); } @Test void multipleTasksAreRunningForTheSameQueue() throws Exception { try (MockedStatic threadUtils = Mockito.mockStatic(ThreadUtils.class)) { - doReturn(1).when(rqueueSchedulerConfig).getProcessingMessageThreadPoolSize(); - doReturn(true).when(rqueueSchedulerConfig).isAutoStart(); - doReturn(true).when(rqueueSchedulerConfig).isEnabled(); - doReturn(true).when(rqueueSchedulerConfig).isRedisEnabled(); AtomicInteger counter = new AtomicInteger(0); doAnswer(invocation -> { counter.incrementAndGet(); @@ -121,18 +100,15 @@ void multipleTasksAreRunningForTheSameQueue() throws Exception { waitFor(() -> counter.get() >= 2, "scripts are getting executed"); sleep(10); messageScheduler.destroy(); - assertTrue(scheduler.submittedTasks() >= 3); + assertEquals(1, scheduler.submittedTasks()); } } @Test void taskShouldBeScheduledOnFailure() throws Exception { try (MockedStatic threadUtils = Mockito.mockStatic(ThreadUtils.class)) { - doReturn(1).when(rqueueSchedulerConfig).getProcessingMessageThreadPoolSize(); - doReturn(true).when(rqueueSchedulerConfig).isAutoStart(); - doReturn(true).when(rqueueSchedulerConfig).isEnabled(); - doReturn(true).when(rqueueSchedulerConfig).isRedisEnabled(); doReturn(10000L).when(rqueueSchedulerConfig).getMaxMessageMoverDelay(); + doReturn(100L).when(rqueueSchedulerConfig).minMessageMoveDelay(); AtomicInteger counter = new AtomicInteger(0); doAnswer(invocation -> { counter.incrementAndGet(); @@ -143,21 +119,18 @@ void taskShouldBeScheduledOnFailure() throws Exception { threadUtils.when(() -> ThreadUtils.createTaskScheduler(1, "processingQueueMsgScheduler-", 60)) .thenReturn(scheduler); messageScheduler.onApplicationEvent(new RqueueBootstrapEvent("Test", true)); - waitFor(() -> counter.get() >= 2, "scripts are getting executed"); + waitFor(() -> counter.get() >= 3, "scripts are getting executed"); sleep(10); messageScheduler.destroy(); - assertTrue(scheduler.submittedTasks() >= 3); + assertEquals(1, scheduler.submittedTasks()); } } @Test void continuousTaskFailure() throws Exception { try (MockedStatic threadUtils = Mockito.mockStatic(ThreadUtils.class)) { - doReturn(1).when(rqueueSchedulerConfig).getProcessingMessageThreadPoolSize(); - doReturn(true).when(rqueueSchedulerConfig).isAutoStart(); - doReturn(true).when(rqueueSchedulerConfig).isEnabled(); - doReturn(true).when(rqueueSchedulerConfig).isRedisEnabled(); - doReturn(100L).when(rqueueSchedulerConfig).getMaxMessageMoverDelay(); + doReturn(500L).when(rqueueSchedulerConfig).getMaxMessageMoverDelay(); + doReturn(100L).when(rqueueSchedulerConfig).minMessageMoveDelay(); AtomicInteger counter = new AtomicInteger(0); doAnswer( invocation -> { @@ -182,7 +155,7 @@ void continuousTaskFailure() throws Exception { waitFor(() -> counter.get() >= 5, "scripts are getting executed"); sleep(10); messageScheduler.destroy(); - assertTrue(scheduler.submittedTasks() >= 6); + assertEquals(1, scheduler.submittedTasks()); } } diff --git a/rqueue-core/src/test/java/com/github/sonus21/rqueue/core/RedisAndNormalSchedulingTest.java b/rqueue-core/src/test/java/com/github/sonus21/rqueue/core/RedisAndNormalSchedulingTest.java index 3cb395304..45ab84db2 100644 --- a/rqueue-core/src/test/java/com/github/sonus21/rqueue/core/RedisAndNormalSchedulingTest.java +++ b/rqueue-core/src/test/java/com/github/sonus21/rqueue/core/RedisAndNormalSchedulingTest.java @@ -68,7 +68,7 @@ void redisAndNormalScheduling() throws Exception { long totalTime = 2000L; long minDelay = 10L; //15% buffer due to short polling intervals - double buffer = 0.13; + double buffer = 0.15; String channelName = messageScheduler.getChannelName(slowQueue); doReturn(1).when(rqueueSchedulerConfig).getScheduledMessageThreadPoolSize(); doReturn(true).when(rqueueSchedulerConfig).isAutoStart(); From 457a0eaf16115d35f4b0edb40db30028f3ac0b04 Mon Sep 17 00:00:00 2001 From: Sonu Kumar Date: Sat, 27 May 2023 14:01:00 +0530 Subject: [PATCH 7/9] 25% buffer --- .../sonus21/rqueue/core/RedisAndNormalSchedulingTest.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/rqueue-core/src/test/java/com/github/sonus21/rqueue/core/RedisAndNormalSchedulingTest.java b/rqueue-core/src/test/java/com/github/sonus21/rqueue/core/RedisAndNormalSchedulingTest.java index 45ab84db2..684e9ed9b 100644 --- a/rqueue-core/src/test/java/com/github/sonus21/rqueue/core/RedisAndNormalSchedulingTest.java +++ b/rqueue-core/src/test/java/com/github/sonus21/rqueue/core/RedisAndNormalSchedulingTest.java @@ -67,8 +67,8 @@ void redisAndNormalScheduling() throws Exception { AtomicBoolean generateMessage = new AtomicBoolean(true); long totalTime = 2000L; long minDelay = 10L; - //15% buffer due to short polling intervals - double buffer = 0.15; + //25% buffer due to short polling intervals, IN CI it runs slowly + double buffer = 0.25; String channelName = messageScheduler.getChannelName(slowQueue); doReturn(1).when(rqueueSchedulerConfig).getScheduledMessageThreadPoolSize(); doReturn(true).when(rqueueSchedulerConfig).isAutoStart(); From 958546d80696745f28e5899229a616903ed0920f Mon Sep 17 00:00:00 2001 From: Sonu Kumar Date: Sat, 27 May 2023 14:18:44 +0530 Subject: [PATCH 8/9] doc updated. --- docs/message-handling/producer-consumer.md | 7 +++++ .../rqueue/config/RqueueSchedulerConfig.java | 4 ++- .../spring-configuration-metadata.json | 26 +++++++++++++++++-- 3 files changed, 34 insertions(+), 3 deletions(-) diff --git a/docs/message-handling/producer-consumer.md b/docs/message-handling/producer-consumer.md index 3bac670f1..1e0d146b2 100644 --- a/docs/message-handling/producer-consumer.md +++ b/docs/message-handling/producer-consumer.md @@ -179,6 +179,13 @@ supported configurations. instances when large number of messages are scheduled to be run in next 5 minutes. In such cases Rqueue can pull message from scheduled queue to normal queue at higher rate. By default, it copies 100 messages from scheduled/processing queue to normal queue. +* `rqueue.scheduler.max.message.mover.delay=60000` Rqueue continuously move scheduled messages from + processing/scheduled queue to normal queue so that we can process them asap. Due to failure, it + can load the Redis system, in such cases it uses exponential backoff to limit the damage. This + time indicates maximum time for which it should wait before making Redis calls. +* `rqueue.scheduler.min.message.mover.delay=200` Rqueue continuously move scheduled messages from + processing/scheduled queue to normal queue so that we can process them asap. It periodically + fetches the messages, the minium delay in such cases can be configured using this variable. ### Dead Letter Queue Consumer/Listener diff --git a/rqueue-core/src/main/java/com/github/sonus21/rqueue/config/RqueueSchedulerConfig.java b/rqueue-core/src/main/java/com/github/sonus21/rqueue/config/RqueueSchedulerConfig.java index b2406f38a..328831797 100644 --- a/rqueue-core/src/main/java/com/github/sonus21/rqueue/config/RqueueSchedulerConfig.java +++ b/rqueue-core/src/main/java/com/github/sonus21/rqueue/config/RqueueSchedulerConfig.java @@ -69,7 +69,8 @@ public class RqueueSchedulerConfig { @Value("${rqueue.scheduler.scheduled.message.time.interval:2000}") private long scheduledMessageTimeIntervalInMilli; - @Value("${rqueue.scheduler.termination.wait.time:1000}") + // How long the application should wait for task termination + @Value("${rqueue.scheduler.termination.wait.time:200}") private long terminationWaitTime; // Maximum delay for message mover task due to failure @@ -80,6 +81,7 @@ public class RqueueSchedulerConfig { @Value("${rqueue.scheduler.min.message.mover.delay:100}") private long minMessageMoverDelay; + // Maximum number of messages that should be copied from scheduled to normal queue @Value("${rqueue.scheduler.max.message.count:100}") private long maxMessageCount; diff --git a/rqueue-core/src/main/resources/META-INF/spring-configuration-metadata.json b/rqueue-core/src/main/resources/META-INF/spring-configuration-metadata.json index 93a28d51f..7f7c484b2 100644 --- a/rqueue-core/src/main/resources/META-INF/spring-configuration-metadata.json +++ b/rqueue-core/src/main/resources/META-INF/spring-configuration-metadata.json @@ -263,14 +263,14 @@ { "sourceType" : "com.github.sonus21.rqueue.config.RqueueSchedulerConfig", "name" : "rqueue.scheduler.redis.enabled", - "description" : "Scheduler Should auto start or not", + "description" : "Scheduler Should use redis PUB/SUB", "type" : "java.lang.Boolean", "defaultValue" : true }, { "sourceType" : "com.github.sonus21.rqueue.config.RqueueSchedulerConfig", "name" : "rqueue.scheduler.scheduled.message.thread.pool.size", - "description" : "Number of threads should be used to pull scheduled message to original queue", + "description" : "Number of threads that will be used to pull scheduled message to normal queue", "type" : "java.lang.Integer", "defaultValue" : 3 }, @@ -295,6 +295,28 @@ "type" : "java.lang.Long", "defaultValue" : 5000 }, + { + "sourceType" : "com.github.sonus21.rqueue.config.RqueueSchedulerConfig", + "name" : "rqueue.scheduler.termination.wait.time", + "description" : "How long the application would wait for task termination", + "type" : "java.lang.Long", + "defaultValue" : 200 + }, + + { + "sourceType" : "com.github.sonus21.rqueue.config.RqueueSchedulerConfig", + "name" : "rqueue.scheduler.max.message.mover.delay", + "description" : "Maximum delay for message mover task due to failure", + "type" : "java.lang.Long", + "defaultValue" : 60000 + }, + { + "sourceType" : "com.github.sonus21.rqueue.config.RqueueSchedulerConfig", + "name" : "rqueue.scheduler.min.message.mover.delay", + "description" : "Minimum amount of time between two consecutive message move calls", + "type" : "java.lang.Long", + "defaultValue" : 100 + }, { "sourceType" : "com.github.sonus21.rqueue.config.RqueueWebConfig", "name" : "rqueue.web.enable", From a29def282e5e91f9d37506868a47956014bf1a00 Mon Sep 17 00:00:00 2001 From: Sonu Kumar Date: Wed, 31 May 2023 17:26:26 +0530 Subject: [PATCH 9/9] new user --- docs/index.md | 1 + docs/static/users/vonage.png | Bin 0 -> 6291 bytes 2 files changed, 1 insertion(+) create mode 100644 docs/static/users/vonage.png diff --git a/docs/index.md b/docs/index.md index 4f25d8c65..561df71a8 100644 --- a/docs/index.md +++ b/docs/index.md @@ -314,6 +314,7 @@ Rqueue is stable and production ready, it's processing millions of on messages d [![Line](static/users/line.png){: width="70" height="60" alt="Line Chat" }](https://line.me){:target="_blank" style="margin:10px"} [![Aviva](static/users/aviva.jpeg){: width="70" height="60" alt="Aviva" }](https://www.aviva.com/){:target="_blank" style="margin:10px"} [![Diamler Truck](static/users/mercedes.png){: width="80" height="60" alt="Daimler Truck (Mercedes)" }](https://www.daimlertruck.com/en){:target="_blank" style="margin:10px"} +[![Vonage](static/users/vonage.png){: width="250" height="60" alt="Vonage" }](http://vonage.com){:target="_blank" style="margin:10px"} [![Poker Stars](static/users/pokerstars.png){: width="210" height="60" alt="PokerStars" }](https://www.pokerstarssports.eu){:target="_blank" style="margin:10px"} [![Tune You](static/users/tuneyou.png){: width="140" height="60" alt="TuneYou" }](https://tuneyou.com){:target="_blank" style="margin:10px"} [![Bit bot](static/users/bitbot.png){: width="80" height="60" alt="BitBot" }](https://bitbot.plus){:target="_blank" style="margin:10px"} diff --git a/docs/static/users/vonage.png b/docs/static/users/vonage.png new file mode 100644 index 0000000000000000000000000000000000000000..afd0325beac64f9f3b8b820b7fc1e4d4eab5eb33 GIT binary patch literal 6291 zcmc&(^B?bm*rKnUpP${u+U3EG)Xrybag^kF@#@{A z%%^XM19yK+5s7-PY7wq#Vn-ej2(TrG z&I>no>-F28%N4Xi8BN+>ZluW`toE-F;So}nq{>`>FL6w!eXQ#`aC;#i3b(-=kpmCV z4Z2z|aW^pQ@c}>mwbwMzXlNd|RU)%H-{rfvusgv&a7Gt4Z(wB5)PG&;AMD8ScP9k~ zD0bT!YmTXI=+)P_HIAPhi{`tP1#Ty97EQ6t%E{UfCH)Ag|0HrZYD)*5fp%fk1b|4n zK*$&Z4K?+M*qyD?f&jjR-84AO*BWTkcetUZm1%11ivmZ2mjbu**dI$z3A`m-3x&h} z@(2ZVg9F~Z;KUdDbvh;6mmz-8wK4a2r*LDne}U-usl!UA>&Iqxq)(>A?b&S9b@~X# zKD7COUSt9jhpWA>d#=7EMWs8G>-yaMY)k7o9-)a^xS3X=uNkH=;rXxZ6wpcG$2(?W zzqX>%yEZe`^_$H5BCo5JHl}3nOnFg9gYqzpI1pyR`yKxk9X-GtNJ!Xw*x%c8Iw{#l zV(d^FGiP5My1m+ONTemXFYLQpQ-wG<&{DK^@S^f?b`x0`>pcqR`-{`mM zyI3NV=pyC1EG&Xyu+HN7?R;=+efK{Zsw<9}7m8>cZ@!@YO5e1^VSz0p3u?g$3m=mq^hUyT}haAZ(t;IqdAHZzJ zs4E&UY+NxzG|*S4MIM9TKaaGz66#Gtj)$YJ<5zQ`IMeR^LZh_aEPI3YPZj?;(nG)- z;O}`O3r;VV=&N%&za^{P8FhYFA;~jLQ`52Vl-ofjg~Elj{mDQ3H3tF71YUicqERc+ ztI{-hG%5aSk;z(2Vz>I$96ACgswZ^6V{wkFuAX!*PM8AY&R)!rm~e~0E%J^LYG)v$ z?}D-Yy*i;twxj1MvbSe?9pB56aD6*2*Mjj>Chs9i@ zv9kw@&uToX`L$DQ!P|@m)Y?vDr*y>AwxuYblV4Bp^;QQTbhD3syzBRwe7vJC`}-)$ z|Kp12uU^$x(<5w*bLx?CaxkofYSBeH0!*U&tOUyBiy{12DnfN$XQ%`7e4 zRUEz)vr8^&?9n4?N(mo_<}VEd;+4*Dfp(10@#Y6sT-rX&*`!HR2&D`zFM9=~Z}AC# zLb-mEvXoMkkpCj%_F25j*TQ$+An{8QfB$!%&}*iV-L<~plT(t5oVXh zM$vpU+T9nMAh5i9*Z-zM}o{Z4AQ=ZQRIES2m5Q{+sr%7LE`+xR~>%8M` zLJ4I--1~oREAk(zqa>JQPR5wYf7P7ADX`tUzQ)QM*f`0fA&allWSmSz9rb>lrSZ00 zre7zI(gnMG!{6hoU|XjI8ct35E7ZZ%*mMG3pi?TXHRe`dgxIb5e+gZ_d5O+PMQ?g- zJUy z*e$zA=U-ri#fsNysv*$Ne;=!Sn^MOWw2&vh8e>ko61{X&DSiL(sq5wq>!G_Ih11M_ z*GB3xr!E&4%gI#x!CPWkQzB|%_{|QFe_t8h=iTv5Mk^?w(*0~!=6ah6oU7NgXMgj_ zzPN#5k5V^DdUJiVY=BP$yUiFAD(wWozt&KUViQHd2Q_ZJ`=W<@c6qkhEI$JhShJ4~ zUKtV{Kf$&^9;Lgr3{&V-S=X`Lt_2^m#tI?RhsNe}$KQ@}_GU`F&Bc7R(S2))3=4a652_LNXiBOrU)-}SXs_XQW^qzi zwWgXDI*Sc~pX$VnEs@~i#QL2aM#%pDG#z$dC&$dsxS^6(_&(!|3aw{!C-K-~CH;Fz z0rr~y&>N?MQYf9MOYVjH5@|E%QZ$)l7g~Z^$ z)Pn~lge0am9?Z2dL)AxucJC=;uq#P`Il}r%j;|ZueHL73XrEEj$sQcvt||_jxrnJ9 z-wcXnUE8L5^EQ%%=u!z95rRp3PR?%j=7Yw{n^M5{s1rkM{={&`=s&kNcX(XzZdly? z@d+m0O$B+89qvx)&8=}pZucQh;Dl;XPo_jSCI%(P5j4Ig@qi0A@hitbjxZ^#X~E)o z=$FT=`>j4v@<<~QwhV~zRR7(d;6pzDwZI()dWXXZ9O~t=TWf>HPAXfmwE&y*IXx{$ zqx6@yQwC-wds2H@xifmenOqirY?#~oMT0~JOU*u%w$5ewRxFJh{A`eonNy~1f9175 z@{jD|rwd52d&u+97{bI5*1ZMTyr1r7K1w|Q9yU(>hvfsnZ=YF2)t~D<{v7qF#+Fc~ zUcV4bLSZLYpjLI^bF9b7Y5%CG5DFlVO%++tzUI zi_6`-wu7O|a|s*#2gZF+amce(byW|YlxB?;0O(6Nh6Ojqk)#o82H zF;14)XG8pHf^x8}P7n4`j9*8Jb7sax6U{^0gr$p!g=Z9BhKS~$rSo}j$8N>uJ2!Tl z^0{X&F2&D`CQdQwkd+uMf)9dzuWkm}d1Y+O4gKJ(%uOENOt!S}wPs!qyS-*z2ytPM z@O*qb%x39R@x5cV>``+DH*iJ$>4|c6zL(?Z@5tjm;(T)KJ4}qjmu8?N>Tg{O+PgQ8 ze3YVIP7ID#h{ogNeyk3^>u^X4HcX?FZioGfWnURgomLX6)O6uXR@wx2W0G33)kY!E zw|G?)u>H*a=J$3M)txixms}P$>;X)Tl&}o5SD&si?wWjXi)+Y+bPnRJlp<52t28Zs zX{Uq^_c|G8r;(j7mEX|oH0$+i%#k#U3zH~+FoQ{2*@cl*kY?I=e{cZXia5#ML1%>s#p|y@yXI+C;ix_bIY41s%E@fS}UNH z_h-P>(O1JM%1Hl!mWm3v4UZG9rk$hU4`(w_Zofd+SNcY>NWJ)eMd_UECqS1$uj?OG zTBcJj34R)#B&NNsMKoDX#KGfu^daj$D?PDk_P&)9MAJNAai3|7zM8HvwPo9p-?b#4 zJ!Q~I^s`FIksWBLVcyVYD4N>rZPs|Jt=jq^40&#um-!X?W7j@$L5TI0C#$7LTPTso z%)nN*2|lxBmgoF%uIe(}OaLEO|Lk35vpWThg_CPaGNao^MThg=A5Q)rbO?JMoiVEAEIawCg$i|O zQ<2S9ObNZ_0OJp#0>qgm{T7}+)%WgZR?&+pDGS6;h32}p)R_Nmg`B(90b>O~&k-ckDm1s5RD+p*g{apz3y>0?@7Zd%#QAuq=>1oG-N*qnodNJ`H; zYHyNs_!8iyAIaJemFH{C4CjHX2&bgV2|eS8k5959;jA+mK#((6v+}c7-*m`{eTK6A-Th@*8@gDX2`7C%r_d?e~1!I`9Kx+>)f2yP3r@)1> zBpP(18KT1x##1NXzIjQj&LQSSH0Es+mtZqk5KnH5j-1t4St`BfG{xA9*@tOe!zxl? zY?F8CYVq2kB6Yb|@H<3J18aE31oN#&VXgXP-4i2B6^%bF@WPJ%t2ED zmscjqR~M|+xh2H_iU^sMu~r|W(>0T3g=qd?GvKn8GnuS=?BXN-HCw)qgE!^mermVFc~m2HY{g3GL%4e=UdV6 zvgUdt}Xm>6imj%UMV0 zH1mX5!D~!#y!m8jWf6KnDtSEi{l+(xBIh`o?&*y_;IGw{z^{RBw1kRaP2ML06gR}G- zWbb}^Z&M?HF6vxb=ENooZSUpvpxR0}IJQq+jAu>`7$0+7eI<&lbsUTaFMLC>6fSNp zA2F_@0&*K)hmx1C%)Jc9&eIdfPyi7VhmUd`Z;<)s;Eq;AYt$p0Yi~7*q#k?H%(2={ ze+Znb)=9(@$tI)~*Vd70g?}}lwmq6jp~?#XL4Omio}osDU0^yjCM)Q&)Lh-1M`)oh zpfWG-Mu0q`eHiq;RA)!Qp5+fa1a_H84)BXEAiM9b%2L_XP}7;a(VnRW?JCNE_jD&A zEG$wGB&io??&=O?Jz{|;%mRikfJ!Zsb#>cN34roE>_M3($W9)Q#zfq6S6^3uje&5P zWzL-Cx#D}K-i?tIbwF=)IF4Vp1U6x|cMd~J{Z|*?$<}4VKWhWe@5&1^$L^RcPYbpV zQAZ5XPsrV8yI_OB!E{Z**z{7OTMn))DP16T7IUN&RNp`}AvGLvPfRsgG1{8MS~swS zzB7f^ zo}O#AV;ueN?Tetc9{!6|>CM2Yjl?)a?Czo#DBi+Mf9W{p{wir|r1;znqCih#9UF!*dbb+9@#_dF z7@4YVHd}q3aXL20zOWSV3G!k%jVIkgKXTF+F%ZNezCV@9KVWJ8=r7?42DM^20u;mk zm-sDNb77dXZ{3(U1_6Z7M8vZFEqlipl^U|cu>WsI zIP*sqaKP?djnn2>#xCYcY=i+^#|%2lTjfmOkHu@