Browse Source

WIP: cluster mode implementation

pull/12510/head
Andrii Shvaika 2 years ago
parent
commit
85119d0247
  1. 4
      application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldExecutionService.java
  2. 11
      application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldExecutionService.java
  3. 6
      application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCalculatedFieldConsumerService.java
  4. 42
      application/src/main/java/org/thingsboard/server/service/queue/DefaultTbClusterService.java
  5. 3
      application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionService.java
  6. 2
      application/src/main/resources/thingsboard.yml

4
application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldExecutionService.java

@ -23,6 +23,8 @@ import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldEntit
import org.thingsboard.server.gen.transport.TransportProtos.ComponentLifecycleMsgProto; import org.thingsboard.server.gen.transport.TransportProtos.ComponentLifecycleMsgProto;
import org.thingsboard.server.service.cf.telemetry.CalculatedFieldTelemetryUpdateRequest; import org.thingsboard.server.service.cf.telemetry.CalculatedFieldTelemetryUpdateRequest;
import java.util.List;
public interface CalculatedFieldExecutionService { public interface CalculatedFieldExecutionService {
/** /**
@ -32,7 +34,7 @@ public interface CalculatedFieldExecutionService {
*/ */
void pushRequestToQueue(TimeseriesSaveRequest request, TimeseriesSaveResult result); void pushRequestToQueue(TimeseriesSaveRequest request, TimeseriesSaveResult result);
void pushRequestToQueue(AttributesSaveRequest request); void pushRequestToQueue(AttributesSaveRequest request, List<Long> result);
// void pushEntityUpdateMsg(TransportProtos.CalculatedFieldEntityUpdateMsgProto proto, TbCallback callback); // void pushEntityUpdateMsg(TransportProtos.CalculatedFieldEntityUpdateMsgProto proto, TbCallback callback);

11
application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldExecutionService.java

@ -197,14 +197,16 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas
} }
@Override @Override
public void pushRequestToQueue(AttributesSaveRequest request) { public void pushRequestToQueue(AttributesSaveRequest request, List<Long> result) {
var tenantId = request.getTenantId(); var tenantId = request.getTenantId();
var entityId = request.getEntityId(); var entityId = request.getEntityId();
checkEntityAndPushToQueue(tenantId, entityId, cf -> cf.matches(request.getEntries(), request.getScope()), cf -> cf.linkMatches(entityId, request.getEntries(), request.getScope()), checkEntityAndPushToQueue(tenantId, entityId, cf -> cf.matches(request.getEntries(), request.getScope()), cf -> cf.linkMatches(entityId, request.getEntries(), request.getScope()),
() -> toCalculatedFieldTelemetryMsgProto(request), request.getCallback()); () -> toCalculatedFieldTelemetryMsgProto(request, result), request.getCallback());
} }
private void checkEntityAndPushToQueue(TenantId tenantId, EntityId entityId, Predicate<CalculatedFieldCtx> mainEntityFilter, Predicate<CalculatedFieldCtx> linkedEntityFilter, Supplier<ToCalculatedFieldMsg> msg, FutureCallback<Void> callback) { private void checkEntityAndPushToQueue(TenantId tenantId, EntityId entityId,
Predicate<CalculatedFieldCtx> mainEntityFilter, Predicate<CalculatedFieldCtx> linkedEntityFilter,
Supplier<ToCalculatedFieldMsg> msg, FutureCallback<Void> callback) {
boolean send = checkEntityForCalculatedFields(tenantId, entityId, mainEntityFilter, linkedEntityFilter); boolean send = checkEntityForCalculatedFields(tenantId, entityId, mainEntityFilter, linkedEntityFilter);
if (send) { if (send) {
clusterService.pushMsgToCalculatedFields(tenantId, entityId, msg.get(), wrap(callback)); clusterService.pushMsgToCalculatedFields(tenantId, entityId, msg.get(), wrap(callback));
@ -827,7 +829,8 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas
return msg.build(); return msg.build();
} }
private ToCalculatedFieldMsg toCalculatedFieldTelemetryMsgProto(AttributesSaveRequest request) { private ToCalculatedFieldMsg toCalculatedFieldTelemetryMsgProto(AttributesSaveRequest request, List<Long> result) {
//TODO: IM Use result in both methods to update the versions of telemetry/attributes.
ToCalculatedFieldMsg.Builder msg = ToCalculatedFieldMsg.newBuilder(); ToCalculatedFieldMsg.Builder msg = ToCalculatedFieldMsg.newBuilder();
CalculatedFieldTelemetryMsgProto.Builder telemetryMsg = buildTelemetryMsgProto(request.getTenantId(), request.getEntityId(), request.getPreviousCalculatedFieldIds()); CalculatedFieldTelemetryMsgProto.Builder telemetryMsg = buildTelemetryMsgProto(request.getTenantId(), request.getEntityId(), request.getPreviousCalculatedFieldIds());

6
application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCalculatedFieldConsumerService.java

@ -62,9 +62,9 @@ import java.util.UUID;
@Slf4j @Slf4j
public class DefaultTbCalculatedFieldConsumerService extends AbstractConsumerService<ToCalculatedFieldNotificationMsg> implements TbCalculatedFieldConsumerService { public class DefaultTbCalculatedFieldConsumerService extends AbstractConsumerService<ToCalculatedFieldNotificationMsg> implements TbCalculatedFieldConsumerService {
@Value("${queue.calculated_fields.poll_interval}") @Value("${queue.calculated_fields.poll_interval:25}")
private long pollInterval; private long pollInterval;
@Value("${queue.calculated_fields.pack_processing_timeout}") @Value("${queue.calculated_fields.pack_processing_timeout:60000}")
private long packProcessingTimeout; private long packProcessingTimeout;
@Value("${queue.calculated_fields.consumer_per_partition:true}") @Value("${queue.calculated_fields.consumer_per_partition:true}")
private boolean consumerPerPartition; private boolean consumerPerPartition;
@ -102,7 +102,7 @@ public class DefaultTbCalculatedFieldConsumerService extends AbstractConsumerSer
this.calculatedFieldsExecutor = MoreExecutors.listeningDecorator(ThingsBoardExecutors.newWorkStealingPool(poolSize, "tb-cf-executor")); // TODO: multiple threads. this.calculatedFieldsExecutor = MoreExecutors.listeningDecorator(ThingsBoardExecutors.newWorkStealingPool(poolSize, "tb-cf-executor")); // TODO: multiple threads.
this.mainConsumer = MainQueueConsumerManager.<TbProtoQueueMsg<ToCalculatedFieldMsg>, CalculatedFieldQueueConfig>builder() this.mainConsumer = MainQueueConsumerManager.<TbProtoQueueMsg<ToCalculatedFieldMsg>, CalculatedFieldQueueConfig>builder()
.queueKey(new QueueKey(ServiceType.TB_CORE)) .queueKey(new QueueKey(ServiceType.TB_RULE_ENGINE))
.config(CalculatedFieldQueueConfig.of(consumerPerPartition, (int) pollInterval)) .config(CalculatedFieldQueueConfig.of(consumerPerPartition, (int) pollInterval))
.msgPackProcessor(this::processMsgs) .msgPackProcessor(this::processMsgs)
.consumerCreator((config, partitionId) -> queueFactory.createToCalculatedFieldMsgConsumer()) .consumerCreator((config, partitionId) -> queueFactory.createToCalculatedFieldMsgConsumer())

42
application/src/main/java/org/thingsboard/server/service/queue/DefaultTbClusterService.java

@ -102,6 +102,7 @@ import org.thingsboard.server.service.ota.OtaPackageStateService;
import org.thingsboard.server.service.profile.TbAssetProfileCache; import org.thingsboard.server.service.profile.TbAssetProfileCache;
import org.thingsboard.server.service.profile.TbDeviceProfileCache; import org.thingsboard.server.service.profile.TbDeviceProfileCache;
import java.util.HashSet;
import java.util.List; import java.util.List;
import java.util.Objects; import java.util.Objects;
import java.util.Optional; import java.util.Optional;
@ -570,7 +571,9 @@ public class DefaultTbClusterService implements TbClusterService {
private void broadcast(ComponentLifecycleMsg msg) { private void broadcast(ComponentLifecycleMsg msg) {
ComponentLifecycleMsgProto componentLifecycleMsgProto = toProto(msg); ComponentLifecycleMsgProto componentLifecycleMsgProto = toProto(msg);
TbQueueProducer<TbProtoQueueMsg<ToRuleEngineNotificationMsg>> toRuleEngineProducer = producerProvider.getRuleEngineNotificationsMsgProducer(); TbQueueProducer<TbProtoQueueMsg<ToRuleEngineNotificationMsg>> toRuleEngineProducer = producerProvider.getRuleEngineNotificationsMsgProducer();
TbQueueProducer<TbProtoQueueMsg<ToCalculatedFieldNotificationMsg>> toCalculatedFieldProducer = producerProvider.getCalculatedFieldsNotificationsMsgProducer();
Set<String> tbRuleEngineServices = partitionService.getAllServiceIds(ServiceType.TB_RULE_ENGINE); Set<String> tbRuleEngineServices = partitionService.getAllServiceIds(ServiceType.TB_RULE_ENGINE);
Set<String> tbCalculatedFieldServices = new HashSet<>(tbRuleEngineServices);
EntityType entityType = msg.getEntityId().getEntityType(); EntityType entityType = msg.getEntityId().getEntityType();
if (entityType.equals(EntityType.TENANT) if (entityType.equals(EntityType.TENANT)
|| entityType.equals(EntityType.TENANT_PROFILE) || entityType.equals(EntityType.TENANT_PROFILE)
@ -581,7 +584,6 @@ public class DefaultTbClusterService implements TbClusterService {
|| (entityType.equals(EntityType.DEVICE) && msg.getEvent() == ComponentLifecycleEvent.UPDATED) || (entityType.equals(EntityType.DEVICE) && msg.getEvent() == ComponentLifecycleEvent.UPDATED)
|| entityType.equals(EntityType.ENTITY_VIEW) || entityType.equals(EntityType.ENTITY_VIEW)
|| entityType.equals(EntityType.NOTIFICATION_RULE) || entityType.equals(EntityType.NOTIFICATION_RULE)
|| entityType.equals(EntityType.CALCULATED_FIELD)
) { ) {
TbQueueProducer<TbProtoQueueMsg<ToCoreNotificationMsg>> toCoreNfProducer = producerProvider.getTbCoreNotificationsMsgProducer(); TbQueueProducer<TbProtoQueueMsg<ToCoreNotificationMsg>> toCoreNfProducer = producerProvider.getTbCoreNotificationsMsgProducer();
Set<String> tbCoreServices = partitionService.getAllServiceIds(ServiceType.TB_CORE); Set<String> tbCoreServices = partitionService.getAllServiceIds(ServiceType.TB_CORE);
@ -594,6 +596,14 @@ public class DefaultTbClusterService implements TbClusterService {
// No need to push notifications twice // No need to push notifications twice
tbRuleEngineServices.removeAll(tbCoreServices); tbRuleEngineServices.removeAll(tbCoreServices);
} }
if (entityType.equals(EntityType.CALCULATED_FIELD)) {
for (String serviceId : tbCalculatedFieldServices) {
TopicPartitionInfo tpi = topicService.getCalculatedFieldNotificationsTopic(serviceId);
ToCalculatedFieldNotificationMsg toCfNotificationMsg = ToCalculatedFieldNotificationMsg.newBuilder().setComponentLifecycle(componentLifecycleMsgProto).build();
toCalculatedFieldProducer.send(tpi, new TbProtoQueueMsg<>(msg.getEntityId().getId(), toCfNotificationMsg), null);
toRuleEngineNfs.incrementAndGet(); // TODO: add separate counter when we will have new ServiceType.CALCULATED_FIELDS
}
}
for (String serviceId : tbRuleEngineServices) { for (String serviceId : tbRuleEngineServices) {
TopicPartitionInfo tpi = topicService.getNotificationsTopic(ServiceType.TB_RULE_ENGINE, serviceId); TopicPartitionInfo tpi = topicService.getNotificationsTopic(ServiceType.TB_RULE_ENGINE, serviceId);
ToRuleEngineNotificationMsg toRuleEngineMsg = ToRuleEngineNotificationMsg.newBuilder().setComponentLifecycle(componentLifecycleMsgProto).build(); ToRuleEngineNotificationMsg toRuleEngineMsg = ToRuleEngineNotificationMsg.newBuilder().setComponentLifecycle(componentLifecycleMsgProto).build();
@ -641,11 +651,11 @@ public class DefaultTbClusterService implements TbClusterService {
if (deviceNameChanged) { if (deviceNameChanged) {
gatewayNotificationsService.onDeviceUpdated(device, old); gatewayNotificationsService.onDeviceUpdated(device, old);
} }
boolean deviceTypeChanged = !device.getType().equals(old.getType()); boolean deviceProfileChanged = !device.getDeviceProfileId().equals(old.getDeviceProfileId());
if (deviceTypeChanged) { if (deviceProfileChanged) {
handleEntityProfileUpdatedEvent(device.getTenantId(), device.getId(), old.getDeviceProfileId(), device.getDeviceProfileId()); handleEntityProfileUpdatedEvent(device.getTenantId(), device.getId(), old.getDeviceProfileId(), device.getDeviceProfileId());
} }
if (deviceNameChanged || deviceTypeChanged) { if (deviceNameChanged || deviceProfileChanged) {
pushMsgToCore(new DeviceNameOrTypeUpdateMsg(device.getTenantId(), device.getId(), device.getName(), device.getType()), null); pushMsgToCore(new DeviceNameOrTypeUpdateMsg(device.getTenantId(), device.getId(), device.getName(), device.getType()), null);
} }
} else { } else {
@ -802,37 +812,13 @@ public class DefaultTbClusterService implements TbClusterService {
@Override @Override
public void onCalculatedFieldUpdated(CalculatedField calculatedField, CalculatedField oldCalculatedField) { public void onCalculatedFieldUpdated(CalculatedField calculatedField, CalculatedField oldCalculatedField) {
var created = oldCalculatedField == null; var created = oldCalculatedField == null;
broadcastEntityChangeToTransport(calculatedField.getTenantId(), calculatedField.getId(), calculatedField, null);
broadcastEntityStateChangeEvent(calculatedField.getTenantId(), calculatedField.getId(), created ? ComponentLifecycleEvent.CREATED : ComponentLifecycleEvent.UPDATED); broadcastEntityStateChangeEvent(calculatedField.getTenantId(), calculatedField.getId(), created ? ComponentLifecycleEvent.CREATED : ComponentLifecycleEvent.UPDATED);
sendCalculatedFieldEvent(calculatedField.getTenantId(), calculatedField.getId(), created, !created, false);
} }
@Override @Override
public void onCalculatedFieldDeleted(TenantId tenantId, CalculatedField calculatedField, TbQueueCallback callback) { public void onCalculatedFieldDeleted(TenantId tenantId, CalculatedField calculatedField, TbQueueCallback callback) {
CalculatedFieldId calculatedFieldId = calculatedField.getId(); CalculatedFieldId calculatedFieldId = calculatedField.getId();
broadcastEntityDeleteToTransport(tenantId, calculatedFieldId, calculatedField.getName(), callback);
broadcastEntityStateChangeEvent(tenantId, calculatedFieldId, ComponentLifecycleEvent.DELETED); broadcastEntityStateChangeEvent(tenantId, calculatedFieldId, ComponentLifecycleEvent.DELETED);
sendCalculatedFieldEvent(tenantId, calculatedFieldId, false, false, true);
}
private void sendCalculatedFieldEvent(TenantId tenantId, CalculatedFieldId calculatedFieldId, boolean added, boolean updated, boolean deleted) {
ComponentLifecycleMsgProto.Builder builder = ComponentLifecycleMsgProto.newBuilder();
builder.setTenantIdMSB(tenantId.getId().getMostSignificantBits());
builder.setTenantIdLSB(tenantId.getId().getLeastSignificantBits());
builder.setEntityType(TransportProtos.EntityTypeProto.CALCULATED_FIELD);
builder.setEntityIdMSB(calculatedFieldId.getId().getMostSignificantBits());
builder.setEntityIdLSB(calculatedFieldId.getId().getLeastSignificantBits());
TransportProtos.ComponentLifecycleEvent event;
if (added) {
event = TransportProtos.ComponentLifecycleEvent.CREATED;
} else if (updated) {
event = TransportProtos.ComponentLifecycleEvent.UPDATED;
} else {
event = TransportProtos.ComponentLifecycleEvent.DELETED;
}
builder.setEvent(event);
pushNotificationToCalculatedFields(tenantId, calculatedFieldId, ToCalculatedFieldNotificationMsg.newBuilder().setComponentLifecycle(builder).build(), null);
} }
private void handleEntityProfileUpdatedEvent(TenantId tenantId, EntityId entityId, EntityId oldProfileId, EntityId newProfileId) { private void handleEntityProfileUpdatedEvent(TenantId tenantId, EntityId entityId, EntityId oldProfileId, EntityId newProfileId) {

3
application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionService.java

@ -167,9 +167,8 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer
public void saveAttributesInternal(AttributesSaveRequest request) { public void saveAttributesInternal(AttributesSaveRequest request) {
log.trace("Executing saveInternal [{}]", request); log.trace("Executing saveInternal [{}]", request);
ListenableFuture<List<Long>> saveFuture = attrService.save(request.getTenantId(), request.getEntityId(), request.getScope(), request.getEntries()); ListenableFuture<List<Long>> saveFuture = attrService.save(request.getTenantId(), request.getEntityId(), request.getScope(), request.getEntries());
addMainCallback(saveFuture, request.getCallback());
DonAsynchron.withCallback(saveFuture, result -> { DonAsynchron.withCallback(saveFuture, result -> {
calculatedFieldExecutionService.pushRequestToQueue(request); calculatedFieldExecutionService.pushRequestToQueue(request, result);
}, safeCallback(request.getCallback()), tsCallBackExecutor); }, safeCallback(request.getCallback()), tsCallBackExecutor);
addWsCallback(saveFuture, success -> onAttributesUpdate(request.getTenantId(), request.getEntityId(), request.getScope().name(), request.getEntries(), request.isNotifyDevice())); addWsCallback(saveFuture, success -> onAttributesUpdate(request.getTenantId(), request.getEntityId(), request.getScope().name(), request.getEntries(), request.isNotifyDevice()));
addCallback(saveFuture, success -> calculatedFieldExecutionService.onTelemetryUpdate(new CalculatedFieldAttributeUpdateRequest(request)), tsCallBackExecutor); addCallback(saveFuture, success -> calculatedFieldExecutionService.onTelemetryUpdate(new CalculatedFieldAttributeUpdateRequest(request)), tsCallBackExecutor);

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

@ -1750,7 +1750,7 @@ queue:
# Amount of partitions used by CF microservices # Amount of partitions used by CF microservices
partitions: "${TB_QUEUE_CF_PARTITIONS:10}" partitions: "${TB_QUEUE_CF_PARTITIONS:10}"
# Timeout for processing a message pack by CF microservices # Timeout for processing a message pack by CF microservices
pack_processing_timeout: "${TB_QUEUE_CF_PACK_PROCESSING_TIMEOUT_MS:2000}" pack_processing_timeout: "${TB_QUEUE_CF_PACK_PROCESSING_TIMEOUT_MS:60000}"
# Enable/disable a separate consumer per partition for CF queue # Enable/disable a separate consumer per partition for CF queue
consumer_per_partition: "${TB_QUEUE_CF_CONSUMER_PER_PARTITION:true}" consumer_per_partition: "${TB_QUEUE_CF_CONSUMER_PER_PARTITION:true}"
# Thread pool size for processing of the incoming messages # Thread pool size for processing of the incoming messages

Loading…
Cancel
Save