Browse Source

moved pubsub queue executor to PubSubSettings

pull/9619/head
dashevchenko 3 years ago
parent
commit
fa12ba544d
  1. 6
      application/src/main/resources/thingsboard.yml
  2. 2
      common/queue/src/main/java/org/thingsboard/server/queue/pubsub/TbPubSubConsumerTemplate.java
  3. 2
      common/queue/src/main/java/org/thingsboard/server/queue/pubsub/TbPubSubProducerTemplate.java
  4. 60
      common/queue/src/main/java/org/thingsboard/server/queue/pubsub/TbPubSubQueueExecutorProvider.java
  5. 25
      common/queue/src/main/java/org/thingsboard/server/queue/pubsub/TbPubSubSettings.java
  6. 2
      msa/vc-executor/src/main/resources/tb-vc-executor.yml
  7. 2
      transport/coap/src/main/resources/tb-coap-transport.yml
  8. 2
      transport/http/src/main/resources/tb-http-transport.yml
  9. 2
      transport/lwm2m/src/main/resources/tb-lwm2m-transport.yml
  10. 2
      transport/mqtt/src/main/resources/tb-mqtt-transport.yml
  11. 2
      transport/snmp/src/main/resources/tb-snmp-transport.yml

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

@ -1411,8 +1411,8 @@ queue:
max_msg_size: "${TB_QUEUE_PUBSUB_MAX_MSG_SIZE:1048576}"
# Number of messages per consumer
max_messages: "${TB_QUEUE_PUBSUB_MAX_MESSAGES:1000}"
# Thread pool size for pubsub queue executor provider. If not set - default pubsub executor provider value will be used (5 * number of available processors)
executor_thread_pool_size: "${TB_QUEUE_PUBSUB_EXECUTOR_THREAD_POOL_SIZE:}"
# Thread pool size for pubsub queue executor provider. If set to 0 - default pubsub executor provider value will be used (5 * number of available processors)
executor_thread_pool_size: "${TB_QUEUE_PUBSUB_EXECUTOR_THREAD_POOL_SIZE:0}"
queue-properties:
# Pub/Sub properties for Rule Engine subscribers, messages which will commit after ackDeadlineInSec period can be consumed again
rule-engine: "${TB_QUEUE_PUBSUB_RE_QUEUE_PROPERTIES:ackDeadlineInSec:30;messageRetentionInSec:604800}"
@ -1643,7 +1643,7 @@ service:
assigned_tenant_profiles: "${TB_RULE_ENGINE_ASSIGNED_TENANT_PROFILES:}"
pubsub:
# Thread pool size for pubsub rule node executor provider. If not set - default pubsub executor provider value will be used (5 * number of available processors)
executor_thread_pool_size: "${TB_RULE_ENGINE_PUBSUB_EXECUTOR_THREAD_POOL_SIZE:}"
executor_thread_pool_size: "${TB_RULE_ENGINE_PUBSUB_EXECUTOR_THREAD_POOL_SIZE:0}"
# Metrics parameters
metrics:

2
common/queue/src/main/java/org/thingsboard/server/queue/pubsub/TbPubSubConsumerTemplate.java

@ -77,7 +77,7 @@ public class TbPubSubConsumerTemplate<T extends TbQueueMsg> extends AbstractPara
SubscriberStubSettings.defaultGrpcTransportProviderBuilder()
.setMaxInboundMessageSize(pubSubSettings.getMaxMsgSize())
.build())
.setExecutorProvider(FixedExecutorProvider.create(pubSubSettings.getTbPubSubQueueExecutorProvider().getExecutor()))
.setExecutorProvider(pubSubSettings.getExecutorProvider())
.build();
this.subscriber = GrpcSubscriberStub.create(subscriberStubSettings);
} catch (IOException e) {

2
common/queue/src/main/java/org/thingsboard/server/queue/pubsub/TbPubSubProducerTemplate.java

@ -123,7 +123,7 @@ public class TbPubSubProducerTemplate<T extends TbQueueMsg> implements TbQueuePr
ProjectTopicName topicName = ProjectTopicName.of(pubSubSettings.getProjectId(), topic);
Publisher publisher = Publisher.newBuilder(topicName)
.setCredentialsProvider(pubSubSettings.getCredentialsProvider())
.setExecutorProvider(FixedExecutorProvider.create(pubSubSettings.getTbPubSubQueueExecutorProvider().getExecutor()))
.setExecutorProvider(pubSubSettings.getExecutorProvider())
.build();
publisherMap.put(topic, publisher);
return publisher;

60
common/queue/src/main/java/org/thingsboard/server/queue/pubsub/TbPubSubQueueExecutorProvider.java

@ -1,60 +0,0 @@
/**
* Copyright © 2016-2023 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.pubsub;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression;
import org.springframework.stereotype.Component;
import org.thingsboard.common.util.ExecutorProvider;
import org.thingsboard.common.util.ThingsBoardThreadFactory;
import javax.annotation.PostConstruct;
import javax.annotation.PreDestroy;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
@ConditionalOnExpression("'${queue.type:null}'=='pubsub'")
@Component
public class TbPubSubQueueExecutorProvider implements ExecutorProvider {
@Value("${queue.pubsub.executor_thread_pool_size}")
private Integer threadPoolSize;
/**
* Refers to com.google.cloud.pubsub.v1.Publisher default executor configuration
*/
private static final int THREADS_PER_CPU = 5;
private ScheduledExecutorService executor;
@PostConstruct
public void init() {
if (threadPoolSize == null) {
threadPoolSize = THREADS_PER_CPU * Runtime.getRuntime().availableProcessors();
}
executor = Executors.newScheduledThreadPool(threadPoolSize, ThingsBoardThreadFactory.forName("pubsub-queue-executor"));
}
@Override
public ScheduledExecutorService getExecutor() {
return executor;
}
@PreDestroy
private void destroy() {
if (executor != null) {
executor.shutdownNow();
}
}
}

