Browse Source

added consumerPerPartition to Queue entity

pull/6134/head
YevhenBondarenko 5 years ago
parent
commit
5c2ee4434c
  1. 1
      common/data/src/main/java/org/thingsboard/server/common/data/queue/Queue.java
  2. 1
      dao/src/main/java/org/thingsboard/server/dao/model/ModelConstants.java
  3. 5
      dao/src/main/java/org/thingsboard/server/dao/model/sql/QueueEntity.java
  4. 1
      dao/src/main/resources/sql/schema-entities.sql
  5. 2
      dao/src/test/java/org/thingsboard/server/dao/service/BaseQueueServiceTest.java
  6. 11
      ui-ngx/src/app/modules/home/pages/admin/queue/queue.component.html
  7. 2
      ui-ngx/src/app/modules/home/pages/admin/queue/queue.component.ts
  8. 1
      ui-ngx/src/app/shared/models/queue.models.ts
  9. 2
      ui-ngx/src/assets/locale/locale.constant-en_US.json

1
common/data/src/main/java/org/thingsboard/server/common/data/queue/Queue.java

@ -29,6 +29,7 @@ public class Queue extends BaseData<QueueId> implements HasName, HasTenantId {
private String topic;
private int pollInterval;
private int partitions;
private boolean consumerPerPartition;
private long packProcessingTimeout;
private SubmitStrategy submitStrategy;
private ProcessingStrategy processingStrategy;

1
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";

5
dao/src/main/java/org/thingsboard/server/dao/model/sql/QueueEntity.java

@ -59,6 +59,9 @@ public class QueueEntity extends BaseSqlEntity<Queue> {
@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<Queue> {
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> {
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));

1
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)

2
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");

11
ui-ngx/src/app/modules/home/pages/admin/queue/queue.component.html

@ -56,6 +56,16 @@
{{ 'queue.partitions-min-value' | translate }}
</mat-error>
</mat-form-field>
<!-- <tb-checkbox formControlName="consumerPerPartition" style="display: block; padding-bottom: 16px;">-->
<!-- {{ 'queue.consumer-per-partition' | translate }}-->
<!-- </tb-checkbox>-->
<mat-checkbox class="hinted-checkbox" formControlName="consumerPerPartition">
<div>{{ 'queue.consumer-per-partition' | translate }}</div>
<div class="tb-hint">{{'queue.consumer-per-partition-hint' | translate}}</div>
</mat-checkbox>
<mat-form-field class="mat-block">
<mat-label translate>queue.processing-timeout</mat-label>
<input type="number" matInput formControlName="packProcessingTimeout" required>
@ -69,7 +79,6 @@
</mat-form-field>
<div class="mat-accordion-container">
<mat-accordion [multi]="true">
<mat-expansion-panel #panel1 hideToggle>
<mat-expansion-panel-header>

2
ui-ngx/src/app/modules/home/pages/admin/queue/queue.component.ts

@ -71,6 +71,7 @@ export class QueueComponent extends EntityComponent<QueueInfo> {
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<QueueInfo> {
name: entity.name,
pollInterval: entity.pollInterval,
partitions: entity.partitions,
consumerPerPartition: entity.consumerPerPartition,
packProcessingTimeout: entity.packProcessingTimeout,
submitStrategy: {
type: entity.submitStrategy?.type,

1
ui-ngx/src/app/shared/models/queue.models.ts

@ -45,6 +45,7 @@ export enum QueueProcessingStrategyTypes {
export interface QueueInfo extends BaseData<QueueId> {
packProcessingTimeout: number;
partitions: number;
consumerPerPartition: boolean,
pollInterval: number;
processingStrategy: {
type: QueueProcessingStrategyTypes,

2
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)",

Loading…
Cancel
Save