mirror of
https://github.com/apache/rocketmq.git
synced 2026-09-24 16:04:00 +08:00
This commit is contained in:
@@ -27,6 +27,8 @@ import java.util.concurrent.locks.ReadWriteLock;
|
||||
import java.util.concurrent.locks.ReentrantLock;
|
||||
import java.util.concurrent.locks.ReentrantReadWriteLock;
|
||||
|
||||
import org.apache.commons.lang3.StringUtils;
|
||||
|
||||
import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer;
|
||||
import org.apache.rocketmq.client.log.ClientLogger;
|
||||
import org.apache.rocketmq.logging.InternalLogger;
|
||||
@@ -85,10 +87,14 @@ public class ProcessQueue {
|
||||
try {
|
||||
this.treeMapLock.readLock().lockInterruptibly();
|
||||
try {
|
||||
if (!msgTreeMap.isEmpty() && System.currentTimeMillis() - Long.parseLong(MessageAccessor.getConsumeStartTimeStamp(msgTreeMap.firstEntry().getValue())) > pushConsumer.getConsumeTimeout() * 60 * 1000) {
|
||||
msg = msgTreeMap.firstEntry().getValue();
|
||||
if (!msgTreeMap.isEmpty()) {
|
||||
String consumeStartTimeStamp = MessageAccessor.getConsumeStartTimeStamp(msgTreeMap.firstEntry().getValue());
|
||||
if (StringUtils.isNotEmpty(consumeStartTimeStamp) && System.currentTimeMillis() - Long.parseLong(consumeStartTimeStamp) > pushConsumer.getConsumeTimeout() * 60 * 1000) {
|
||||
msg = msgTreeMap.firstEntry().getValue();
|
||||
} else {
|
||||
break;
|
||||
}
|
||||
} else {
|
||||
|
||||
break;
|
||||
}
|
||||
} finally {
|
||||
|
||||
Reference in New Issue
Block a user