25
common/queue/src/main/java/org/thingsboard/server/queue/pubsub/TbPubSubSettings.java

@ -17,16 +17,20 @@ package org.thingsboard.server.queue.pubsub;
import com.google.api.gax.core.CredentialsProvider;
import com.google.api.gax.core.FixedCredentialsProvider;
import com.google.api.gax.core.FixedExecutorProvider;
import com.google.auth.oauth2.ServiceAccountCredentials;
import lombok.Data;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression;
import org.springframework.stereotype.Component;
import org.thingsboard.common.util.ThingsBoardThreadFactory;
import javax.annotation.PostConstruct;
import javax.annotation.PreDestroy;
import java.io.ByteArrayInputStream;
import java.io.IOException;
import java.util.concurrent.Executors;
@Slf4j
@ConditionalOnExpression("'${queue.type:null}'=='pubsub'")
@ -46,7 +50,15 @@ public class TbPubSubSettings {
@Value("${queue.pubsub.max_messages}")
private int maxMessages;
private final TbPubSubQueueExecutorProvider tbPubSubQueueExecutorProvider;
@Value("${queue.pubsub.executor_thread_pool_size:0}")
private int threadPoolSize;
/**
* Refers to com.google.cloud.pubsub.v1.Publisher default executor configuration
*/
private static final int THREADS_PER_CPU = 5;
private FixedExecutorProvider executorProvider;
private CredentialsProvider credentialsProvider;
@ -55,6 +67,17 @@ public class TbPubSubSettings {
ServiceAccountCredentials credentials = ServiceAccountCredentials.fromStream(
new ByteArrayInputStream(serviceAccount.getBytes()));
credentialsProvider = FixedCredentialsProvider.create(credentials);
if (threadPoolSize == 0) {
threadPoolSize = THREADS_PER_CPU * Runtime.getRuntime().availableProcessors();
}
executorProvider = FixedExecutorProvider
.create(Executors.newScheduledThreadPool(threadPoolSize, ThingsBoardThreadFactory.forName("pubsub-queue-executor")));
}
@PreDestroy
private void destroy() {
if (executorProvider != null) {
executorProvider.getExecutor().shutdownNow();
}
}
}

2
msa/vc-executor/src/main/resources/tb-vc-executor.yml

@ -173,7 +173,7 @@ queue:
# Number of messages per consumer
max_messages: "${TB_QUEUE_PUBSUB_MAX_MESSAGES:1000}"
# Thread pool size for pubsub queue executor provider. If not set - default pubsub executor provider value will be used (5 * number of available processors)
executor_thread_pool_size: "${TB_QUEUE_PUBSUB_EXECUTOR_THREAD_POOL_SIZE:}"
executor_thread_pool_size: "${TB_QUEUE_PUBSUB_EXECUTOR_THREAD_POOL_SIZE:0}"
queue-properties:
# Pub/Sub properties for Core subscribers, messages which will commit after ackDeadlineInSec period can be consumed again
core: "${TB_QUEUE_PUBSUB_CORE_QUEUE_PROPERTIES:ackDeadlineInSec:30;messageRetentionInSec:604800}"

2
transport/coap/src/main/resources/tb-coap-transport.yml

