Browse Source

Introduce PartitionedQueueResponseTemplate with consumer per partition; use it for EDQS requests processing

pull/13362/head
ViacheslavKlimov 1 year ago
parent
commit
f8bf512a0a
  1. 2
      application/src/main/resources/thingsboard.yml
  2. 3
      common/cluster-api/src/main/java/org/thingsboard/server/queue/TbQueueProducer.java
  3. 46
      common/edqs/src/main/java/org/thingsboard/server/edqs/processor/EdqsProcessor.java
  4. 3
      common/edqs/src/main/java/org/thingsboard/server/edqs/state/EdqsStateService.java
  5. 14
      common/edqs/src/main/java/org/thingsboard/server/edqs/state/KafkaEdqsStateService.java
  6. 8
      common/edqs/src/main/java/org/thingsboard/server/edqs/state/LocalEdqsStateService.java
  7. 1
      common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java
  8. 5
      common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueResponseTemplate.java
  9. 164
      common/queue/src/main/java/org/thingsboard/server/queue/common/PartitionedQueueResponseTemplate.java
  10. 5
      common/queue/src/main/java/org/thingsboard/server/queue/common/consumer/MainQueueConsumerManager.java
  11. 3
      common/queue/src/main/java/org/thingsboard/server/queue/common/consumer/PartitionedQueueConsumerManager.java
  12. 4
      common/queue/src/main/java/org/thingsboard/server/queue/common/state/DefaultQueueStateService.java
  13. 8
      common/queue/src/main/java/org/thingsboard/server/queue/common/state/KafkaQueueStateService.java
  14. 14
      common/queue/src/main/java/org/thingsboard/server/queue/common/state/QueueStateService.java
  15. 2
      common/queue/src/main/java/org/thingsboard/server/queue/edqs/EdqsConfig.java
  16. 70
      common/queue/src/main/java/org/thingsboard/server/queue/edqs/EdqsExecutors.java
  17. 5
      common/queue/src/main/java/org/thingsboard/server/queue/edqs/EdqsQueueFactory.java
  18. 26
      common/queue/src/main/java/org/thingsboard/server/queue/edqs/InMemoryEdqsQueueFactory.java
  19. 35
      common/queue/src/main/java/org/thingsboard/server/queue/edqs/KafkaEdqsQueueFactory.java
  20. 4
      common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaProducerTemplate.java
  21. 5
      common/queue/src/main/java/org/thingsboard/server/queue/memory/InMemoryTbQueueProducer.java
  22. 19
      common/queue/src/test/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplateTest.java
  23. 2
      edqs/src/main/resources/edqs.yml

2
application/src/main/resources/thingsboard.yml

@ -1778,6 +1778,8 @@ queue:
max_pending_requests: "${TB_EDQS_MAX_PENDING_REQUESTS:10000}"
# Maximum timeout for requests to EDQS
max_request_timeout: "${TB_EDQS_MAX_REQUEST_TIMEOUT:20000}"
# Thread pool size for EDQS requests executor
request_executor_size: "${TB_EDQS_REQUEST_EXECUTOR_SIZE:50}"
# Strings longer than this threshold will be compressed
string_compression_length_threshold: "${TB_EDQS_STRING_COMPRESSION_LENGTH_THRESHOLD:512}"
stats:

3
common/cluster-api/src/main/java/org/thingsboard/server/queue/TbQueueProducer.java

@ -19,11 +19,10 @@ import org.thingsboard.server.common.msg.queue.TopicPartitionInfo;
public interface TbQueueProducer<T extends TbQueueMsg> {
void init();
String getDefaultTopic();
void send(TopicPartitionInfo tpi, T msg, TbQueueCallback callback);
void stop();
}

46
common/edqs/src/main/java/org/thingsboard/server/edqs/processor/EdqsProcessor.java

@ -18,7 +18,6 @@ package org.thingsboard.server.edqs.processor;
import com.google.common.collect.Sets;
import com.google.common.util.concurrent.ListenableFuture;
import com.google.common.util.concurrent.ListeningExecutorService;
import com.google.common.util.concurrent.MoreExecutors;
import jakarta.annotation.PostConstruct;
import jakarta.annotation.PreDestroy;
import lombok.Getter;
@ -29,7 +28,6 @@ import org.springframework.context.event.EventListener;
import org.springframework.stereotype.Service;
import org.thingsboard.common.util.ExceptionUtil;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.common.util.ThingsBoardExecutors;
import org.thingsboard.common.util.ThingsBoardThreadFactory;
import org.thingsboard.server.common.data.ObjectType;
import org.thingsboard.server.common.data.edqs.EdqsEvent;
@ -54,7 +52,7 @@ import org.thingsboard.server.gen.transport.TransportProtos.EdqsEventMsg;
import org.thingsboard.server.gen.transport.TransportProtos.FromEdqsMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToEdqsMsg;
import org.thingsboard.server.queue.TbQueueHandler;
import org.thingsboard.server.queue.TbQueueResponseTemplate;
import org.thingsboard.server.queue.common.PartitionedQueueResponseTemplate;
import org.thingsboard.server.queue.common.TbProtoQueueMsg;
import org.thingsboard.server.queue.common.consumer.PartitionedQueueConsumerManager;
import org.thingsboard.server.queue.discovery.QueueKey;
@ -62,15 +60,14 @@ import org.thingsboard.server.queue.discovery.event.PartitionChangeEvent;
import org.thingsboard.server.queue.edqs.EdqsComponent;
import org.thingsboard.server.queue.edqs.EdqsConfig;
import org.thingsboard.server.queue.edqs.EdqsConfig.EdqsPartitioningStrategy;
import org.thingsboard.server.queue.edqs.EdqsExecutors;
import org.thingsboard.server.queue.edqs.EdqsQueueFactory;
import org.thingsboard.server.queue.util.AfterStartUp;
import java.util.List;
import java.util.Objects;
import java.util.Set;
import java.util.UUID;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.function.Consumer;
import java.util.stream.Collectors;
@ -87,20 +84,16 @@ public class EdqsProcessor implements TbQueueHandler<TbProtoQueueMsg<ToEdqsMsg>,
private final EdqsConverter converter;
private final EdqsRepository repository;
private final EdqsConfig config;
private final EdqsExecutors edqsExecutors;
private final EdqsPartitionService partitionService;
private final ConfigurableApplicationContext applicationContext;
private final EdqsStateService stateService;
private PartitionedQueueConsumerManager<TbProtoQueueMsg<ToEdqsMsg>> eventConsumer;
private TbQueueResponseTemplate<TbProtoQueueMsg<ToEdqsMsg>, TbProtoQueueMsg<FromEdqsMsg>> responseTemplate;
private ExecutorService consumersExecutor;
private ExecutorService taskExecutor;
private ScheduledExecutorService scheduler;
private PartitionedQueueResponseTemplate<TbProtoQueueMsg<ToEdqsMsg>, TbProtoQueueMsg<FromEdqsMsg>> responseTemplate;
private ListeningExecutorService requestExecutor;
private final VersionsStore versionsStore = new VersionsStore();
private final AtomicInteger counter = new AtomicInteger();
@Getter
@ -108,10 +101,6 @@ public class EdqsProcessor implements TbQueueHandler<TbProtoQueueMsg<ToEdqsMsg>,
@PostConstruct
private void init() {
consumersExecutor = Executors.newCachedThreadPool(ThingsBoardThreadFactory.forName("edqs-consumer"));
taskExecutor = ThingsBoardExecutors.newWorkStealingPool(4, "edqs-consumer-task-executor");
scheduler = ThingsBoardExecutors.newSingleThreadScheduledExecutor("edqs-scheduler");
requestExecutor = MoreExecutors.listeningDecorator(ThingsBoardExecutors.newWorkStealingPool(12, "edqs-requests"));
errorHandler = error -> {
if (error instanceof OutOfMemoryError) {
log.error("OOM detected, shutting down");
@ -120,6 +109,7 @@ public class EdqsProcessor implements TbQueueHandler<TbProtoQueueMsg<ToEdqsMsg>,
.execute(applicationContext::close);
}
};
requestExecutor = edqsExecutors.getRequestExecutor();
eventConsumer = PartitionedQueueConsumerManager.<TbProtoQueueMsg<ToEdqsMsg>>create()
.queueKey(new QueueKey(ServiceType.EDQS, config.getEventsTopic()))
@ -141,19 +131,14 @@ public class EdqsProcessor implements TbQueueHandler<TbProtoQueueMsg<ToEdqsMsg>,
})
.consumerCreator((config, tpi) -> queueFactory.createEdqsEventsConsumer())
.queueAdmin(queueFactory.getEdqsQueueAdmin())
.consumerExecutor(consumersExecutor)
.taskExecutor(taskExecutor)
.scheduler(scheduler)
.consumerExecutor(edqsExecutors.getConsumersExecutor())
.taskExecutor(edqsExecutors.getConsumerTaskExecutor())
.scheduler(edqsExecutors.getScheduler())
.uncaughtErrorHandler(errorHandler)
.build();
stateService.init(eventConsumer);
responseTemplate = queueFactory.createEdqsResponseTemplate();
}
responseTemplate = queueFactory.createEdqsResponseTemplate(this);
@AfterStartUp(order = 1)
public void start() {
responseTemplate.launch(this);
stateService.init(eventConsumer, List.of(responseTemplate.getRequestConsumer()));
}
@EventListener
@ -163,10 +148,8 @@ public class EdqsProcessor implements TbQueueHandler<TbProtoQueueMsg<ToEdqsMsg>,
}
try {
Set<TopicPartitionInfo> newPartitions = event.getNewPartitions().get(new QueueKey(ServiceType.EDQS));
stateService.process(withTopic(newPartitions, config.getStateTopic()));
// eventsConsumer's partitions are updated by stateService
responseTemplate.subscribe(withTopic(newPartitions, config.getRequestsTopic())); // TODO: we subscribe to partitions before we are ready. implement consumer-per-partition version for request template
// partitions for event and request consumers are updated by stateService
Set<TopicPartitionInfo> oldPartitions = event.getOldPartitions().get(new QueueKey(ServiceType.EDQS));
if (CollectionsUtil.isNotEmpty(oldPartitions)) {
@ -290,11 +273,6 @@ public class EdqsProcessor implements TbQueueHandler<TbProtoQueueMsg<ToEdqsMsg>,
eventConsumer.awaitStop();
responseTemplate.stop();
stateService.stop();
consumersExecutor.shutdownNow();
taskExecutor.shutdownNow();
scheduler.shutdownNow();
requestExecutor.shutdownNow();
}
}

3
common/edqs/src/main/java/org/thingsboard/server/edqs/state/EdqsStateService.java

