diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index abafd81e7c..6ad3386c4f 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/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: diff --git a/common/cluster-api/src/main/java/org/thingsboard/server/queue/TbQueueProducer.java b/common/cluster-api/src/main/java/org/thingsboard/server/queue/TbQueueProducer.java index 51fd735a9d..a374928b16 100644 --- a/common/cluster-api/src/main/java/org/thingsboard/server/queue/TbQueueProducer.java +++ b/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 { - void init(); - String getDefaultTopic(); void send(TopicPartitionInfo tpi, T msg, TbQueueCallback callback); void stop(); + } diff --git a/common/edqs/src/main/java/org/thingsboard/server/edqs/processor/EdqsProcessor.java b/common/edqs/src/main/java/org/thingsboard/server/edqs/processor/EdqsProcessor.java index 78d17ef368..7b30fbaaa5 100644 --- a/common/edqs/src/main/java/org/thingsboard/server/edqs/processor/EdqsProcessor.java +++ b/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, 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> eventConsumer; - private TbQueueResponseTemplate, TbProtoQueueMsg> responseTemplate; - - private ExecutorService consumersExecutor; - private ExecutorService taskExecutor; - private ScheduledExecutorService scheduler; + private PartitionedQueueResponseTemplate, TbProtoQueueMsg> 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, @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, .execute(applicationContext::close); } }; + requestExecutor = edqsExecutors.getRequestExecutor(); eventConsumer = PartitionedQueueConsumerManager.>create() .queueKey(new QueueKey(ServiceType.EDQS, config.getEventsTopic())) @@ -141,19 +131,14 @@ public class EdqsProcessor implements TbQueueHandler, }) .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, } try { Set 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 oldPartitions = event.getOldPartitions().get(new QueueKey(ServiceType.EDQS)); if (CollectionsUtil.isNotEmpty(oldPartitions)) { @@ -290,11 +273,6 @@ public class EdqsProcessor implements TbQueueHandler, eventConsumer.awaitStop(); responseTemplate.stop(); stateService.stop(); - - consumersExecutor.shutdownNow(); - taskExecutor.shutdownNow(); - scheduler.shutdownNow(); - requestExecutor.shutdownNow(); } } diff --git a/common/edqs/src/main/java/org/thingsboard/server/edqs/state/EdqsStateService.java b/common/edqs/src/main/java/org/thingsboard/server/edqs/state/EdqsStateService.java index ee7b058a8a..050fa74e6b 100644 --- a/common/edqs/src/main/java/org/thingsboard/server/edqs/state/EdqsStateService.java +++ b/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> eventConsumer); + void init(PartitionedQueueConsumerManager> eventConsumer, List> otherConsumers); void process(Set partitions); diff --git a/common/edqs/src/main/java/org/thingsboard/server/edqs/state/KafkaEdqsStateService.java b/common/edqs/src/main/java/org/thingsboard/server/edqs/state/KafkaEdqsStateService.java index 0efe6e7d3b..e8ede39d2c 100644 --- a/common/edqs/src/main/java/org/thingsboard/server/edqs/state/KafkaEdqsStateService.java +++ b/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> eventConsumer) { + public void init(PartitionedQueueConsumerManager> eventConsumer, List> otherConsumers) { TbKafkaAdmin queueAdmin = queueFactory.getEdqsQueueAdmin(); stateConsumer = PartitionedQueueConsumerManager.>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>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 diff --git a/common/edqs/src/main/java/org/thingsboard/server/edqs/state/LocalEdqsStateService.java b/common/edqs/src/main/java/org/thingsboard/server/edqs/state/LocalEdqsStateService.java index cde21edfaf..c9923e6194 100644 --- a/common/edqs/src/main/java/org/thingsboard/server/edqs/state/LocalEdqsStateService.java +++ b/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> eventConsumer; + private List> otherConsumers; private Set partitions; @Override - public void init(PartitionedQueueConsumerManager> eventConsumer) { + public void init(PartitionedQueueConsumerManager> eventConsumer, List> 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; } diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java b/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java index 1beb505595..c6f76a5000 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java @@ -91,7 +91,6 @@ public class DefaultTbQueueRequestTemplate requestTemplate; private final TbQueueProducer responseTemplate; - private final ConcurrentMap pendingRequests; private final ExecutorService loopExecutor; private final ScheduledExecutorService timeoutExecutor; private final ExecutorService callbackExecutor; @@ -67,7 +64,6 @@ public class DefaultTbQueueResponseTemplate(); this.maxPendingRequests = maxPendingRequests; this.pollInterval = pollInterval; this.requestTimeout = requestTimeout; @@ -89,7 +85,6 @@ public class DefaultTbQueueResponseTemplate handler) { - this.responseTemplate.init(); loopExecutor.submit(() -> { while (!stopped) { try { diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/common/PartitionedQueueResponseTemplate.java b/common/queue/src/main/java/org/thingsboard/server/queue/common/PartitionedQueueResponseTemplate.java new file mode 100644 index 0000000000..d1e1be708f --- /dev/null +++ b/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 extends AbstractTbQueueTemplate { + + @Getter + private final PartitionedQueueConsumerManager requestConsumer; + private final TbQueueProducer responseProducer; + + private final TbQueueHandler 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 handler, + String requestsTopic, + Function> consumerCreator, + TbQueueProducer 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.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 requests, TbQueueConsumer 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 partitions) { + requestConsumer.update(partitions); + } + + public void stop() { + if (requestConsumer != null) { + requestConsumer.stop(); + requestConsumer.awaitStop(); + } + if (responseProducer != null) { + responseProducer.stop(); + } + if (scheduler != null) { + scheduler.shutdownNow(); + } + } + +} diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/common/consumer/MainQueueConsumerManager.java b/common/queue/src/main/java/org/thingsboard/server/queue/common/consumer/MainQueueConsumerManager.java index 2233855a37..86e0721dd8 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/common/consumer/MainQueueConsumerManager.java +++ b/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 { @Getter - protected final QueueKey queueKey; + protected final Object queueKey; @Getter protected C config; protected final MsgPackProcessor msgPackProcessor; @@ -72,7 +71,7 @@ public class MainQueueConsumerManager msgPackProcessor, BiFunction> consumerCreator, ExecutorService consumerExecutor, diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/common/consumer/PartitionedQueueConsumerManager.java b/common/queue/src/main/java/org/thingsboard/server/queue/common/consumer/PartitionedQueueConsumerManager.java index 1b19fcab0e..42d5b6bb18 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/common/consumer/PartitionedQueueConsumerManager.java +++ b/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 extends MainQ private final String topic; @Builder(builderMethodName = "create") // not to conflict with super.builder() - public PartitionedQueueConsumerManager(QueueKey queueKey, String topic, long pollInterval, MsgPackProcessor msgPackProcessor, + public PartitionedQueueConsumerManager(Object queueKey, String topic, long pollInterval, MsgPackProcessor msgPackProcessor, BiFunction> consumerCreator, TbQueueAdmin queueAdmin, ExecutorService consumerExecutor, ScheduledExecutorService scheduler, ExecutorService taskExecutor, Consumer uncaughtErrorHandler) { diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/common/state/DefaultQueueStateService.java b/common/queue/src/main/java/org/thingsboard/server/queue/common/state/DefaultQueueStateService.java index be019caaa7..be379fb76d 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/common/state/DefaultQueueStateService.java +++ b/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 extends QueueStateService { public DefaultQueueStateService(PartitionedQueueConsumerManager eventConsumer) { - super(eventConsumer); + super(eventConsumer, Collections.emptyList()); } } diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/common/state/KafkaQueueStateService.java b/common/queue/src/main/java/org/thingsboard/server/queue/common/state/KafkaQueueStateService.java index bf02afe86c..d8d0c8e0d2 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/common/state/KafkaQueueStateService.java +++ b/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 @Builder public KafkaQueueStateService(PartitionedQueueConsumerManager eventConsumer, PartitionedQueueConsumerManager stateConsumer, + List> otherConsumers, Supplier> eventsStartOffsetsProvider) { - super(eventConsumer); + super(eventConsumer, otherConsumers != null ? otherConsumers : Collections.emptyList()); this.stateConsumer = stateConsumer; this.eventsStartOffsetsProvider = eventsStartOffsetsProvider; } @@ -62,6 +65,9 @@ public class KafkaQueueStateService 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(); diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/common/state/QueueStateService.java b/common/queue/src/main/java/org/thingsboard/server/queue/common/state/QueueStateService.java index 29426fab63..61b81707ce 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/common/state/QueueStateService.java +++ b/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 { protected final PartitionedQueueConsumerManager eventConsumer; + protected final List> otherConsumers; @Getter protected final Map> partitions = new HashMap<>(); @@ -45,8 +47,9 @@ public abstract class QueueStateService eventConsumer) { + protected QueueStateService(PartitionedQueueConsumerManager eventConsumer, List> otherConsumers) { this.eventConsumer = eventConsumer; + this.otherConsumers = otherConsumers; } public void update(QueueKey queueKey, Set newPartitions) { @@ -78,10 +81,16 @@ public abstract class QueueStateService partitions) { eventConsumer.addPartitions(partitions); + for (PartitionedQueueConsumerManager consumer : otherConsumers) { + consumer.addPartitions(withTopic(partitions, consumer.getTopic())); + } } protected void removePartitions(QueueKey queueKey, Set partitions) { eventConsumer.removePartitions(partitions); + for (PartitionedQueueConsumerManager consumer : otherConsumers) { + consumer.removePartitions(withTopic(partitions, consumer.getTopic())); + } } public void delete(Set partitions) { @@ -100,6 +109,9 @@ public abstract class QueueStateService partitions) { eventConsumer.delete(withTopic(partitions, eventConsumer.getTopic())); + for (PartitionedQueueConsumerManager consumer : otherConsumers) { + consumer.removePartitions(withTopic(partitions, consumer.getTopic())); + } } public Set getPartitionsInProgress() { diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/edqs/EdqsConfig.java b/common/queue/src/main/java/org/thingsboard/server/queue/edqs/EdqsConfig.java index 3c927f135b..4a9b9f2a28 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/edqs/EdqsConfig.java +++ b/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) { diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/edqs/EdqsExecutors.java b/common/queue/src/main/java/org/thingsboard/server/queue/edqs/EdqsExecutors.java new file mode 100644 index 0000000000..8e804410fe --- /dev/null +++ b/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(); + } + } + +} diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/edqs/EdqsQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/edqs/EdqsQueueFactory.java index 5c0d68779a..ed021e3380 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/edqs/EdqsQueueFactory.java +++ b/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> createEdqsStateProducer(); - TbQueueResponseTemplate, TbProtoQueueMsg> createEdqsResponseTemplate(); + PartitionedQueueResponseTemplate, TbProtoQueueMsg> createEdqsResponseTemplate(TbQueueHandler, TbProtoQueueMsg> handler); TbQueueAdmin getEdqsQueueAdmin(); diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/edqs/InMemoryEdqsQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/edqs/InMemoryEdqsQueueFactory.java index 8c670e66c0..6c7618f7db 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/edqs/InMemoryEdqsQueueFactory.java +++ b/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> createEdqsResponseTemplate() { - TbQueueConsumer> requestConsumer = new InMemoryTbQueueConsumer<>(storage, edqsConfig.getRequestsTopic()); + public PartitionedQueueResponseTemplate, TbProtoQueueMsg> createEdqsResponseTemplate(TbQueueHandler, TbProtoQueueMsg> handler) { TbQueueProducer> responseProducer = new InMemoryTbQueueProducer<>(storage, edqsConfig.getResponsesTopic()); - return DefaultTbQueueResponseTemplate., TbProtoQueueMsg>builder() - .requestTemplate(requestConsumer) - .responseTemplate(responseProducer) - .maxPendingRequests(edqsConfig.getMaxPendingRequests()) - .requestTimeout(edqsConfig.getMaxRequestTimeout()) + return PartitionedQueueResponseTemplate., TbProtoQueueMsg>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(); } diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/edqs/KafkaEdqsQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/edqs/KafkaEdqsQueueFactory.java index ab88943b10..f071c30942 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/edqs/KafkaEdqsQueueFactory.java +++ b/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> createEdqsResponseTemplate() { - var requestConsumer = createEdqsMsgConsumer(edqsConfig.getRequestsTopic(), - "edqs-requests-consumer-" + serviceInfoProvider.getServiceId(), - "edqs-requests-consumer-group", - false, edqsRequestsAdmin); + public PartitionedQueueResponseTemplate, TbProtoQueueMsg> createEdqsResponseTemplate(TbQueueHandler, TbProtoQueueMsg> handler) { var responseProducer = TbKafkaProducerTemplate.>builder() .settings(kafkaSettings) .clientId("edqs-response-producer-" + serviceInfoProvider.getServiceId()) .defaultTopic(topicService.buildTopicName(edqsConfig.getResponsesTopic())) .admin(edqsRequestsAdmin) .build(); - return DefaultTbQueueResponseTemplate., TbProtoQueueMsg>builder() - .requestTemplate(requestConsumer) - .responseTemplate(responseProducer) - .maxPendingRequests(edqsConfig.getMaxPendingRequests()) - .requestTimeout(edqsConfig.getMaxRequestTimeout()) + return PartitionedQueueResponseTemplate., TbProtoQueueMsg>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(); } diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaProducerTemplate.java b/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaProducerTemplate.java index cac6f2ea1e..f0b3694d74 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaProducerTemplate.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaProducerTemplate.java @@ -75,10 +75,6 @@ public class TbKafkaProducerTemplate implements TbQueuePro topics = ConcurrentHashMap.newKeySet(); } - @Override - public void init() { - } - void addAnalyticHeaders(List
headers) { headers.add(new RecordHeader("_producerId", getClientId().getBytes(StandardCharsets.UTF_8))); headers.add(new RecordHeader("_threadName", Thread.currentThread().getName().getBytes(StandardCharsets.UTF_8))); diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/memory/InMemoryTbQueueProducer.java b/common/queue/src/main/java/org/thingsboard/server/queue/memory/InMemoryTbQueueProducer.java index 49c52ed21a..43d820bc7e 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/memory/InMemoryTbQueueProducer.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/memory/InMemoryTbQueueProducer.java @@ -33,11 +33,6 @@ public class InMemoryTbQueueProducer 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); diff --git a/common/queue/src/test/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplateTest.java b/common/queue/src/test/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplateTest.java index 9e493220ca..e5bb62bc04 100644 --- a/common/queue/src/test/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplateTest.java +++ b/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)); diff --git a/edqs/src/main/resources/edqs.yml b/edqs/src/main/resources/edqs.yml index 1cc32a4230..0815187704 100644 --- a/edqs/src/main/resources/edqs.yml +++ b/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: