diff --git a/application/src/main/java/org/thingsboard/server/service/subscription/TbAbstractDataSubCtx.java b/application/src/main/java/org/thingsboard/server/service/subscription/TbAbstractDataSubCtx.java index 26e97b6f2c..271de08733 100644 --- a/application/src/main/java/org/thingsboard/server/service/subscription/TbAbstractDataSubCtx.java +++ b/application/src/main/java/org/thingsboard/server/service/subscription/TbAbstractDataSubCtx.java @@ -18,6 +18,7 @@ package org.thingsboard.server.service.subscription; import lombok.Getter; import lombok.extern.slf4j.Slf4j; import org.thingsboard.server.common.data.id.EntityId; +import org.thingsboard.server.common.data.kv.Aggregation; import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.query.AbstractDataQuery; import org.thingsboard.server.common.data.query.EntityData; @@ -197,6 +198,9 @@ public abstract class TbAbstractDataSubCtx { long ts = Arrays.stream(v).map(TsValue::getTs).max(Long::compareTo).orElse(0L); log.trace("[{}][{}] Updating key: {} with ts: {}", serviceId, cmdId, k, ts); + if (!Aggregation.NONE.equals(getCurrentAggregation()) && ts < endTs) { + ts = endTs; + } keyStates.put(k, ts); }); } @@ -247,4 +251,5 @@ public abstract class TbAbstractDataSubCtx { } } + @Override + protected Aggregation getCurrentAggregation() { + return Aggregation.NONE; + } + private void sendWsMsg(String sessionId, AlarmSubscriptionUpdate subscriptionUpdate) { Alarm alarm = subscriptionUpdate.getAlarm(); AlarmId alarmId = alarm.getId(); diff --git a/application/src/main/java/org/thingsboard/server/service/subscription/TbEntityDataSubCtx.java b/application/src/main/java/org/thingsboard/server/service/subscription/TbEntityDataSubCtx.java index 385661f177..016ed198ec 100644 --- a/application/src/main/java/org/thingsboard/server/service/subscription/TbEntityDataSubCtx.java +++ b/application/src/main/java/org/thingsboard/server/service/subscription/TbEntityDataSubCtx.java @@ -19,6 +19,7 @@ import lombok.Getter; import lombok.Setter; import lombok.extern.slf4j.Slf4j; import org.thingsboard.server.common.data.id.EntityId; +import org.thingsboard.server.common.data.kv.Aggregation; import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.query.EntityData; import org.thingsboard.server.common.data.query.EntityDataQuery; @@ -86,6 +87,11 @@ public class TbEntityDataSubCtx extends TbAbstractDataSubCtx { } } + @Override + protected Aggregation getCurrentAggregation() { + return (this.curTsCmd == null || this.curTsCmd.getAgg() == null) ? Aggregation.NONE : this.curTsCmd.getAgg(); + } + private void sendLatestWsMsg(EntityId entityId, String sessionId, TelemetrySubscriptionUpdate subscriptionUpdate, EntityKeyType keyType) { Map latestUpdate = new HashMap<>(); subscriptionUpdate.getData().forEach((k, v) -> { diff --git a/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java b/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java index 737ef8d8f6..1b17ec1816 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java @@ -115,7 +115,6 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD private PreparedStatement[] fetchStmtsAsc; private PreparedStatement[] fetchStmtsDesc; private PreparedStatement deleteStmt; - private PreparedStatement deletePartitionStmt; private final Lock stmtCreationLock = new ReentrantLock(); private boolean isInstall() { @@ -584,51 +583,6 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD return deleteStmt; } - private void deletePartitionAsync(TenantId tenantId, final QueryCursor cursor, final SimpleListenableFuture resultFuture) { - if (!cursor.hasNextPartition()) { - resultFuture.set(null); - } else { - PreparedStatement proto = getDeletePartitionStmt(); - BoundStatementBuilder stmtBuilder = new BoundStatementBuilder(proto.bind()); - stmtBuilder.setString(0, cursor.getEntityType()); - stmtBuilder.setUuid(1, cursor.getEntityId()); - stmtBuilder.setLong(2, cursor.getNextPartition()); - stmtBuilder.setString(3, cursor.getKey()); - - BoundStatement stmt = stmtBuilder.build(); - - Futures.addCallback(executeAsyncWrite(tenantId, stmt), new FutureCallback() { - @Override - public void onSuccess(@Nullable AsyncResultSet result) { - deletePartitionAsync(tenantId, cursor, resultFuture); - } - - @Override - public void onFailure(Throwable t) { - log.error("[{}][{}] Failed to delete data for query {}-{}", stmt, t); - } - }, readResultsProcessingExecutor); - } - } - - private PreparedStatement getDeletePartitionStmt() { - if (deletePartitionStmt == null) { - stmtCreationLock.lock(); - try { - if (deletePartitionStmt == null) { - deletePartitionStmt = prepare("DELETE FROM " + ModelConstants.TS_KV_PARTITIONS_CF + - " WHERE " + ModelConstants.ENTITY_TYPE_COLUMN + EQUALS_PARAM - + "AND " + ModelConstants.ENTITY_ID_COLUMN + EQUALS_PARAM - + "AND " + ModelConstants.PARTITION_COLUMN + EQUALS_PARAM - + "AND " + ModelConstants.KEY_COLUMN + EQUALS_PARAM); - } - } finally { - stmtCreationLock.unlock(); - } - } - return deletePartitionStmt; - } - private PreparedStatement getSaveStmt(DataType dataType) { if (saveStmts == null) { stmtCreationLock.lock();