Browse Source

queue consumer: not going to sleep after pull if time left less then 1 millisecond. topic added to logs

pull/4520/head
Sergey Matvienko 6 years ago
committed by Andrew Shvayka
parent
commit
5ad113ad1a
  1. 12
      common/queue/src/main/java/org/thingsboard/server/queue/common/AbstractTbQueueConsumerTemplate.java

12
common/queue/src/main/java/org/thingsboard/server/queue/common/AbstractTbQueueConsumerTemplate.java

@ -38,6 +38,7 @@ import static java.util.Collections.emptyList;
@Slf4j @Slf4j
public abstract class AbstractTbQueueConsumerTemplate<R, T extends TbQueueMsg> implements TbQueueConsumer<T> { public abstract class AbstractTbQueueConsumerTemplate<R, T extends TbQueueMsg> implements TbQueueConsumer<T> {
public static final long ONE_MILLISECOND_IN_NANOS = TimeUnit.MILLISECONDS.toNanos(1);
private volatile boolean subscribed; private volatile boolean subscribed;
protected volatile boolean stopped = false; protected volatile boolean stopped = false;
protected volatile Set<TopicPartitionInfo> partitions; protected volatile Set<TopicPartitionInfo> partitions;
@ -83,7 +84,7 @@ public abstract class AbstractTbQueueConsumerTemplate<R, T extends TbQueueMsg> i
} }
if (consumerLock.isLocked()) { if (consumerLock.isLocked()) {
log.error("poll. consumerLock is locked. will wait with no timeout. it looks like a race conditions or deadlock", new RuntimeException("stacktrace")); log.error("poll. consumerLock is locked. will wait with no timeout. it looks like a race conditions or deadlock topic " + topic, new RuntimeException("stacktrace"));
} }
consumerLock.lock(); consumerLock.lock();
@ -131,9 +132,12 @@ public abstract class AbstractTbQueueConsumerTemplate<R, T extends TbQueueMsg> i
List<T> sleepAndReturnEmpty(final long startNanos, final long durationInMillis) { List<T> sleepAndReturnEmpty(final long startNanos, final long durationInMillis) {
long durationNanos = TimeUnit.MILLISECONDS.toNanos(durationInMillis); long durationNanos = TimeUnit.MILLISECONDS.toNanos(durationInMillis);
long spentNanos = System.nanoTime() - startNanos; long spentNanos = System.nanoTime() - startNanos;
if (spentNanos < durationNanos) { long nanosLeft = durationNanos - spentNanos;
if (nanosLeft >= ONE_MILLISECOND_IN_NANOS) {
try { try {
Thread.sleep(Math.max(TimeUnit.NANOSECONDS.toMillis(durationNanos - spentNanos), 1)); long sleepMs = TimeUnit.NANOSECONDS.toMillis(nanosLeft);
log.trace("Going to sleep after poll: topic {} for {}ms", topic, sleepMs);
Thread.sleep(sleepMs);
} catch (InterruptedException e) { } catch (InterruptedException e) {
if (!stopped) { if (!stopped) {
log.error("Failed to wait", e); log.error("Failed to wait", e);
@ -146,7 +150,7 @@ public abstract class AbstractTbQueueConsumerTemplate<R, T extends TbQueueMsg> i
@Override @Override
public void commit() { public void commit() {
if (consumerLock.isLocked()) { if (consumerLock.isLocked()) {
log.error("commit. consumerLock is locked. will wait with no timeout. it looks like a race conditions or deadlock", new RuntimeException("stacktrace")); log.error("commit. consumerLock is locked. will wait with no timeout. it looks like a race conditions or deadlock topic " + topic, new RuntimeException("stacktrace"));
} }
consumerLock.lock(); consumerLock.lock();
try { try {

Loading…
Cancel
Save