diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/queue/Queue.java b/common/data/src/main/java/org/thingsboard/server/common/data/queue/Queue.java index 84776c88f4..004bd0057d 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/queue/Queue.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/queue/Queue.java @@ -29,6 +29,7 @@ public class Queue extends BaseData implements HasName, HasTenantId { private String topic; private int pollInterval; private int partitions; + private boolean consumerPerPartition; private long packProcessingTimeout; private SubmitStrategy submitStrategy; private ProcessingStrategy processingStrategy; diff --git a/dao/src/main/java/org/thingsboard/server/dao/model/ModelConstants.java b/dao/src/main/java/org/thingsboard/server/dao/model/ModelConstants.java index cb209c2ee7..73c19cab3c 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/model/ModelConstants.java +++ b/dao/src/main/java/org/thingsboard/server/dao/model/ModelConstants.java @@ -592,6 +592,7 @@ public class ModelConstants { public static final String QUEUE_TOPIC_PROPERTY = "topic"; public static final String QUEUE_POLL_INTERVAL_PROPERTY = "poll_interval"; public static final String QUEUE_PARTITIONS_PROPERTY = "partitions"; + public static final String QUEUE_CONSUMER_PER_PARTITION = "consumer_per_partition"; public static final String QUEUE_PACK_PROCESSING_TIMEOUT_PROPERTY = "pack_processing_timeout"; public static final String QUEUE_SUBMIT_STRATEGY_PROPERTY = "submit_strategy"; public static final String QUEUE_PROCESSING_STRATEGY_PROPERTY = "processing_strategy"; diff --git a/dao/src/main/java/org/thingsboard/server/dao/model/sql/QueueEntity.java b/dao/src/main/java/org/thingsboard/server/dao/model/sql/QueueEntity.java index 3791ab4ac3..fefe686c9e 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/model/sql/QueueEntity.java +++ b/dao/src/main/java/org/thingsboard/server/dao/model/sql/QueueEntity.java @@ -59,6 +59,9 @@ public class QueueEntity extends BaseSqlEntity { @Column(name = ModelConstants.QUEUE_PARTITIONS_PROPERTY) private int partitions; + @Column(name = ModelConstants.QUEUE_CONSUMER_PER_PARTITION) + private boolean consumerPerPartition; + @Column(name = ModelConstants.QUEUE_PACK_PROCESSING_TIMEOUT_PROPERTY) private long packProcessingTimeout; @@ -83,6 +86,7 @@ public class QueueEntity extends BaseSqlEntity { this.topic = queue.getTopic(); this.pollInterval = queue.getPollInterval(); this.partitions = queue.getPartitions(); + this.consumerPerPartition = queue.isConsumerPerPartition(); this.packProcessingTimeout = queue.getPackProcessingTimeout(); this.submitStrategy = mapper.valueToTree(queue.getSubmitStrategy()); this.processingStrategy = mapper.valueToTree(queue.getProcessingStrategy()); @@ -97,6 +101,7 @@ public class QueueEntity extends BaseSqlEntity { queue.setTopic(topic); queue.setPollInterval(pollInterval); queue.setPartitions(partitions); + queue.setConsumerPerPartition(consumerPerPartition); queue.setPackProcessingTimeout(packProcessingTimeout); queue.setSubmitStrategy(mapper.convertValue(this.submitStrategy, SubmitStrategy.class)); queue.setProcessingStrategy(mapper.convertValue(this.processingStrategy, ProcessingStrategy.class)); diff --git a/dao/src/main/resources/sql/schema-entities.sql b/dao/src/main/resources/sql/schema-entities.sql index 5335ce6f5f..54c3cbfa93 100644 --- a/dao/src/main/resources/sql/schema-entities.sql +++ b/dao/src/main/resources/sql/schema-entities.sql @@ -642,6 +642,7 @@ CREATE TABLE IF NOT EXISTS queue( topic varchar(255), poll_interval int, partitions int, + consumer_per_partition boolean, pack_processing_timeout bigint, submit_strategy varchar(255), processing_strategy varchar(255) diff --git a/dao/src/test/java/org/thingsboard/server/dao/service/BaseQueueServiceTest.java b/dao/src/test/java/org/thingsboard/server/dao/service/BaseQueueServiceTest.java index 173f596287..8899224bd5 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/service/BaseQueueServiceTest.java +++ b/dao/src/test/java/org/thingsboard/server/dao/service/BaseQueueServiceTest.java @@ -150,7 +150,7 @@ public abstract class BaseQueueServiceTest extends AbstractServiceTest { } @Test(expected = DataValidationException.class) - public void testSaveQueueWithEmptyPoolInterval() { + public void testSaveQueueWithEmptyPollInterval() { Queue queue = new Queue(); queue.setTenantId(tenantId); queue.setName("Test"); diff --git a/ui-ngx/src/app/modules/home/pages/admin/queue/queue.component.html b/ui-ngx/src/app/modules/home/pages/admin/queue/queue.component.html index ae163d9ba7..ae8f08e348 100644 --- a/ui-ngx/src/app/modules/home/pages/admin/queue/queue.component.html +++ b/ui-ngx/src/app/modules/home/pages/admin/queue/queue.component.html @@ -56,6 +56,16 @@ {{ 'queue.partitions-min-value' | translate }} + + + + + + +
{{ 'queue.consumer-per-partition' | translate }}
+
{{'queue.consumer-per-partition-hint' | translate}}
+
+ queue.processing-timeout @@ -69,7 +79,6 @@
- diff --git a/ui-ngx/src/app/modules/home/pages/admin/queue/queue.component.ts b/ui-ngx/src/app/modules/home/pages/admin/queue/queue.component.ts index b694e21ddf..f9a60773df 100644 --- a/ui-ngx/src/app/modules/home/pages/admin/queue/queue.component.ts +++ b/ui-ngx/src/app/modules/home/pages/admin/queue/queue.component.ts @@ -71,6 +71,7 @@ export class QueueComponent extends EntityComponent { entity && entity.partitions ? entity.partitions : 10, [Validators.min(1), Validators.required] ], + consumerPerPartition: [entity ? entity.consumerPerPartition : false, []], packProcessingTimeout: [ entity && entity.packProcessingTimeout ? entity.packProcessingTimeout : 2000, [Validators.min(1), Validators.required] @@ -118,6 +119,7 @@ export class QueueComponent extends EntityComponent { name: entity.name, pollInterval: entity.pollInterval, partitions: entity.partitions, + consumerPerPartition: entity.consumerPerPartition, packProcessingTimeout: entity.packProcessingTimeout, submitStrategy: { type: entity.submitStrategy?.type, diff --git a/ui-ngx/src/app/shared/models/queue.models.ts b/ui-ngx/src/app/shared/models/queue.models.ts index 2fe2f511f4..508cf2aee1 100644 --- a/ui-ngx/src/app/shared/models/queue.models.ts +++ b/ui-ngx/src/app/shared/models/queue.models.ts @@ -45,6 +45,7 @@ export enum QueueProcessingStrategyTypes { export interface QueueInfo extends BaseData { packProcessingTimeout: number; partitions: number; + consumerPerPartition: boolean, pollInterval: number; processingStrategy: { type: QueueProcessingStrategyTypes, diff --git a/ui-ngx/src/assets/locale/locale.constant-en_US.json b/ui-ngx/src/assets/locale/locale.constant-en_US.json index 75dbdf4f5f..4e30b773b1 100644 --- a/ui-ngx/src/assets/locale/locale.constant-en_US.json +++ b/ui-ngx/src/assets/locale/locale.constant-en_US.json @@ -2714,6 +2714,8 @@ "processing-strategy": "Processing Strategy", "poll-interval": "Poll interval", "partitions": "Partitions", + "consumer-per-partition": "Consumer per partition", + "consumer-per-partition-hint": "Enable separate consumer(s) per each partition", "processing-timeout": "Processing timeout", "batch-size": "Batch size", "retries": "Retries (0 - unlimited)",