From 1298f6b13047627e9199e2812ef5b7221fba1910 Mon Sep 17 00:00:00 2001 From: Andrii Shvaika Date: Mon, 29 Nov 2021 16:51:43 +0200 Subject: [PATCH] Fix delete attributes cluster notification --- .../DefaultSubscriptionManagerService.java | 17 +++++++++++------ 1 file changed, 11 insertions(+), 6 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/subscription/DefaultSubscriptionManagerService.java b/application/src/main/java/org/thingsboard/server/service/subscription/DefaultSubscriptionManagerService.java index 15cf6f1a17..018d80f013 100644 --- a/application/src/main/java/org/thingsboard/server/service/subscription/DefaultSubscriptionManagerService.java +++ b/application/src/main/java/org/thingsboard/server/service/subscription/DefaultSubscriptionManagerService.java @@ -222,7 +222,7 @@ public class DefaultSubscriptionManagerService extends TbApplicationEventListene } } return subscriptionUpdate; - }); + }, true); if (entityId.getEntityType() == EntityType.DEVICE) { updateDeviceInactivityTimeout(tenantId, entityId, ts); } @@ -256,7 +256,7 @@ public class DefaultSubscriptionManagerService extends TbApplicationEventListene } } return subscriptionUpdate; - }); + }, true); if (entityId.getEntityType() == EntityType.DEVICE) { if (TbAttributeSubscriptionScope.SERVER_SCOPE.name().equalsIgnoreCase(scope)) { updateDeviceInactivityTimeout(tenantId, entityId, attributes); @@ -333,14 +333,15 @@ public class DefaultSubscriptionManagerService extends TbApplicationEventListene } } return subscriptionUpdate; - }); + }, false); callback.onSuccess(); } private void onLocalTelemetrySubUpdate(EntityId entityId, Function castFunction, Predicate filterFunction, - Function> processFunction) { + Function> processFunction, + boolean ignoreEmptyUpdates) { Set entitySubscriptions = subscriptionsByEntityId.get(entityId); if (entitySubscriptions != null) { entitySubscriptions.stream().map(castFunction).filter(Objects::nonNull).filter(filterFunction).forEach(s -> { @@ -351,7 +352,7 @@ public class DefaultSubscriptionManagerService extends TbApplicationEventListene localSubscriptionService.onSubscriptionUpdate(s.getSessionId(), update, TbCallback.EMPTY); } else { TopicPartitionInfo tpi = partitionService.getNotificationsTopic(ServiceType.TB_CORE, s.getServiceId()); - toCoreNotificationsProducer.send(tpi, toProto(s, subscriptionUpdate), null); + toCoreNotificationsProducer.send(tpi, toProto(s, subscriptionUpdate, ignoreEmptyUpdates), null); } } }); @@ -467,6 +468,10 @@ public class DefaultSubscriptionManagerService extends TbApplicationEventListene } private TbProtoQueueMsg toProto(TbSubscription subscription, List updates) { + return toProto(subscription, updates, true); + } + + private TbProtoQueueMsg toProto(TbSubscription subscription, List updates, boolean ignoreEmptyUpdates) { TbSubscriptionUpdateProto.Builder builder = TbSubscriptionUpdateProto.newBuilder(); builder.setSessionId(subscription.getSessionId()); @@ -494,7 +499,7 @@ public class DefaultSubscriptionManagerService extends TbApplicationEventListene dataBuilder.addValue(strVal); } } - if (hasData) { + if (!ignoreEmptyUpdates || hasData) { builder.addData(dataBuilder.build()); } });