@ -23,11 +23,12 @@ import org.thingsboard.server.gen.transport.TransportProtos.ToEdqsMsg;
import org.thingsboard.server.queue.common.TbProtoQueueMsg;
import org.thingsboard.server.queue.common.consumer.PartitionedQueueConsumerManager;
import java.util.List;
import java.util.Set;
public interface EdqsStateService {
void init(PartitionedQueueConsumerManager<TbProtoQueueMsg<ToEdqsMsg>> eventConsumer);
void init(PartitionedQueueConsumerManager<TbProtoQueueMsg<ToEdqsMsg>> eventConsumer, List<PartitionedQueueConsumerManager<?>> otherConsumers);
void process(Set<TopicPartitionInfo> partitions);

14
common/edqs/src/main/java/org/thingsboard/server/edqs/state/KafkaEdqsStateService.java

@ -37,12 +37,14 @@ import org.thingsboard.server.queue.common.state.KafkaQueueStateService;
import org.thingsboard.server.queue.common.state.QueueStateService;
import org.thingsboard.server.queue.discovery.QueueKey;
import org.thingsboard.server.queue.edqs.EdqsConfig;
import org.thingsboard.server.queue.edqs.EdqsExecutors;
import org.thingsboard.server.queue.edqs.KafkaEdqsComponent;
import org.thingsboard.server.queue.edqs.KafkaEdqsQueueFactory;
import org.thingsboard.server.queue.kafka.TbKafkaAdmin;
import org.thingsboard.server.queue.kafka.TbKafkaConsumerTemplate;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.UUID;
@ -59,6 +61,7 @@ public class KafkaEdqsStateService implements EdqsStateService {
private final EdqsConfig config;
private final EdqsPartitionService partitionService;
private final KafkaEdqsQueueFactory queueFactory;
private final EdqsExecutors edqsExecutors;
@Autowired
@Lazy
private EdqsProcessor edqsProcessor;
@ -74,7 +77,7 @@ public class KafkaEdqsStateService implements EdqsStateService {
private Boolean ready;
@Override
public void init(PartitionedQueueConsumerManager<TbProtoQueueMsg<ToEdqsMsg>> eventConsumer) {
public void init(PartitionedQueueConsumerManager<TbProtoQueueMsg<ToEdqsMsg>> eventConsumer, List<PartitionedQueueConsumerManager<?>> otherConsumers) {
TbKafkaAdmin queueAdmin = queueFactory.getEdqsQueueAdmin();
stateConsumer = PartitionedQueueConsumerManager.<TbProtoQueueMsg<ToEdqsMsg>>create()
.queueKey(new QueueKey(ServiceType.EDQS, config.getStateTopic()))
@ -96,9 +99,9 @@ public class KafkaEdqsStateService implements EdqsStateService {
})
.consumerCreator((config, tpi) -> queueFactory.createEdqsStateConsumer())
.queueAdmin(queueAdmin)
.consumerExecutor(eventConsumer.getConsumerExecutor())
.taskExecutor(eventConsumer.getTaskExecutor())
.scheduler(eventConsumer.getScheduler())
.consumerExecutor(edqsExecutors.getConsumersExecutor())
.taskExecutor(edqsExecutors.getConsumerTaskExecutor())
.scheduler(edqsExecutors.getScheduler())
.uncaughtErrorHandler(edqsProcessor.getErrorHandler())
.build();
@ -141,7 +144,7 @@ public class KafkaEdqsStateService implements EdqsStateService {
consumer.commit();
})
.consumerCreator(() -> eventsToBackupKafkaConsumer)
.consumerExecutor(eventConsumer.getConsumerExecutor())
.consumerExecutor(edqsExecutors.getConsumersExecutor())
.threadPrefix("edqs-events-to-backup")
.build();
@ -153,6 +156,7 @@ public class KafkaEdqsStateService implements EdqsStateService {
queueStateService = KafkaQueueStateService.<TbProtoQueueMsg<ToEdqsMsg>, TbProtoQueueMsg<ToEdqsMsg>>builder()
.eventConsumer(eventConsumer)
.stateConsumer(stateConsumer)
.otherConsumers(otherConsumers)
.eventsStartOffsetsProvider(() -> {
// taking start offsets for events topics from the events-to-backup consumer group,
// since eventConsumer doesn't use consumer group management and thus offset tracking

8
common/edqs/src/main/java/org/thingsboard/server/edqs/state/LocalEdqsStateService.java

@ -31,6 +31,7 @@ import org.thingsboard.server.queue.common.TbProtoQueueMsg;
import org.thingsboard.server.queue.common.consumer.PartitionedQueueConsumerManager;
import org.thingsboard.server.queue.edqs.InMemoryEdqsComponent;
import java.util.List;
import java.util.Set;
import static org.thingsboard.server.common.msg.queue.TopicPartitionInfo.withTopic;
@ -46,11 +47,13 @@ public class LocalEdqsStateService implements EdqsStateService {
private EdqsProcessor processor;
private PartitionedQueueConsumerManager<TbProtoQueueMsg<ToEdqsMsg>> eventConsumer;
private List<PartitionedQueueConsumerManager<?>> otherConsumers;
private Set<TopicPartitionInfo> partitions;
@Override
public void init(PartitionedQueueConsumerManager<TbProtoQueueMsg<ToEdqsMsg>> eventConsumer) {
public void init(PartitionedQueueConsumerManager<TbProtoQueueMsg<ToEdqsMsg>> eventConsumer, List<PartitionedQueueConsumerManager<?>> otherConsumers) {
this.eventConsumer = eventConsumer;
this.otherConsumers = otherConsumers;
}
@Override
@ -68,6 +71,9 @@ public class LocalEdqsStateService implements EdqsStateService {
log.info("Restore completed");
}
eventConsumer.update(withTopic(partitions, eventConsumer.getTopic()));
for (PartitionedQueueConsumerManager<?> consumer : otherConsumers) {
consumer.update(withTopic(partitions, consumer.getTopic()));
}
this.partitions = partitions;
}

1
common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java

@ -91,7 +91,6 @@ public class DefaultTbQueueRequestTemplate<Request extends TbQueueMsg, Response
@Override
public void init() {
queueAdmin.createTopicIfNotExists(responseTemplate.getTopic());
requestTemplate.init();
responseTemplate.subscribe();
executor.submit(this::mainLoop);
}

5
common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueResponseTemplate.java

@ -30,8 +30,6 @@ import org.thingsboard.server.queue.TbQueueResponseTemplate;
import java.util.List;
import java.util.Set;
import java.util.UUID;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
@ -44,7 +42,6 @@ public class DefaultTbQueueResponseTemplate<Request extends TbQueueMsg, Response
private final TbQueueConsumer<Request> requestTemplate;
private final TbQueueProducer<Response> responseTemplate;
private final ConcurrentMap<UUID, String> pendingRequests;
private final ExecutorService loopExecutor;
private final ScheduledExecutorService timeoutExecutor;
private final ExecutorService callbackExecutor;
@ -67,7 +64,6 @@ public class DefaultTbQueueResponseTemplate<Request extends TbQueueMsg, Response
MessagesStats stats) {
this.requestTemplate = requestTemplate;
this.responseTemplate = responseTemplate;
this.pendingRequests = new ConcurrentHashMap<>();
this.maxPendingRequests = maxPendingRequests;
this.pollInterval = pollInterval;
this.requestTimeout = requestTimeout;
@ -89,7 +85,6 @@ public class DefaultTbQueueResponseTemplate<Request extends TbQueueMsg, Response
@Override
public void launch(TbQueueHandler<Request, Response> handler) {
this.responseTemplate.init();
loopExecutor.submit(() -> {
while (!stopped) {
try {

164
common/queue/src/main/java/org/thingsboard/server/queue/common/PartitionedQueueResponseTemplate.java

@ -0,0 +1,164 @@
/**
* Copyright © 2016-2025 The Thingsboard Authors
*
* 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
*
* http://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 org.thingsboard.server.queue.common;
import lombok.Builder;
import lombok.Getter;
import lombok.extern.slf4j.Slf4j;
import org.thingsboard.common.util.ThingsBoardExecutors;
import org.thingsboard.server.common.msg.queue.TopicPartitionInfo;
import org.thingsboard.server.common.stats.MessagesStats;
import org.thingsboard.server.queue.TbQueueConsumer;
import org.thingsboard.server.queue.TbQueueHandler;
import org.thingsboard.server.queue.TbQueueMsg;
import org.thingsboard.server.queue.TbQueueProducer;
import org.thingsboard.server.queue.common.consumer.PartitionedQueueConsumerManager;
import java.util.List;
import java.util.Set;
import java.util.UUID;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeoutException;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.function.Function;
@Slf4j
public class PartitionedQueueResponseTemplate<Request extends TbQueueMsg, Response extends TbQueueMsg> extends AbstractTbQueueTemplate {
@Getter
private final PartitionedQueueConsumerManager<Request> requestConsumer;
private final TbQueueProducer<Response> responseProducer;
private final TbQueueHandler<Request, Response> handler;
private final long pollInterval;
private final int maxPendingRequests;
private final long requestTimeout;
private final MessagesStats stats;
private final ScheduledExecutorService scheduler;
private final ExecutorService callbackExecutor;
private final AtomicInteger pendingRequestCount = new AtomicInteger();
@Builder
public PartitionedQueueResponseTemplate(String key,
TbQueueHandler<Request, Response> handler,
String requestsTopic,
Function<TopicPartitionInfo, TbQueueConsumer<Request>> consumerCreator,
TbQueueProducer<Response> responseProducer,
long pollInterval,
long requestTimeout,
int maxPendingRequests,
ExecutorService consumerExecutor,
ExecutorService callbackExecutor,
ExecutorService consumerTaskExecutor,
MessagesStats stats) {
this.scheduler = ThingsBoardExecutors.newSingleThreadScheduledExecutor(key + "-queue-response-template-scheduler");
this.callbackExecutor = callbackExecutor;
this.handler = handler;
this.requestConsumer = PartitionedQueueConsumerManager.<Request>create()
.queueKey(key + "-requests")
.topic(requestsTopic)
.pollInterval(pollInterval)
.msgPackProcessor((requests, consumer, config) -> processRequests(requests, consumer))
.consumerCreator((config, tpi) -> consumerCreator.apply(tpi))
.consumerExecutor(consumerExecutor)
.scheduler(scheduler)
.taskExecutor(consumerTaskExecutor)
.build();
this.responseProducer = responseProducer;
this.pollInterval = pollInterval;
this.maxPendingRequests = maxPendingRequests;
this.requestTimeout = requestTimeout;
this.stats = stats;
}
private void processRequests(List<Request> requests, TbQueueConsumer<Request> consumer) {
while (pendingRequestCount.get() >= maxPendingRequests) {
try {
Thread.sleep(pollInterval);
} catch (InterruptedException e) {
log.trace("Failed to wait until the server has capacity to handle new requests", e);
}
}
requests.forEach(request -> {
long currentTime = System.currentTimeMillis();
long expireTs = bytesToLong(request.getHeaders().get(EXPIRE_TS_HEADER));
if (expireTs >= currentTime) {
byte[] requestIdHeader = request.getHeaders().get(REQUEST_ID_HEADER);
if (requestIdHeader == null) {
log.error("[{}] Missing requestId in header", request);
return;
}
byte[] responseTopicHeader = request.getHeaders().get(RESPONSE_TOPIC_HEADER);
if (responseTopicHeader == null) {
log.error("[{}] Missing response topic in header", request);
return;
}
UUID requestId = bytesToUuid(requestIdHeader);
String responseTopic = bytesToString(responseTopicHeader);
try {
pendingRequestCount.getAndIncrement();
stats.incrementTotal();
AsyncCallbackTemplate.withCallbackAndTimeout(handler.handle(request),
response -> {
pendingRequestCount.decrementAndGet();
response.getHeaders().put(REQUEST_ID_HEADER, uuidToBytes(requestId));
responseProducer.send(TopicPartitionInfo.builder().topic(responseTopic).build(), response, null);
stats.incrementSuccessful();
},
e -> {
pendingRequestCount.decrementAndGet();
if (e.getCause() != null && e.getCause() instanceof TimeoutException) {
log.warn("[{}] Timeout to process the request: {}", requestId, request, e);
} else {
log.trace("[{}] Failed to process the request: {}", requestId, request, e);
}
stats.incrementFailed();
},
requestTimeout,
scheduler,
callbackExecutor);
} catch (Throwable e) {
pendingRequestCount.decrementAndGet();
log.warn("[{}] Failed to process the request: {}", requestId, request, e);
stats.incrementFailed();
}
}
});
consumer.commit();
}
public void subscribe(Set<TopicPartitionInfo> partitions) {
requestConsumer.update(partitions);
}
public void stop() {
if (requestConsumer != null) {
requestConsumer.stop();
requestConsumer.awaitStop();
}
if (responseProducer != null) {
responseProducer.stop();
}
if (scheduler != null) {
scheduler.shutdownNow();
}
}
}

5
common/queue/src/main/java/org/thingsboard/server/queue/common/consumer/MainQueueConsumerManager.java

@ -25,7 +25,6 @@ import org.thingsboard.server.queue.TbQueueConsumer;
import org.thingsboard.server.queue.TbQueueMsg;
import org.thingsboard.server.queue.common.consumer.TbQueueConsumerManagerTask.UpdateConfigTask;
import org.thingsboard.server.queue.common.consumer.TbQueueConsumerManagerTask.UpdatePartitionsTask;
import org.thingsboard.server.queue.discovery.QueueKey;
import org.thingsboard.server.queue.kafka.TbKafkaConsumerTemplate;
import java.util.Collection;
@ -50,7 +49,7 @@ import java.util.function.Function;
public class MainQueueConsumerManager<M extends TbQueueMsg, C extends QueueConfig> {
@Getter
protected final QueueKey queueKey;
protected final Object queueKey;
@Getter
protected C config;
protected final MsgPackProcessor<M, C> msgPackProcessor;
@ -72,7 +71,7 @@ public class MainQueueConsumerManager<M extends TbQueueMsg, C extends QueueConfi
protected volatile boolean stopped;
@Builder
public MainQueueConsumerManager(QueueKey queueKey, C config,
public MainQueueConsumerManager(Object queueKey, C config,
MsgPackProcessor<M, C> msgPackProcessor,
BiFunction<C, TopicPartitionInfo, TbQueueConsumer<M>> consumerCreator,
ExecutorService consumerExecutor,

3
common/queue/src/main/java/org/thingsboard/server/queue/common/consumer/PartitionedQueueConsumerManager.java

@ -26,7 +26,6 @@ import org.thingsboard.server.queue.TbQueueMsg;
import org.thingsboard.server.queue.common.consumer.TbQueueConsumerManagerTask.AddPartitionsTask;
import org.thingsboard.server.queue.common.consumer.TbQueueConsumerManagerTask.DeletePartitionsTask;
import org.thingsboard.server.queue.common.consumer.TbQueueConsumerManagerTask.RemovePartitionsTask;
import org.thingsboard.server.queue.discovery.QueueKey;
import java.util.Set;
import java.util.concurrent.ExecutorService;
@ -44,7 +43,7 @@ public class PartitionedQueueConsumerManager<M extends TbQueueMsg> extends MainQ
private final String topic;
@Builder(builderMethodName = "create") // not to conflict with super.builder()
public PartitionedQueueConsumerManager(QueueKey queueKey, String topic, long pollInterval, MsgPackProcessor<M, QueueConfig> msgPackProcessor,
public PartitionedQueueConsumerManager(Object queueKey, String topic, long pollInterval, MsgPackProcessor<M, QueueConfig> msgPackProcessor,
BiFunction<QueueConfig, TopicPartitionInfo, TbQueueConsumer<M>> consumerCreator, TbQueueAdmin queueAdmin,
ExecutorService consumerExecutor, ScheduledExecutorService scheduler,
ExecutorService taskExecutor, Consumer<Throwable> uncaughtErrorHandler) {

4
common/queue/src/main/java/org/thingsboard/server/queue/common/state/DefaultQueueStateService.java

@ -18,10 +18,12 @@ package org.thingsboard.server.queue.common.state;
import org.thingsboard.server.queue.TbQueueMsg;
import org.thingsboard.server.queue.common.consumer.PartitionedQueueConsumerManager;
import java.util.Collections;
public class DefaultQueueStateService<E extends TbQueueMsg, S extends TbQueueMsg> extends QueueStateService<E, S> {
public DefaultQueueStateService(PartitionedQueueConsumerManager<E> eventConsumer) {
super(eventConsumer);
super(eventConsumer, Collections.emptyList());
}
}

8
common/queue/src/main/java/org/thingsboard/server/queue/common/state/KafkaQueueStateService.java

@ -22,6 +22,8 @@ import org.thingsboard.server.queue.TbQueueMsg;
import org.thingsboard.server.queue.common.consumer.PartitionedQueueConsumerManager;
import org.thingsboard.server.queue.discovery.QueueKey;
import java.util.Collections;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.function.Supplier;
@ -37,8 +39,9 @@ public class KafkaQueueStateService<E extends TbQueueMsg, S extends TbQueueMsg>
@Builder
public KafkaQueueStateService(PartitionedQueueConsumerManager<E> eventConsumer,
PartitionedQueueConsumerManager<S> stateConsumer,
List<PartitionedQueueConsumerManager<?>> otherConsumers,
Supplier<Map<String, Long>> eventsStartOffsetsProvider) {
super(eventConsumer);
super(eventConsumer, otherConsumers != null ? otherConsumers : Collections.emptyList());
this.stateConsumer = stateConsumer;
this.eventsStartOffsetsProvider = eventsStartOffsetsProvider;
}
@ -62,6 +65,9 @@ public class KafkaQueueStateService<E extends TbQueueMsg, S extends TbQueueMsg>
TopicPartitionInfo eventPartition = statePartition.withTopic(eventConsumer.getTopic());
if (this.partitions.get(queueKey).contains(eventPartition)) {
eventConsumer.addPartitions(Set.of(eventPartition), null, eventsStartOffsets != null ? eventsStartOffsets::get : null);
for (PartitionedQueueConsumerManager<?> consumer : otherConsumers) {
consumer.addPartitions(Set.of(statePartition.withTopic(consumer.getTopic())));
}
}
} finally {
readLock.unlock();

14
common/queue/src/main/java/org/thingsboard/server/queue/common/state/QueueStateService.java

@ -25,6 +25,7 @@ import org.thingsboard.server.queue.discovery.QueueKey;
import java.util.Collections;
import java.util.HashMap;
import java.util.HashSet;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
@ -37,6 +38,7 @@ import static org.thingsboard.server.common.msg.queue.TopicPartitionInfo.withTop
public abstract class QueueStateService<E extends TbQueueMsg, S extends TbQueueMsg> {
protected final PartitionedQueueConsumerManager<E> eventConsumer;
protected final List<PartitionedQueueConsumerManager<?>> otherConsumers;
@Getter
protected final Map<QueueKey, Set<TopicPartitionInfo>> partitions = new HashMap<>();
@ -45,8 +47,9 @@ public abstract class QueueStateService<E extends TbQueueMsg, S extends TbQueueM
protected final ReadWriteLock partitionsLock = new ReentrantReadWriteLock();
protected QueueStateService(PartitionedQueueConsumerManager<E> eventConsumer) {
protected QueueStateService(PartitionedQueueConsumerManager<E> eventConsumer, List<PartitionedQueueConsumerManager<?>> otherConsumers) {
this.eventConsumer = eventConsumer;
this.otherConsumers = otherConsumers;
}
public void update(QueueKey queueKey, Set<TopicPartitionInfo> newPartitions) {
@ -78,10 +81,16 @@ public abstract class QueueStateService<E extends TbQueueMsg, S extends TbQueueM
protected void addPartitions(QueueKey queueKey, Set<TopicPartitionInfo> partitions) {
eventConsumer.addPartitions(partitions);
for (PartitionedQueueConsumerManager<?> consumer : otherConsumers) {
consumer.addPartitions(withTopic(partitions, consumer.getTopic()));
}
}
protected void removePartitions(QueueKey queueKey, Set<TopicPartitionInfo> partitions) {
eventConsumer.removePartitions(partitions);
for (PartitionedQueueConsumerManager<?> consumer : otherConsumers) {
consumer.removePartitions(withTopic(partitions, consumer.getTopic()));
}
}
public void delete(Set<TopicPartitionInfo> partitions) {
@ -100,6 +109,9 @@ public abstract class QueueStateService<E extends TbQueueMsg, S extends TbQueueM
protected void deletePartitions(Set<TopicPartitionInfo> partitions) {
eventConsumer.delete(withTopic(partitions, eventConsumer.getTopic()));
for (PartitionedQueueConsumerManager<?> consumer : otherConsumers) {
consumer.removePartitions(withTopic(partitions, consumer.getTopic()));
}
}
public Set<TopicPartitionInfo> getPartitionsInProgress() {

2
common/queue/src/main/java/org/thingsboard/server/queue/edqs/EdqsConfig.java

@ -44,6 +44,8 @@ public class EdqsConfig {
private int maxPendingRequests;
@Value("${queue.edqs.max_request_timeout:20000}")
private int maxRequestTimeout;
@Value("${queue.edqs.request_executor_size:50}")
private int requestExecutorSize;
public String getLabel() {
if (partitioningStrategy == EdqsPartitioningStrategy.NONE) {

70
common/queue/src/main/java/org/thingsboard/server/queue/edqs/EdqsExecutors.java

@ -0,0 +1,70 @@
/**
* Copyright © 2016-2025 The Thingsboard Authors
*
* 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
*
* http://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 org.thingsboard.server.queue.edqs;
import com.google.common.util.concurrent.ListeningExecutorService;
import com.google.common.util.concurrent.MoreExecutors;
import jakarta.annotation.PostConstruct;
import jakarta.annotation.PreDestroy;
import lombok.Getter;
import lombok.RequiredArgsConstructor;
import org.springframework.context.annotation.Lazy;
import org.springframework.stereotype.Component;
import org.thingsboard.common.util.ThingsBoardExecutors;
import org.thingsboard.common.util.ThingsBoardThreadFactory;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
@Lazy
@Component
@Getter
@RequiredArgsConstructor
public class EdqsExecutors {
private final EdqsConfig edqsConfig;
private ExecutorService consumersExecutor;
private ExecutorService consumerTaskExecutor;
private ScheduledExecutorService scheduler;
private ListeningExecutorService requestExecutor;
@PostConstruct
private void init() {
consumersExecutor = Executors.newCachedThreadPool(ThingsBoardThreadFactory.forName("edqs-consumer"));
consumerTaskExecutor = ThingsBoardExecutors.newWorkStealingPool(4, "edqs-consumer-task-executor");
scheduler = ThingsBoardExecutors.newSingleThreadScheduledExecutor("edqs-scheduler");
requestExecutor = MoreExecutors.listeningDecorator(ThingsBoardExecutors.newWorkStealingPool(edqsConfig.getRequestExecutorSize(), "edqs-requests"));
}
@PreDestroy
private void destroy() {
if (consumersExecutor != null) {
consumersExecutor.shutdownNow();
}
if (consumerTaskExecutor != null) {
consumerTaskExecutor.shutdownNow();
}
if (scheduler != null) {
scheduler.shutdownNow();
}
if (requestExecutor != null) {
requestExecutor.shutdownNow();
}
}
}

5
common/queue/src/main/java/org/thingsboard/server/queue/edqs/EdqsQueueFactory.java

@ -19,8 +19,9 @@ import org.thingsboard.server.gen.transport.TransportProtos.FromEdqsMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToEdqsMsg;
import org.thingsboard.server.queue.TbQueueAdmin;
import org.thingsboard.server.queue.TbQueueConsumer;
import org.thingsboard.server.queue.TbQueueHandler;
import org.thingsboard.server.queue.TbQueueProducer;
import org.thingsboard.server.queue.TbQueueResponseTemplate;
import org.thingsboard.server.queue.common.PartitionedQueueResponseTemplate;
import org.thingsboard.server.queue.common.TbProtoQueueMsg;
public interface EdqsQueueFactory {
@ -33,7 +34,7 @@ public interface EdqsQueueFactory {
TbQueueProducer<TbProtoQueueMsg<ToEdqsMsg>> createEdqsStateProducer();
TbQueueResponseTemplate<TbProtoQueueMsg<ToEdqsMsg>, TbProtoQueueMsg<FromEdqsMsg>> createEdqsResponseTemplate();
PartitionedQueueResponseTemplate<TbProtoQueueMsg<ToEdqsMsg>, TbProtoQueueMsg<FromEdqsMsg>> createEdqsResponseTemplate(TbQueueHandler<TbProtoQueueMsg<ToEdqsMsg>, TbProtoQueueMsg<FromEdqsMsg>> handler);
TbQueueAdmin getEdqsQueueAdmin();

26
common/queue/src/main/java/org/thingsboard/server/queue/edqs/InMemoryEdqsQueueFactory.java

@ -17,16 +17,15 @@ package org.thingsboard.server.queue.edqs;
import lombok.RequiredArgsConstructor;
import org.springframework.stereotype.Component;
import org.thingsboard.common.util.ThingsBoardExecutors;
import org.thingsboard.server.common.stats.StatsFactory;
import org.thingsboard.server.common.stats.StatsType;
import org.thingsboard.server.gen.transport.TransportProtos.FromEdqsMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToEdqsMsg;
import org.thingsboard.server.queue.TbQueueAdmin;
import org.thingsboard.server.queue.TbQueueConsumer;
import org.thingsboard.server.queue.TbQueueHandler;
import org.thingsboard.server.queue.TbQueueProducer;
import org.thingsboard.server.queue.TbQueueResponseTemplate;
import org.thingsboard.server.queue.common.DefaultTbQueueResponseTemplate;
import org.thingsboard.server.queue.common.PartitionedQueueResponseTemplate;
import org.thingsboard.server.queue.common.TbProtoQueueMsg;
import org.thingsboard.server.queue.memory.InMemoryStorage;
import org.thingsboard.server.queue.memory.InMemoryTbQueueConsumer;
@ -39,6 +38,7 @@ public class InMemoryEdqsQueueFactory implements EdqsQueueFactory {
private final InMemoryStorage storage;
private final EdqsConfig edqsConfig;
private final EdqsExecutors edqsExecutors;
private final StatsFactory statsFactory;
private final TbQueueAdmin queueAdmin;
@ -63,17 +63,21 @@ public class InMemoryEdqsQueueFactory implements EdqsQueueFactory {
}
@Override
public TbQueueResponseTemplate<TbProtoQueueMsg<ToEdqsMsg>, TbProtoQueueMsg<FromEdqsMsg>> createEdqsResponseTemplate() {
TbQueueConsumer<TbProtoQueueMsg<ToEdqsMsg>> requestConsumer = new InMemoryTbQueueConsumer<>(storage, edqsConfig.getRequestsTopic());
public PartitionedQueueResponseTemplate<TbProtoQueueMsg<ToEdqsMsg>, TbProtoQueueMsg<FromEdqsMsg>> createEdqsResponseTemplate(TbQueueHandler<TbProtoQueueMsg<ToEdqsMsg>, TbProtoQueueMsg<FromEdqsMsg>> handler) {
TbQueueProducer<TbProtoQueueMsg<FromEdqsMsg>> responseProducer = new InMemoryTbQueueProducer<>(storage, edqsConfig.getResponsesTopic());
return DefaultTbQueueResponseTemplate.<TbProtoQueueMsg<ToEdqsMsg>, TbProtoQueueMsg<FromEdqsMsg>>builder()
.requestTemplate(requestConsumer)
.responseTemplate(responseProducer)
.maxPendingRequests(edqsConfig.getMaxPendingRequests())
.requestTimeout(edqsConfig.getMaxRequestTimeout())
return PartitionedQueueResponseTemplate.<TbProtoQueueMsg<ToEdqsMsg>, TbProtoQueueMsg<FromEdqsMsg>>builder()
.key("edqs")
.handler(handler)
.requestsTopic(edqsConfig.getRequestsTopic())
.consumerCreator(tpi -> new InMemoryTbQueueConsumer<>(storage, edqsConfig.getRequestsTopic()))
.responseProducer(responseProducer)
.pollInterval(edqsConfig.getPollInterval())
.requestTimeout(edqsConfig.getMaxRequestTimeout())
.maxPendingRequests(edqsConfig.getMaxPendingRequests())
.consumerExecutor(edqsExecutors.getConsumersExecutor())
.callbackExecutor(edqsExecutors.getRequestExecutor())
.consumerTaskExecutor(edqsExecutors.getConsumerTaskExecutor())
.stats(statsFactory.createMessagesStats(StatsType.EDQS.getName()))
.executor(ThingsBoardExecutors.newWorkStealingPool(5, "edqs"))
.build();
}

35
common/queue/src/main/java/org/thingsboard/server/queue/edqs/KafkaEdqsQueueFactory.java

@ -16,14 +16,13 @@
package org.thingsboard.server.queue.edqs;
import org.springframework.stereotype.Component;
import org.thingsboard.common.util.ThingsBoardExecutors;
import org.thingsboard.server.common.stats.StatsFactory;
import org.thingsboard.server.common.stats.StatsType;
import org.thingsboard.server.gen.transport.TransportProtos.FromEdqsMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToEdqsMsg;
import org.thingsboard.server.queue.TbQueueHandler;
import org.thingsboard.server.queue.TbQueueProducer;
import org.thingsboard.server.queue.TbQueueResponseTemplate;
import org.thingsboard.server.queue.common.DefaultTbQueueResponseTemplate;
import org.thingsboard.server.queue.common.PartitionedQueueResponseTemplate;
import org.thingsboard.server.queue.common.TbProtoQueueMsg;
import org.thingsboard.server.queue.discovery.TbServiceInfoProvider;
import org.thingsboard.server.queue.discovery.TopicService;
@ -45,6 +44,7 @@ public class KafkaEdqsQueueFactory implements EdqsQueueFactory {
private final TbKafkaAdmin edqsRequestsAdmin;
private final TbKafkaAdmin edqsStateAdmin;
private final EdqsConfig edqsConfig;
private final EdqsExecutors edqsExecutors;
private final TbServiceInfoProvider serviceInfoProvider;
private final TbKafkaConsumerStatsService consumerStatsService;
private final TopicService topicService;
@ -53,7 +53,7 @@ public class KafkaEdqsQueueFactory implements EdqsQueueFactory {
private final AtomicInteger consumerCounter = new AtomicInteger();
public KafkaEdqsQueueFactory(TbKafkaSettings kafkaSettings, TbKafkaTopicConfigs topicConfigs,
EdqsConfig edqsConfig, TbServiceInfoProvider serviceInfoProvider,
EdqsConfig edqsConfig, EdqsExecutors edqsExecutors, TbServiceInfoProvider serviceInfoProvider,
TbKafkaConsumerStatsService consumerStatsService, TopicService topicService,
StatsFactory statsFactory) {
this.edqsEventsAdmin = new TbKafkaAdmin(kafkaSettings, topicConfigs.getEdqsEventsConfigs());
@ -61,6 +61,7 @@ public class KafkaEdqsQueueFactory implements EdqsQueueFactory {
this.edqsStateAdmin = new TbKafkaAdmin(kafkaSettings, topicConfigs.getEdqsStateConfigs());
this.kafkaSettings = kafkaSettings;
this.edqsConfig = edqsConfig;
this.edqsExecutors = edqsExecutors;
this.serviceInfoProvider = serviceInfoProvider;
this.consumerStatsService = consumerStatsService;
this.topicService = topicService;
@ -116,25 +117,29 @@ public class KafkaEdqsQueueFactory implements EdqsQueueFactory {
}
@Override
public TbQueueResponseTemplate<TbProtoQueueMsg<ToEdqsMsg>, TbProtoQueueMsg<FromEdqsMsg>> createEdqsResponseTemplate() {
var requestConsumer = createEdqsMsgConsumer(edqsConfig.getRequestsTopic(),
"edqs-requests-consumer-" + serviceInfoProvider.getServiceId(),
"edqs-requests-consumer-group",
false, edqsRequestsAdmin);
public PartitionedQueueResponseTemplate<TbProtoQueueMsg<ToEdqsMsg>, TbProtoQueueMsg<FromEdqsMsg>> createEdqsResponseTemplate(TbQueueHandler<TbProtoQueueMsg<ToEdqsMsg>, TbProtoQueueMsg<FromEdqsMsg>> handler) {
var responseProducer = TbKafkaProducerTemplate.<TbProtoQueueMsg<FromEdqsMsg>>builder()
.settings(kafkaSettings)
.clientId("edqs-response-producer-" + serviceInfoProvider.getServiceId())
.defaultTopic(topicService.buildTopicName(edqsConfig.getResponsesTopic()))
.admin(edqsRequestsAdmin)
.build();
return DefaultTbQueueResponseTemplate.<TbProtoQueueMsg<ToEdqsMsg>, TbProtoQueueMsg<FromEdqsMsg>>builder()
.requestTemplate(requestConsumer)
.responseTemplate(responseProducer)
.maxPendingRequests(edqsConfig.getMaxPendingRequests())
.requestTimeout(edqsConfig.getMaxRequestTimeout())
return PartitionedQueueResponseTemplate.<TbProtoQueueMsg<ToEdqsMsg>, TbProtoQueueMsg<FromEdqsMsg>>builder()
.key("edqs")
.handler(handler)
.requestsTopic(topicService.buildTopicName(edqsConfig.getRequestsTopic()))
.consumerCreator(tpi -> createEdqsMsgConsumer(edqsConfig.getRequestsTopic(),
"edqs-requests-consumer-" + serviceInfoProvider.getServiceId() + "-" + tpi.getPartition().orElse(999),
"edqs-requests-consumer-group",
false, edqsRequestsAdmin))
.responseProducer(responseProducer)
.pollInterval(edqsConfig.getPollInterval())
.requestTimeout(edqsConfig.getMaxRequestTimeout())
.maxPendingRequests(edqsConfig.getMaxPendingRequests())
.consumerExecutor(edqsExecutors.getConsumersExecutor())
.callbackExecutor(edqsExecutors.getRequestExecutor())
.consumerTaskExecutor(edqsExecutors.getConsumerTaskExecutor())
.stats(statsFactory.createMessagesStats(StatsType.EDQS.getName()))
.executor(ThingsBoardExecutors.newWorkStealingPool(5, "edqs"))
.build();
}

4
common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaProducerTemplate.java

@ -75,10 +75,6 @@ public class TbKafkaProducerTemplate<T extends TbQueueMsg> implements TbQueuePro
topics = ConcurrentHashMap.newKeySet();
}
@Override
public void init() {
}
void addAnalyticHeaders(List<Header> headers) {
headers.add(new RecordHeader("_producerId", getClientId().getBytes(StandardCharsets.UTF_8)));
headers.add(new RecordHeader("_threadName", Thread.currentThread().getName().getBytes(StandardCharsets.UTF_8)));

5
common/queue/src/main/java/org/thingsboard/server/queue/memory/InMemoryTbQueueProducer.java

@ -33,11 +33,6 @@ public class InMemoryTbQueueProducer<T extends TbQueueMsg> implements TbQueuePro
this.defaultTopic = defaultTopic;
}
@Override
public void init() {
}
@Override
public void send(TopicPartitionInfo tpi, T msg, TbQueueCallback callback) {
boolean result = storage.put(tpi.getFullTopicName(), msg);

19
common/queue/src/test/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplateTest.java

@ -37,11 +37,23 @@ import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicLong;
import static org.hamcrest.MatcherAssert.assertThat;
import static org.hamcrest.Matchers.*;
import static org.hamcrest.Matchers.equalTo;
import static org.hamcrest.Matchers.greaterThanOrEqualTo;
import static org.hamcrest.Matchers.is;
import static org.hamcrest.Matchers.lessThan;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.anyLong;
import static org.mockito.BDDMockito.*;
import static org.mockito.Mockito.*;
import static org.mockito.BDDMockito.RETURNS_DEEP_STUBS;
import static org.mockito.BDDMockito.atLeastOnce;
import static org.mockito.BDDMockito.lenient;
import static org.mockito.BDDMockito.mock;
import static org.mockito.BDDMockito.never;
import static org.mockito.BDDMockito.spy;
import static org.mockito.BDDMockito.times;
import static org.mockito.BDDMockito.verify;
import static org.mockito.BDDMockito.willAnswer;
import static org.mockito.BDDMockito.willDoNothing;
import static org.mockito.BDDMockito.willReturn;
import static org.mockito.hamcrest.MockitoHamcrest.longThat;
@Slf4j
@ -96,7 +108,6 @@ public class DefaultTbQueueRequestTemplateTest {
inst.init();
assertThat(inst.nextCleanupNs, equalTo(0L));
verify(queueAdmin, times(1)).createTopicIfNotExists(topic);
verify(requestTemplate, times(1)).init();
verify(responseTemplate, times(1)).subscribe();
verify(executorMock, times(1)).submit(any(Runnable.class));

2
edqs/src/main/resources/edqs.yml

@ -71,6 +71,8 @@ queue:
max_pending_requests: "${TB_EDQS_MAX_PENDING_REQUESTS:10000}"
# Maximum timeout for requests to EDQS
max_request_timeout: "${TB_EDQS_MAX_REQUEST_TIMEOUT:20000}"
# Thread pool size for EDQS requests executor
request_executor_size: "${TB_EDQS_REQUEST_EXECUTOR_SIZE:50}"
# Strings longer than this threshold will be compressed
string_compression_length_threshold: "${TB_EDQS_STRING_COMPRESSION_LENGTH_THRESHOLD:512}"
stats:

Loading…
Cancel
Save