mirror of
https://github.com/apache/rocketmq.git
synced 2026-09-21 05:44:03 +08:00
fix scheduleAtFixedRate bug (#3861)
This commit is contained in:
+5
-1
@@ -93,7 +93,11 @@ public class ConsumeMessageConcurrentlyService implements ConsumeMessageService
|
||||
|
||||
@Override
|
||||
public void run() {
|
||||
cleanExpireMsg();
|
||||
try {
|
||||
cleanExpireMsg();
|
||||
} catch (Throwable e) {
|
||||
log.error("scheduleAtFixedRate cleanExpireMsg exception", e);
|
||||
}
|
||||
}
|
||||
|
||||
}, this.defaultMQPushConsumer.getConsumeTimeout(), this.defaultMQPushConsumer.getConsumeTimeout(), TimeUnit.MINUTES);
|
||||
|
||||
+5
-1
@@ -97,7 +97,11 @@ public class ConsumeMessageOrderlyService implements ConsumeMessageService {
|
||||
this.scheduledExecutorService.scheduleAtFixedRate(new Runnable() {
|
||||
@Override
|
||||
public void run() {
|
||||
ConsumeMessageOrderlyService.this.lockMQPeriodically();
|
||||
try {
|
||||
ConsumeMessageOrderlyService.this.lockMQPeriodically();
|
||||
} catch (Throwable e) {
|
||||
log.error("scheduleAtFixedRate lockMQPeriodically exception", e);
|
||||
}
|
||||
}
|
||||
}, 1000 * 1, ProcessQueue.REBALANCE_LOCK_INTERVAL, TimeUnit.MILLISECONDS);
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user