@ -296,7 +296,7 @@ queue:
# Number of messages per consumer
max_messages: "${TB_QUEUE_PUBSUB_MAX_MESSAGES:1000}"
# Thread pool size for pubsub queue executor provider. If not set - default pubsub executor provider value will be used (5 * number of available processors)
executor_thread_pool_size: "${TB_QUEUE_PUBSUB_EXECUTOR_THREAD_POOL_SIZE:}"
executor_thread_pool_size: "${TB_QUEUE_PUBSUB_EXECUTOR_THREAD_POOL_SIZE:0}"
queue-properties:
# Pub/Sub properties for Rule Engine subscribers, messages which will commit after ackDeadlineInSec period can be consumed again
rule-engine: "${TB_QUEUE_PUBSUB_RE_QUEUE_PROPERTIES:ackDeadlineInSec:30;messageRetentionInSec:604800}"

2
transport/http/src/main/resources/tb-http-transport.yml

@ -279,7 +279,7 @@ queue:
# Number of messages per a consumer
max_messages: "${TB_QUEUE_PUBSUB_MAX_MESSAGES:1000}"
# Thread pool size for pubsub queue executor provider. If not set - default pubsub executor provider value will be used (5 * number of available processors)
executor_thread_pool_size: "${TB_QUEUE_PUBSUB_EXECUTOR_THREAD_POOL_SIZE:}"
executor_thread_pool_size: "${TB_QUEUE_PUBSUB_EXECUTOR_THREAD_POOL_SIZE:0}"
queue-properties:
# Pub/Sub properties for Rule Engine subscribers, messages which will commit after ackDeadlineInSec period can be consume again
rule-engine: "${TB_QUEUE_PUBSUB_RE_QUEUE_PROPERTIES:ackDeadlineInSec:30;messageRetentionInSec:604800}"

2
transport/lwm2m/src/main/resources/tb-lwm2m-transport.yml

@ -375,7 +375,7 @@ queue:
# Number of messages per consumer
max_messages: "${TB_QUEUE_PUBSUB_MAX_MESSAGES:1000}"
# Thread pool size for pubsub queue executor provider. If not set - default pubsub executor provider value will be used (5 * number of available processors)
executor_thread_pool_size: "${TB_QUEUE_PUBSUB_EXECUTOR_THREAD_POOL_SIZE:}"
executor_thread_pool_size: "${TB_QUEUE_PUBSUB_EXECUTOR_THREAD_POOL_SIZE:0}"
queue-properties:
# Pub/Sub properties for Rule Engine subscribers, messages which will commit after ackDeadlineInSec period can be consumed again
rule-engine: "${TB_QUEUE_PUBSUB_RE_QUEUE_PROPERTIES:ackDeadlineInSec:30;messageRetentionInSec:604800}"

2
transport/mqtt/src/main/resources/tb-mqtt-transport.yml

@ -312,7 +312,7 @@ queue:
# Number of messages per consumer
max_messages: "${TB_QUEUE_PUBSUB_MAX_MESSAGES:1000}"
# Thread pool size for pubsub queue executor provider. If not set - default pubsub executor provider value will be used (5 * number of available processors)
executor_thread_pool_size: "${TB_QUEUE_PUBSUB_EXECUTOR_THREAD_POOL_SIZE:}"
executor_thread_pool_size: "${TB_QUEUE_PUBSUB_EXECUTOR_THREAD_POOL_SIZE:0}"
queue-properties:
# Pub/Sub properties for Rule Engine subscribers, messages which will commit after ackDeadlineInSec period can be consumed again
rule-engine: "${TB_QUEUE_PUBSUB_RE_QUEUE_PROPERTIES:ackDeadlineInSec:30;messageRetentionInSec:604800}"

2
transport/snmp/src/main/resources/tb-snmp-transport.yml

@ -265,7 +265,7 @@ queue:
# Number of messages per consumer
max_messages: "${TB_QUEUE_PUBSUB_MAX_MESSAGES:1000}"
# Thread pool size for pubsub queue executor provider. If not set - default pubsub executor provider value will be used (5 * number of available processors)
executor_thread_pool_size: "${TB_QUEUE_PUBSUB_EXECUTOR_THREAD_POOL_SIZE:}"
executor_thread_pool_size: "${TB_QUEUE_PUBSUB_EXECUTOR_THREAD_POOL_SIZE:0}"
queue-properties:
# Pub/Sub properties for Rule Engine subscribers, messages which will commit after ackDeadlineInSec period can be consumed again
rule-engine: "${TB_QUEUE_PUBSUB_RE_QUEUE_PROPERTIES:ackDeadlineInSec:30;messageRetentionInSec:604800}"

Loading…
Cancel
Save