|
|
@ -213,7 +213,7 @@ public class DefaultSubscriptionManagerService extends TbApplicationEventListene |
|
|
}, s -> true, s -> { |
|
|
}, s -> true, s -> { |
|
|
List<TsKvEntry> subscriptionUpdate = null; |
|
|
List<TsKvEntry> subscriptionUpdate = null; |
|
|
for (TsKvEntry kv : ts) { |
|
|
for (TsKvEntry kv : ts) { |
|
|
if (isInTimeRange(s, kv.getTs()) && (s.isAllKeys() || s.getKeyStates().containsKey((kv.getKey())))) { |
|
|
if ((s.isAllKeys() || s.getKeyStates().containsKey((kv.getKey())))) { |
|
|
if (subscriptionUpdate == null) { |
|
|
if (subscriptionUpdate == null) { |
|
|
subscriptionUpdate = new ArrayList<>(); |
|
|
subscriptionUpdate = new ArrayList<>(); |
|
|
} |
|
|
} |
|
|
@ -375,11 +375,6 @@ public class DefaultSubscriptionManagerService extends TbApplicationEventListene |
|
|
} |
|
|
} |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
private boolean isInTimeRange(TbTimeseriesSubscription subscription, long kvTime) { |
|
|
|
|
|
return (subscription.getStartTime() == 0 || subscription.getStartTime() <= kvTime) |
|
|
|
|
|
&& (subscription.getEndTime() == 0 || subscription.getEndTime() >= kvTime); |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
private void removeSubscriptionFromEntityMap(TbSubscription sub) { |
|
|
private void removeSubscriptionFromEntityMap(TbSubscription sub) { |
|
|
Set<TbSubscription> entitySubSet = subscriptionsByEntityId.get(sub.getEntityId()); |
|
|
Set<TbSubscription> entitySubSet = subscriptionsByEntityId.get(sub.getEntityId()); |
|
|
if (entitySubSet != null) { |
|
|
if (entitySubSet != null) { |
|
|
@ -429,18 +424,9 @@ public class DefaultSubscriptionManagerService extends TbApplicationEventListene |
|
|
serviceId, subscription.getSessionId(), subscription.getSubscriptionId(), subscription.getEntityId()); |
|
|
serviceId, subscription.getSessionId(), subscription.getSubscriptionId(), subscription.getEntityId()); |
|
|
|
|
|
|
|
|
long curTs = System.currentTimeMillis(); |
|
|
long curTs = System.currentTimeMillis(); |
|
|
List<ReadTsKvQuery> queries = new ArrayList<>(); |
|
|
|
|
|
subscription.getKeyStates().forEach((key, value) -> { |
|
|
if (subscription.isLatestValues()) { |
|
|
if (curTs > value) { |
|
|
DonAsynchron.withCallback(tsService.findLatest(subscription.getTenantId(), subscription.getEntityId(), subscription.getKeyStates().keySet()), |
|
|
long startTs = subscription.getStartTime() > 0 ? Math.max(subscription.getStartTime(), value + 1L) : (value + 1L); |
|
|
|
|
|
long endTs = subscription.getEndTime() > 0 ? Math.min(subscription.getEndTime(), curTs) : curTs; |
|
|
|
|
|
if (startTs > 1) { |
|
|
|
|
|
queries.add(new BaseReadTsKvQuery(key, startTs, endTs, 0, 1000, Aggregation.NONE)); |
|
|
|
|
|
} |
|
|
|
|
|
} |
|
|
|
|
|
}); |
|
|
|
|
|
if (!queries.isEmpty()) { |
|
|
|
|
|
DonAsynchron.withCallback(tsService.findAll(subscription.getTenantId(), subscription.getEntityId(), queries), |
|
|
|
|
|
missedUpdates -> { |
|
|
missedUpdates -> { |
|
|
if (missedUpdates != null && !missedUpdates.isEmpty()) { |
|
|
if (missedUpdates != null && !missedUpdates.isEmpty()) { |
|
|
TopicPartitionInfo tpi = partitionService.getNotificationsTopic(ServiceType.TB_CORE, subscription.getServiceId()); |
|
|
TopicPartitionInfo tpi = partitionService.getNotificationsTopic(ServiceType.TB_CORE, subscription.getServiceId()); |
|
|
@ -449,6 +435,26 @@ public class DefaultSubscriptionManagerService extends TbApplicationEventListene |
|
|
}, |
|
|
}, |
|
|
e -> log.error("Failed to fetch missed updates.", e), |
|
|
e -> log.error("Failed to fetch missed updates.", e), |
|
|
tsCallBackExecutor); |
|
|
tsCallBackExecutor); |
|
|
|
|
|
} else { |
|
|
|
|
|
List<ReadTsKvQuery> queries = new ArrayList<>(); |
|
|
|
|
|
subscription.getKeyStates().forEach((key, value) -> { |
|
|
|
|
|
if (curTs > value) { |
|
|
|
|
|
long startTs = subscription.getStartTime() > 0 ? Math.max(subscription.getStartTime(), value + 1L) : (value + 1L); |
|
|
|
|
|
long endTs = subscription.getEndTime() > 0 ? Math.min(subscription.getEndTime(), curTs) : curTs; |
|
|
|
|
|
queries.add(new BaseReadTsKvQuery(key, startTs, endTs, 0, 1000, Aggregation.NONE)); |
|
|
|
|
|
} |
|
|
|
|
|
}); |
|
|
|
|
|
if (!queries.isEmpty()) { |
|
|
|
|
|
DonAsynchron.withCallback(tsService.findAll(subscription.getTenantId(), subscription.getEntityId(), queries), |
|
|
|
|
|
missedUpdates -> { |
|
|
|
|
|
if (missedUpdates != null && !missedUpdates.isEmpty()) { |
|
|
|
|
|
TopicPartitionInfo tpi = partitionService.getNotificationsTopic(ServiceType.TB_CORE, subscription.getServiceId()); |
|
|
|
|
|
toCoreNotificationsProducer.send(tpi, toProto(subscription, missedUpdates), null); |
|
|
|
|
|
} |
|
|
|
|
|
}, |
|
|
|
|
|
e -> log.error("Failed to fetch missed updates.", e), |
|
|
|
|
|
tsCallBackExecutor); |
|
|
|
|
|
} |
|
|
} |
|
|
} |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
@ -470,12 +476,19 @@ public class DefaultSubscriptionManagerService extends TbApplicationEventListene |
|
|
data.forEach((key, value) -> { |
|
|
data.forEach((key, value) -> { |
|
|
TbSubscriptionUpdateValueListProto.Builder dataBuilder = TbSubscriptionUpdateValueListProto.newBuilder(); |
|
|
TbSubscriptionUpdateValueListProto.Builder dataBuilder = TbSubscriptionUpdateValueListProto.newBuilder(); |
|
|
dataBuilder.setKey(key); |
|
|
dataBuilder.setKey(key); |
|
|
value.forEach(v -> { |
|
|
boolean hasData = false; |
|
|
|
|
|
for (Object v : value) { |
|
|
Object[] array = (Object[]) v; |
|
|
Object[] array = (Object[]) v; |
|
|
dataBuilder.addTs((long) array[0]); |
|
|
dataBuilder.addTs((long) array[0]); |
|
|
dataBuilder.addValue((String) array[1]); |
|
|
String strVal = (String) array[1]; |
|
|
}); |
|
|
if (strVal != null) { |
|
|
builder.addData(dataBuilder.build()); |
|
|
hasData = true; |
|
|
|
|
|
dataBuilder.addValue(strVal); |
|
|
|
|
|
} |
|
|
|
|
|
} |
|
|
|
|
|
if (hasData) { |
|
|
|
|
|
builder.addData(dataBuilder.build()); |
|
|
|
|
|
} |
|
|
}); |
|
|
}); |
|
|
|
|
|
|
|
|
ToCoreNotificationMsg toCoreMsg = ToCoreNotificationMsg.newBuilder().setToLocalSubscriptionServiceMsg( |
|
|
ToCoreNotificationMsg toCoreMsg = ToCoreNotificationMsg.newBuilder().setToLocalSubscriptionServiceMsg( |
|
|
|