Browse Source

added sequence number for attributes

pull/10975/head
YevhenBondarenko 2 years ago
parent
commit
b666f9499a
  1. 4
      application/src/main/java/org/thingsboard/server/controller/DeviceController.java
  2. 2
      application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionService.java
  3. 2
      common/dao-api/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesService.java
  4. 8
      common/data/src/main/java/org/thingsboard/server/common/data/util/CollectionsUtil.java
  5. 128
      dao/src/main/java/org/thingsboard/server/dao/AbstractSequenceInsertRepository.java
  6. 47
      dao/src/main/java/org/thingsboard/server/dao/sql/TbSqlBlockingQueue.java
  7. 1
      dao/src/main/java/org/thingsboard/server/dao/sql/TbSqlBlockingQueueParams.java
  8. 15
      dao/src/main/java/org/thingsboard/server/dao/sql/TbSqlBlockingQueueWrapper.java
  9. 8
      dao/src/main/java/org/thingsboard/server/dao/sql/TbSqlQueue.java
  10. 6
      dao/src/main/java/org/thingsboard/server/dao/sql/TbSqlQueueElement.java
  11. 226
      dao/src/main/java/org/thingsboard/server/dao/sql/attributes/AttributeKvInsertRepository.java
  12. 5
      dao/src/main/java/org/thingsboard/server/dao/sql/attributes/JpaAttributeDao.java
  13. 27
      dao/src/main/java/org/thingsboard/server/dao/sql/attributes/SqlAttributesInsertRepository.java
  14. 2
      dao/src/main/java/org/thingsboard/server/dao/sql/edge/JpaBaseEdgeEventDao.java
  15. 2
      dao/src/main/java/org/thingsboard/server/dao/sql/event/JpaBaseEventDao.java
  16. 2
      dao/src/main/java/org/thingsboard/server/dao/sqlts/AbstractChunkedAggregationTimeseriesDao.java
  17. 36
      dao/src/main/java/org/thingsboard/server/dao/sqlts/SqlTimeseriesLatestDao.java
  18. 224
      dao/src/main/java/org/thingsboard/server/dao/sqlts/insert/latest/sql/SqlLatestInsertTsRepository.java
  19. 8
      dao/src/main/java/org/thingsboard/server/dao/sqlts/timescale/TimescaleTimeseriesDao.java
  20. 4
      dao/src/main/java/org/thingsboard/server/dao/timeseries/BaseTimeseriesService.java
  21. 2
      dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesLatestDao.java
  22. 2
      dao/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesLatestDao.java
  23. 3
      dao/src/main/resources/sql/schema-entities.sql

4
application/src/main/java/org/thingsboard/server/controller/DeviceController.java

@ -20,6 +20,7 @@ import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.ListenableFuture;
import com.google.common.util.concurrent.MoreExecutors; import com.google.common.util.concurrent.MoreExecutors;
import io.swagger.v3.oas.annotations.Parameter; import io.swagger.v3.oas.annotations.Parameter;
import io.swagger.v3.oas.annotations.media.ArraySchema;
import io.swagger.v3.oas.annotations.media.Content; import io.swagger.v3.oas.annotations.media.Content;
import io.swagger.v3.oas.annotations.media.Schema; import io.swagger.v3.oas.annotations.media.Schema;
import io.swagger.v3.oas.annotations.responses.ApiResponse; import io.swagger.v3.oas.annotations.responses.ApiResponse;
@ -30,6 +31,7 @@ import org.springframework.http.HttpStatus;
import org.springframework.http.MediaType; import org.springframework.http.MediaType;
import org.springframework.http.ResponseEntity; import org.springframework.http.ResponseEntity;
import org.springframework.security.access.prepost.PreAuthorize; import org.springframework.security.access.prepost.PreAuthorize;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.PathVariable; import org.springframework.web.bind.annotation.PathVariable;
import org.springframework.web.bind.annotation.PostMapping; import org.springframework.web.bind.annotation.PostMapping;
import org.springframework.web.bind.annotation.RequestBody; import org.springframework.web.bind.annotation.RequestBody;
@ -488,7 +490,7 @@ public class DeviceController extends BaseController {
@RequestMapping(value = "/devices", params = {"deviceIds"}, method = RequestMethod.GET) @RequestMapping(value = "/devices", params = {"deviceIds"}, method = RequestMethod.GET)
@ResponseBody @ResponseBody
public List<Device> getDevicesByIds( public List<Device> getDevicesByIds(
@Parameter(description = "A list of devices ids, separated by comma ','") @Parameter(description = "A list of devices ids, separated by comma ','", array = @ArraySchema(schema = @Schema(type = "string")))
@RequestParam("deviceIds") String[] strDeviceIds) throws ThingsboardException, ExecutionException, InterruptedException { @RequestParam("deviceIds") String[] strDeviceIds) throws ThingsboardException, ExecutionException, InterruptedException {
checkArrayParameter("deviceIds", strDeviceIds); checkArrayParameter("deviceIds", strDeviceIds);
SecurityUser user = getCurrentUser(); SecurityUser user = getCurrentUser();

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

@ -280,7 +280,7 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer
@Override @Override
public void saveLatestAndNotifyInternal(TenantId tenantId, EntityId entityId, List<TsKvEntry> ts, FutureCallback<Void> callback) { public void saveLatestAndNotifyInternal(TenantId tenantId, EntityId entityId, List<TsKvEntry> ts, FutureCallback<Void> callback) {
ListenableFuture<List<Void>> saveFuture = tsService.saveLatest(tenantId, entityId, ts); ListenableFuture<List<Long>> saveFuture = tsService.saveLatest(tenantId, entityId, ts);
addVoidCallback(saveFuture, callback); addVoidCallback(saveFuture, callback);
addWsCallback(saveFuture, success -> onTimeSeriesUpdate(tenantId, entityId, ts)); addWsCallback(saveFuture, success -> onTimeSeriesUpdate(tenantId, entityId, ts));
} }

2
common/dao-api/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesService.java

@ -50,7 +50,7 @@ public interface TimeseriesService {
ListenableFuture<Integer> saveWithoutLatest(TenantId tenantId, EntityId entityId, List<TsKvEntry> tsKvEntry, long ttl); ListenableFuture<Integer> saveWithoutLatest(TenantId tenantId, EntityId entityId, List<TsKvEntry> tsKvEntry, long ttl);
ListenableFuture<List<Void>> saveLatest(TenantId tenantId, EntityId entityId, List<TsKvEntry> tsKvEntry); ListenableFuture<List<Long>> saveLatest(TenantId tenantId, EntityId entityId, List<TsKvEntry> tsKvEntry);
ListenableFuture<List<TsKvLatestRemovingResult>> remove(TenantId tenantId, EntityId entityId, List<DeleteTsKvQuery> queries); ListenableFuture<List<TsKvLatestRemovingResult>> remove(TenantId tenantId, EntityId entityId, List<DeleteTsKvQuery> queries);

8
common/data/src/main/java/org/thingsboard/server/common/data/util/CollectionsUtil.java

@ -17,6 +17,7 @@ package org.thingsboard.server.common.data.util;
import java.util.Collection; import java.util.Collection;
import java.util.HashMap; import java.util.HashMap;
import java.util.List;
import java.util.Map; import java.util.Map;
import java.util.Set; import java.util.Set;
import java.util.stream.Collectors; import java.util.stream.Collectors;
@ -37,6 +38,13 @@ public class CollectionsUtil {
return b.stream().filter(p -> !a.contains(p)).collect(Collectors.toSet()); return b.stream().filter(p -> !a.contains(p)).collect(Collectors.toSet());
} }
/**
* Returns new list with elements that are present in list B(new) but absent in list A(old).
*/
public static <T> List<T> diffLists(List<T> a, List<T> b) {
return b.stream().filter(p -> !a.contains(p)).collect(Collectors.toList());
}
public static <T> boolean contains(Collection<T> collection, T element) { public static <T> boolean contains(Collection<T> collection, T element) {
return isNotEmpty(collection) && collection.contains(element); return isNotEmpty(collection) && collection.contains(element);
} }

128
dao/src/main/java/org/thingsboard/server/dao/AbstractSequenceInsertRepository.java

@ -0,0 +1,128 @@
/**
* Copyright © 2016-2024 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.dao;
import org.springframework.jdbc.core.BatchPreparedStatementSetter;
import org.springframework.jdbc.core.PreparedStatementCreator;
import org.springframework.jdbc.core.SqlProvider;
import org.springframework.jdbc.support.GeneratedKeyHolder;
import org.springframework.jdbc.support.KeyHolder;
import org.thingsboard.server.dao.sqlts.insert.AbstractInsertRepository;
import java.sql.Connection;
import java.sql.PreparedStatement;
import java.sql.SQLException;
import java.util.ArrayList;
import java.util.List;
import java.util.Map;
public abstract class AbstractSequenceInsertRepository<T> extends AbstractInsertRepository {
public static final String SEQ_NUMBER = "seq_number";
public List<Long> saveOrUpdate(List<T> entities) {
return transactionTemplate.execute(status -> {
List<Long> seqNumbers = new ArrayList<>(entities.size());
KeyHolder keyHolder = new GeneratedKeyHolder();
int[] updateResult = onBatchUpdate(entities, keyHolder);
List<Map<String, Object>> seqNumbersList = keyHolder.getKeyList();
int notUpdatedCount = entities.size() - seqNumbersList.size();
List<Integer> toInsertIndexes = new ArrayList<>(notUpdatedCount);
List<T> insertEntities = new ArrayList<>(notUpdatedCount);
int keyHolderIndex = 0;
for (int i = 0; i < updateResult.length; i++) {
if (updateResult[i] == 0) {
insertEntities.add(entities.get(i));
seqNumbers.add(0L);
toInsertIndexes.add(i);
} else {
seqNumbers.add((Long) seqNumbersList.get(keyHolderIndex).get(SEQ_NUMBER));
keyHolderIndex++;
}
}
if (insertEntities.isEmpty()) {
return seqNumbers;
}
onInsertOrUpdate(insertEntities, keyHolder);
seqNumbersList = keyHolder.getKeyList();
for (int i = 0; i < seqNumbersList.size(); i++) {
seqNumbers.set(toInsertIndexes.get(i), (Long) seqNumbersList.get(i).get(SEQ_NUMBER));
}
return seqNumbers;
});
}
private int[] onBatchUpdate(List<T> entities, KeyHolder keyHolder) {
return jdbcTemplate.batchUpdate(new SequencePreparedStatementCreator(getBatchUpdateQuery()), new BatchPreparedStatementSetter() {
@Override
public void setValues(PreparedStatement ps, int i) throws SQLException {
setOnBatchUpdateValues(ps, i, entities);
}
@Override
public int getBatchSize() {
return entities.size();
}
}, keyHolder);
}
private void onInsertOrUpdate(List<T> insertEntities, KeyHolder keyHolder) {
jdbcTemplate.batchUpdate(new SequencePreparedStatementCreator(getInsertOrUpdateQuery()), new BatchPreparedStatementSetter() {
@Override
public void setValues(PreparedStatement ps, int i) throws SQLException {
setOnInsertOrUpdateValues(ps, i, insertEntities);
}
@Override
public int getBatchSize() {
return insertEntities.size();
}
}, keyHolder);
}
protected abstract void setOnBatchUpdateValues(PreparedStatement ps, int i, List<T> entities) throws SQLException;
protected abstract void setOnInsertOrUpdateValues(PreparedStatement ps, int i, List<T> entities) throws SQLException;
protected abstract String getBatchUpdateQuery();
protected abstract String getInsertOrUpdateQuery();
private record SequencePreparedStatementCreator(String sql) implements PreparedStatementCreator, SqlProvider {
private static final String[] COLUMNS = {SEQ_NUMBER};
@Override
public PreparedStatement createPreparedStatement(Connection con) throws SQLException {
return con.prepareStatement(sql, COLUMNS);
}
@Override
public String getSql() {
return this.sql;
}
}
}

47
dao/src/main/java/org/thingsboard/server/dao/sql/TbSqlBlockingQueue.java

@ -19,6 +19,7 @@ import com.google.common.util.concurrent.ListenableFuture;
import com.google.common.util.concurrent.SettableFuture; import com.google.common.util.concurrent.SettableFuture;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.thingsboard.common.util.ThingsBoardThreadFactory; import org.thingsboard.common.util.ThingsBoardThreadFactory;
import org.thingsboard.server.common.data.util.CollectionsUtil;
import org.thingsboard.server.common.stats.MessagesStats; import org.thingsboard.server.common.stats.MessagesStats;
import java.util.ArrayList; import java.util.ArrayList;
@ -29,14 +30,13 @@ import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors; import java.util.concurrent.Executors;
import java.util.concurrent.LinkedBlockingQueue; import java.util.concurrent.LinkedBlockingQueue;
import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeUnit;
import java.util.function.Consumer; import java.util.function.Function;
import java.util.stream.Collectors; import java.util.stream.Collectors;
import java.util.stream.Stream;
@Slf4j @Slf4j
public class TbSqlBlockingQueue<E> implements TbSqlQueue<E> { public class TbSqlBlockingQueue<E, R> implements TbSqlQueue<E, R> {
private final BlockingQueue<TbSqlQueueElement<E>> queue = new LinkedBlockingQueue<>(); private final BlockingQueue<TbSqlQueueElement<E, R>> queue = new LinkedBlockingQueue<>();
private final TbSqlBlockingQueueParams params; private final TbSqlBlockingQueueParams params;
private ExecutorService executor; private ExecutorService executor;
@ -48,17 +48,17 @@ public class TbSqlBlockingQueue<E> implements TbSqlQueue<E> {
} }
@Override @Override
public void init(ScheduledLogExecutorComponent logExecutor, Consumer<List<E>> saveFunction, Comparator<E> batchUpdateComparator, int index) { public void init(ScheduledLogExecutorComponent logExecutor, Function<List<E>, List<R>> saveFunction, Comparator<E> batchUpdateComparator, Function<List<TbSqlQueueElement<E, R>>, List<TbSqlQueueElement<E, R>>> filter, int index) {
executor = Executors.newSingleThreadExecutor(ThingsBoardThreadFactory.forName("sql-queue-" + index + "-" + params.getLogName().toLowerCase())); executor = Executors.newSingleThreadExecutor(ThingsBoardThreadFactory.forName("sql-queue-" + index + "-" + params.getLogName().toLowerCase()));
executor.submit(() -> { executor.submit(() -> {
String logName = params.getLogName(); String logName = params.getLogName();
int batchSize = params.getBatchSize(); int batchSize = params.getBatchSize();
long maxDelay = params.getMaxDelay(); long maxDelay = params.getMaxDelay();
final List<TbSqlQueueElement<E>> entities = new ArrayList<>(batchSize); final List<TbSqlQueueElement<E, R>> entities = new ArrayList<>(batchSize);
while (!Thread.interrupted()) { while (!Thread.interrupted()) {
try { try {
long currentTs = System.currentTimeMillis(); long currentTs = System.currentTimeMillis();
TbSqlQueueElement<E> attr = queue.poll(maxDelay, TimeUnit.MILLISECONDS); TbSqlQueueElement<E, R> attr = queue.poll(maxDelay, TimeUnit.MILLISECONDS);
if (attr == null) { if (attr == null) {
continue; continue;
} else { } else {
@ -70,12 +70,27 @@ public class TbSqlBlockingQueue<E> implements TbSqlQueue<E> {
log.debug("[{}] Going to save {} entities", logName, entities.size()); log.debug("[{}] Going to save {} entities", logName, entities.size());
log.trace("[{}] Going to save entities: {}", logName, entities); log.trace("[{}] Going to save entities: {}", logName, entities);
} }
Stream<E> entitiesStream = entities.stream().map(TbSqlQueueElement::getEntity);
saveFunction.accept( List<TbSqlQueueElement<E, R>> entitiesToSave = filter.apply(entities);
(params.isBatchSortEnabled() ? entitiesStream.sorted(batchUpdateComparator) : entitiesStream)
.collect(Collectors.toList()) if (params.isBatchSortEnabled()) {
); entitiesToSave = entitiesToSave.stream().sorted((o1, o2) -> batchUpdateComparator.compare(o1.getEntity(), o2.getEntity())).toList();
entities.forEach(v -> v.getFuture().set(null)); }
List<R> result = saveFunction.apply(entitiesToSave.stream().map(TbSqlQueueElement::getEntity).collect(Collectors.toList()));
if (params.isWithResponse()) {
for (int i = 0; i < entitiesToSave.size(); i++) {
entitiesToSave.get(i).getFuture().set(result.get(i));
}
if (entities.size() > entitiesToSave.size()) {
CollectionsUtil.diffLists(entitiesToSave, entities).forEach(v -> v.getFuture().set(null));
}
} else {
entities.forEach(v -> v.getFuture().set(null));
}
stats.incrementSuccessful(entities.size()); stats.incrementSuccessful(entities.size());
if (!fullPack) { if (!fullPack) {
long remainingDelay = maxDelay - (System.currentTimeMillis() - currentTs); long remainingDelay = maxDelay - (System.currentTimeMillis() - currentTs);
@ -104,7 +119,7 @@ public class TbSqlBlockingQueue<E> implements TbSqlQueue<E> {
}); });
logExecutor.scheduleAtFixedRate(() -> { logExecutor.scheduleAtFixedRate(() -> {
if (queue.size() > 0 || stats.getTotal() > 0 || stats.getSuccessful() > 0 || stats.getFailed() > 0) { if (!queue.isEmpty() || stats.getTotal() > 0 || stats.getSuccessful() > 0 || stats.getFailed() > 0) {
log.info("Queue-{} [{}] queueSize [{}] totalAdded [{}] totalSaved [{}] totalFailed [{}]", index, log.info("Queue-{} [{}] queueSize [{}] totalAdded [{}] totalSaved [{}] totalFailed [{}]", index,
params.getLogName(), queue.size(), stats.getTotal(), stats.getSuccessful(), stats.getFailed()); params.getLogName(), queue.size(), stats.getTotal(), stats.getSuccessful(), stats.getFailed());
stats.reset(); stats.reset();
@ -120,8 +135,8 @@ public class TbSqlBlockingQueue<E> implements TbSqlQueue<E> {
} }
@Override @Override
public ListenableFuture<Void> add(E element) { public ListenableFuture<R> add(E element) {
SettableFuture<Void> future = SettableFuture.create(); SettableFuture<R> future = SettableFuture.create();
queue.add(new TbSqlQueueElement<>(future, element)); queue.add(new TbSqlQueueElement<>(future, element));
stats.incrementTotal(); stats.incrementTotal();
return future; return future;

1
dao/src/main/java/org/thingsboard/server/dao/sql/TbSqlBlockingQueueParams.java

@ -30,4 +30,5 @@ public class TbSqlBlockingQueueParams {
private final long statsPrintIntervalMs; private final long statsPrintIntervalMs;
private final String statsNamePrefix; private final String statsNamePrefix;
private final boolean batchSortEnabled; private final boolean batchSortEnabled;
private final boolean withResponse;
} }

15
dao/src/main/java/org/thingsboard/server/dao/sql/TbSqlBlockingQueueWrapper.java

@ -29,10 +29,9 @@ import java.util.function.Function;
@Slf4j @Slf4j
@Data @Data
public class TbSqlBlockingQueueWrapper<E> { public class TbSqlBlockingQueueWrapper<E, R> {
private final CopyOnWriteArrayList<TbSqlBlockingQueue<E>> queues = new CopyOnWriteArrayList<>(); private final CopyOnWriteArrayList<TbSqlBlockingQueue<E, R>> queues = new CopyOnWriteArrayList<>();
private final TbSqlBlockingQueueParams params; private final TbSqlBlockingQueueParams params;
private ScheduledLogExecutorComponent logExecutor;
private final Function<E, Integer> hashCodeFunction; private final Function<E, Integer> hashCodeFunction;
private final int maxThreads; private final int maxThreads;
private final StatsFactory statsFactory; private final StatsFactory statsFactory;
@ -46,15 +45,19 @@ public class TbSqlBlockingQueueWrapper<E> {
* NOTE: you must use all of primary key parts in your comparator * NOTE: you must use all of primary key parts in your comparator
*/ */
public void init(ScheduledLogExecutorComponent logExecutor, Consumer<List<E>> saveFunction, Comparator<E> batchUpdateComparator) { public void init(ScheduledLogExecutorComponent logExecutor, Consumer<List<E>> saveFunction, Comparator<E> batchUpdateComparator) {
init(logExecutor, l -> { saveFunction.accept(l); return null; }, batchUpdateComparator, l -> l);
}
public void init(ScheduledLogExecutorComponent logExecutor, Function<List<E>, List<R>> saveFunction, Comparator<E> batchUpdateComparator, Function<List<TbSqlQueueElement<E, R>>, List<TbSqlQueueElement<E, R>>> filter) {
for (int i = 0; i < maxThreads; i++) { for (int i = 0; i < maxThreads; i++) {
MessagesStats stats = statsFactory.createMessagesStats(params.getStatsNamePrefix() + ".queue." + i); MessagesStats stats = statsFactory.createMessagesStats(params.getStatsNamePrefix() + ".queue." + i);
TbSqlBlockingQueue<E> queue = new TbSqlBlockingQueue<>(params, stats); TbSqlBlockingQueue<E, R> queue = new TbSqlBlockingQueue<>(params, stats);
queues.add(queue); queues.add(queue);
queue.init(logExecutor, saveFunction, batchUpdateComparator, i); queue.init(logExecutor, saveFunction, batchUpdateComparator, filter, i);
} }
} }
public ListenableFuture<Void> add(E element) { public ListenableFuture<R> add(E element) {
int queueIndex = element != null ? (hashCodeFunction.apply(element) & 0x7FFFFFFF) % maxThreads : 0; int queueIndex = element != null ? (hashCodeFunction.apply(element) & 0x7FFFFFFF) % maxThreads : 0;
return queues.get(queueIndex).add(element); return queues.get(queueIndex).add(element);
} }

8
dao/src/main/java/org/thingsboard/server/dao/sql/TbSqlQueue.java

@ -19,13 +19,13 @@ import com.google.common.util.concurrent.ListenableFuture;
import java.util.Comparator; import java.util.Comparator;
import java.util.List; import java.util.List;
import java.util.function.Consumer; import java.util.function.Function;
public interface TbSqlQueue<E> { public interface TbSqlQueue<E, R> {
void init(ScheduledLogExecutorComponent logExecutor, Consumer<List<E>> saveFunction, Comparator<E> batchUpdateComparator, int queueIndex); void init(ScheduledLogExecutorComponent logExecutor, Function<List<E>, List<R>> saveFunction, Comparator<E> batchUpdateComparator, Function<List<TbSqlQueueElement<E, R>>, List<TbSqlQueueElement<E, R>>> filter, int queueIndex);
void destroy(); void destroy();
ListenableFuture<Void> add(E element); ListenableFuture<R> add(E element);
} }

6
dao/src/main/java/org/thingsboard/server/dao/sql/TbSqlQueueElement.java

@ -20,13 +20,13 @@ import lombok.Getter;
import lombok.ToString; import lombok.ToString;
@ToString(exclude = "future") @ToString(exclude = "future")
public final class TbSqlQueueElement<E> { public final class TbSqlQueueElement<E, R> {
@Getter @Getter
private final SettableFuture<Void> future; private final SettableFuture<R> future;
@Getter @Getter
private final E entity; private final E entity;
public TbSqlQueueElement(SettableFuture<Void> future, E entity) { public TbSqlQueueElement(SettableFuture<R> future, E entity) {
this.future = future; this.future = future;
this.entity = entity; this.entity = entity;
} }

226
dao/src/main/java/org/thingsboard/server/dao/sql/attributes/AttributeKvInsertRepository.java

@ -15,162 +15,110 @@
*/ */
package org.thingsboard.server.dao.sql.attributes; package org.thingsboard.server.dao.sql.attributes;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.jdbc.core.BatchPreparedStatementSetter;
import org.springframework.jdbc.core.JdbcTemplate;
import org.springframework.stereotype.Repository; import org.springframework.stereotype.Repository;
import org.springframework.transaction.TransactionStatus; import org.springframework.transaction.annotation.Transactional;
import org.springframework.transaction.support.TransactionCallbackWithoutResult; import org.thingsboard.server.dao.AbstractSequenceInsertRepository;
import org.springframework.transaction.support.TransactionTemplate;
import org.thingsboard.server.dao.model.sql.AttributeKvEntity; import org.thingsboard.server.dao.model.sql.AttributeKvEntity;
import org.thingsboard.server.dao.util.SqlDao; import org.thingsboard.server.dao.util.SqlDao;
import java.sql.PreparedStatement; import java.sql.PreparedStatement;
import java.sql.SQLException; import java.sql.SQLException;
import java.sql.Types; import java.sql.Types;
import java.util.ArrayList;
import java.util.List; import java.util.List;
import java.util.regex.Pattern;
@Repository @Repository
@Slf4j @Transactional
@SqlDao @SqlDao
public abstract class AttributeKvInsertRepository { public class AttributeKvInsertRepository extends AbstractSequenceInsertRepository<AttributeKvEntity> {
private static final ThreadLocal<Pattern> PATTERN_THREAD_LOCAL = ThreadLocal.withInitial(() -> Pattern.compile(String.valueOf(Character.MIN_VALUE))); private static final String BATCH_UPDATE = "UPDATE attribute_kv SET str_v = ?, long_v = ?, dbl_v = ?, bool_v = ?, json_v = cast(? AS json), last_update_ts = ?, seq_number = nextval('attribute_kv_latest_seq') " +
private static final String EMPTY_STR = ""; "WHERE entity_id = ? and attribute_type =? and attribute_key = ? RETURNING seq_number;";
private static final String BATCH_UPDATE = "UPDATE attribute_kv SET str_v = ?, long_v = ?, dbl_v = ?, bool_v = ?, json_v = cast(? AS json), last_update_ts = ? " +
"WHERE entity_id = ? and attribute_type =? and attribute_key = ?;";
private static final String INSERT_OR_UPDATE = private static final String INSERT_OR_UPDATE =
"INSERT INTO attribute_kv (entity_id, attribute_type, attribute_key, str_v, long_v, dbl_v, bool_v, json_v, last_update_ts) " + "INSERT INTO attribute_kv (entity_id, attribute_type, attribute_key, str_v, long_v, dbl_v, bool_v, json_v, last_update_ts, seq_number) " +
"VALUES(?, ?, ?, ?, ?, ?, ?, cast(? AS json), ?) " + "VALUES(?, ?, ?, ?, ?, ?, ?, cast(? AS json), ?, nextval('attribute_kv_latest_seq')) " +
"ON CONFLICT (entity_id, attribute_type, attribute_key) " + "ON CONFLICT (entity_id, attribute_type, attribute_key) " +
"DO UPDATE SET str_v = ?, long_v = ?, dbl_v = ?, bool_v = ?, json_v = cast(? AS json), last_update_ts = ?;"; "DO UPDATE SET str_v = ?, long_v = ?, dbl_v = ?, bool_v = ?, json_v = cast(? AS json), last_update_ts = ?, seq_number = nextval('attribute_kv_latest_seq') RETURNING seq_number;";
@Autowired @Override
protected JdbcTemplate jdbcTemplate; protected void setOnBatchUpdateValues(PreparedStatement ps, int i, List<AttributeKvEntity> entities) throws SQLException {
AttributeKvEntity kvEntity = entities.get(i);
@Autowired ps.setString(1, replaceNullChars(kvEntity.getStrValue()));
private TransactionTemplate transactionTemplate;
if (kvEntity.getLongValue() != null) {
@Value("${sql.remove_null_chars:true}") ps.setLong(2, kvEntity.getLongValue());
private boolean removeNullChars; } else {
ps.setNull(2, Types.BIGINT);
public void saveOrUpdate(List<AttributeKvEntity> entities) { }
transactionTemplate.execute(new TransactionCallbackWithoutResult() {
@Override if (kvEntity.getDoubleValue() != null) {
protected void doInTransactionWithoutResult(TransactionStatus status) { ps.setDouble(3, kvEntity.getDoubleValue());
int[] result = jdbcTemplate.batchUpdate(BATCH_UPDATE, new BatchPreparedStatementSetter() { } else {
@Override ps.setNull(3, Types.DOUBLE);
public void setValues(PreparedStatement ps, int i) throws SQLException { }
AttributeKvEntity kvEntity = entities.get(i);
ps.setString(1, replaceNullChars(kvEntity.getStrValue())); if (kvEntity.getBooleanValue() != null) {
ps.setBoolean(4, kvEntity.getBooleanValue());
if (kvEntity.getLongValue() != null) { } else {
ps.setLong(2, kvEntity.getLongValue()); ps.setNull(4, Types.BOOLEAN);
} else { }
ps.setNull(2, Types.BIGINT);
} ps.setString(5, replaceNullChars(kvEntity.getJsonValue()));
if (kvEntity.getDoubleValue() != null) { ps.setLong(6, kvEntity.getLastUpdateTs());
ps.setDouble(3, kvEntity.getDoubleValue()); ps.setObject(7, kvEntity.getId().getEntityId());
} else { ps.setInt(8, kvEntity.getId().getAttributeType());
ps.setNull(3, Types.DOUBLE); ps.setInt(9, kvEntity.getId().getAttributeKey());
}
if (kvEntity.getBooleanValue() != null) {
ps.setBoolean(4, kvEntity.getBooleanValue());
} else {
ps.setNull(4, Types.BOOLEAN);
}
ps.setString(5, replaceNullChars(kvEntity.getJsonValue()));
ps.setLong(6, kvEntity.getLastUpdateTs());
ps.setObject(7, kvEntity.getId().getEntityId());
ps.setInt(8, kvEntity.getId().getAttributeType());
ps.setInt(9, kvEntity.getId().getAttributeKey());
}
@Override
public int getBatchSize() {
return entities.size();
}
});
int updatedCount = 0;
for (int i = 0; i < result.length; i++) {
if (result[i] == 0) {
updatedCount++;
}
}
List<AttributeKvEntity> insertEntities = new ArrayList<>(updatedCount);
for (int i = 0; i < result.length; i++) {
if (result[i] == 0) {
insertEntities.add(entities.get(i));
}
}
jdbcTemplate.batchUpdate(INSERT_OR_UPDATE, new BatchPreparedStatementSetter() {
@Override
public void setValues(PreparedStatement ps, int i) throws SQLException {
AttributeKvEntity kvEntity = insertEntities.get(i);
ps.setObject(1, kvEntity.getId().getEntityId());
ps.setInt(2, kvEntity.getId().getAttributeType());
ps.setInt(3, kvEntity.getId().getAttributeKey());
ps.setString(4, replaceNullChars(kvEntity.getStrValue()));
ps.setString(10, replaceNullChars(kvEntity.getStrValue()));
if (kvEntity.getLongValue() != null) {
ps.setLong(5, kvEntity.getLongValue());
ps.setLong(11, kvEntity.getLongValue());
} else {
ps.setNull(5, Types.BIGINT);
ps.setNull(11, Types.BIGINT);
}
if (kvEntity.getDoubleValue() != null) {
ps.setDouble(6, kvEntity.getDoubleValue());
ps.setDouble(12, kvEntity.getDoubleValue());
} else {
ps.setNull(6, Types.DOUBLE);
ps.setNull(12, Types.DOUBLE);
}
if (kvEntity.getBooleanValue() != null) {
ps.setBoolean(7, kvEntity.getBooleanValue());
ps.setBoolean(13, kvEntity.getBooleanValue());
} else {
ps.setNull(7, Types.BOOLEAN);
ps.setNull(13, Types.BOOLEAN);
}
ps.setString(8, replaceNullChars(kvEntity.getJsonValue()));
ps.setString(14, replaceNullChars(kvEntity.getJsonValue()));
ps.setLong(9, kvEntity.getLastUpdateTs());
ps.setLong(15, kvEntity.getLastUpdateTs());
}
@Override
public int getBatchSize() {
return insertEntities.size();
}
});
}
});
} }
private String replaceNullChars(String strValue) { @Override
if (removeNullChars && strValue != null) { protected void setOnInsertOrUpdateValues(PreparedStatement ps, int i, List<AttributeKvEntity> insertEntities) throws SQLException {
return PATTERN_THREAD_LOCAL.get().matcher(strValue).replaceAll(EMPTY_STR); AttributeKvEntity kvEntity = insertEntities.get(i);
ps.setObject(1, kvEntity.getId().getEntityId());
ps.setInt(2, kvEntity.getId().getAttributeType());
ps.setInt(3, kvEntity.getId().getAttributeKey());
ps.setString(4, replaceNullChars(kvEntity.getStrValue()));
ps.setString(10, replaceNullChars(kvEntity.getStrValue()));
if (kvEntity.getLongValue() != null) {
ps.setLong(5, kvEntity.getLongValue());
ps.setLong(11, kvEntity.getLongValue());
} else {
ps.setNull(5, Types.BIGINT);
ps.setNull(11, Types.BIGINT);
} }
return strValue;
if (kvEntity.getDoubleValue() != null) {
ps.setDouble(6, kvEntity.getDoubleValue());
ps.setDouble(12, kvEntity.getDoubleValue());
} else {
ps.setNull(6, Types.DOUBLE);
ps.setNull(12, Types.DOUBLE);
}
if (kvEntity.getBooleanValue() != null) {
ps.setBoolean(7, kvEntity.getBooleanValue());
ps.setBoolean(13, kvEntity.getBooleanValue());
} else {
ps.setNull(7, Types.BOOLEAN);
ps.setNull(13, Types.BOOLEAN);
}
ps.setString(8, replaceNullChars(kvEntity.getJsonValue()));
ps.setString(14, replaceNullChars(kvEntity.getJsonValue()));
ps.setLong(9, kvEntity.getLastUpdateTs());
ps.setLong(15, kvEntity.getLastUpdateTs());
}
@Override
protected String getBatchUpdateQuery() {
return BATCH_UPDATE;
}
@Override
protected String getInsertOrUpdateQuery() {
return INSERT_OR_UPDATE;
} }
} }

5
dao/src/main/java/org/thingsboard/server/dao/sql/attributes/JpaAttributeDao.java

@ -88,7 +88,7 @@ public class JpaAttributeDao extends JpaAbstractDaoListeningExecutorService impl
@Value("${sql.batch_sort:true}") @Value("${sql.batch_sort:true}")
private boolean batchSortEnabled; private boolean batchSortEnabled;
private TbSqlBlockingQueueWrapper<AttributeKvEntity> queue; private TbSqlBlockingQueueWrapper<AttributeKvEntity, Long> queue;
@PostConstruct @PostConstruct
private void init() { private void init() {
@ -99,6 +99,7 @@ public class JpaAttributeDao extends JpaAbstractDaoListeningExecutorService impl
.statsPrintIntervalMs(statsPrintIntervalMs) .statsPrintIntervalMs(statsPrintIntervalMs)
.statsNamePrefix("attributes") .statsNamePrefix("attributes")
.batchSortEnabled(batchSortEnabled) .batchSortEnabled(batchSortEnabled)
.withResponse(true)
.build(); .build();
Function<AttributeKvEntity, Integer> hashcodeFunction = entity -> entity.getId().getEntityId().hashCode(); Function<AttributeKvEntity, Integer> hashcodeFunction = entity -> entity.getId().getEntityId().hashCode();
@ -106,7 +107,7 @@ public class JpaAttributeDao extends JpaAbstractDaoListeningExecutorService impl
queue.init(logExecutor, v -> attributeKvInsertRepository.saveOrUpdate(v), queue.init(logExecutor, v -> attributeKvInsertRepository.saveOrUpdate(v),
Comparator.comparing((AttributeKvEntity attributeKvEntity) -> attributeKvEntity.getId().getEntityId()) Comparator.comparing((AttributeKvEntity attributeKvEntity) -> attributeKvEntity.getId().getEntityId())
.thenComparing(attributeKvEntity -> attributeKvEntity.getId().getAttributeType()) .thenComparing(attributeKvEntity -> attributeKvEntity.getId().getAttributeType())
.thenComparing(attributeKvEntity -> attributeKvEntity.getId().getAttributeKey()) .thenComparing(attributeKvEntity -> attributeKvEntity.getId().getAttributeKey()), l -> l
); );
} }

27
dao/src/main/java/org/thingsboard/server/dao/sql/attributes/SqlAttributesInsertRepository.java

@ -1,27 +0,0 @@
/**
* Copyright © 2016-2024 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.dao.sql.attributes;
import org.springframework.stereotype.Repository;
import org.springframework.transaction.annotation.Transactional;
import org.thingsboard.server.dao.util.SqlDao;
@Repository
@Transactional
@SqlDao
public class SqlAttributesInsertRepository extends AttributeKvInsertRepository {
}

2
dao/src/main/java/org/thingsboard/server/dao/sql/edge/JpaBaseEdgeEventDao.java

@ -90,7 +90,7 @@ public class JpaBaseEdgeEventDao extends JpaPartitionedAbstractDao<EdgeEventEnti
private static final String TABLE_NAME = ModelConstants.EDGE_EVENT_TABLE_NAME; private static final String TABLE_NAME = ModelConstants.EDGE_EVENT_TABLE_NAME;
private TbSqlBlockingQueueWrapper<EdgeEventEntity> queue; private TbSqlBlockingQueueWrapper<EdgeEventEntity, Void> queue;
@Override @Override
protected Class<EdgeEventEntity> getEntityClass() { protected Class<EdgeEventEntity> getEntityClass() {

2
dao/src/main/java/org/thingsboard/server/dao/sql/event/JpaBaseEventDao.java

@ -110,7 +110,7 @@ public class JpaBaseEventDao implements EventDao {
@Value("${sql.batch_sort:true}") @Value("${sql.batch_sort:true}")
private boolean batchSortEnabled; private boolean batchSortEnabled;
private TbSqlBlockingQueueWrapper<Event> queue; private TbSqlBlockingQueueWrapper<Event, Void> queue;
private final Map<EventType, EventRepository<?, ?>> repositories = new ConcurrentHashMap<>(); private final Map<EventType, EventRepository<?, ?>> repositories = new ConcurrentHashMap<>();

2
dao/src/main/java/org/thingsboard/server/dao/sqlts/AbstractChunkedAggregationTimeseriesDao.java

@ -60,7 +60,7 @@ public abstract class AbstractChunkedAggregationTimeseriesDao extends AbstractSq
@Autowired @Autowired
protected InsertTsRepository<TsKvEntity> insertRepository; protected InsertTsRepository<TsKvEntity> insertRepository;
protected TbSqlBlockingQueueWrapper<TsKvEntity> tsQueue; protected TbSqlBlockingQueueWrapper<TsKvEntity, Void> tsQueue;
@Autowired @Autowired
private StatsFactory statsFactory; private StatsFactory statsFactory;

36
dao/src/main/java/org/thingsboard/server/dao/sqlts/SqlTimeseriesLatestDao.java

@ -46,6 +46,7 @@ import org.thingsboard.server.dao.model.sqlts.latest.TsKvLatestEntity;
import org.thingsboard.server.dao.sql.ScheduledLogExecutorComponent; import org.thingsboard.server.dao.sql.ScheduledLogExecutorComponent;
import org.thingsboard.server.dao.sql.TbSqlBlockingQueueParams; import org.thingsboard.server.dao.sql.TbSqlBlockingQueueParams;
import org.thingsboard.server.dao.sql.TbSqlBlockingQueueWrapper; import org.thingsboard.server.dao.sql.TbSqlBlockingQueueWrapper;
import org.thingsboard.server.dao.sql.TbSqlQueueElement;
import org.thingsboard.server.dao.sqlts.insert.latest.InsertLatestTsRepository; import org.thingsboard.server.dao.sqlts.insert.latest.InsertLatestTsRepository;
import org.thingsboard.server.dao.sqlts.latest.SearchTsKvLatestRepository; import org.thingsboard.server.dao.sqlts.latest.SearchTsKvLatestRepository;
import org.thingsboard.server.dao.sqlts.latest.TsKvLatestRepository; import org.thingsboard.server.dao.sqlts.latest.TsKvLatestRepository;
@ -81,7 +82,7 @@ public class SqlTimeseriesLatestDao extends BaseAbstractSqlTimeseriesDao impleme
@Autowired @Autowired
private InsertLatestTsRepository insertLatestTsRepository; private InsertLatestTsRepository insertLatestTsRepository;
private TbSqlBlockingQueueWrapper<TsKvLatestEntity> tsLatestQueue; private TbSqlBlockingQueueWrapper<TsKvLatestEntity, Long> tsLatestQueue;
@Value("${sql.ts_latest.batch_size:1000}") @Value("${sql.ts_latest.batch_size:1000}")
private int tsLatestBatchSize; private int tsLatestBatchSize;
@ -115,25 +116,26 @@ public class SqlTimeseriesLatestDao extends BaseAbstractSqlTimeseriesDao impleme
.maxDelay(tsLatestMaxDelay) .maxDelay(tsLatestMaxDelay)
.statsPrintIntervalMs(tsLatestStatsPrintIntervalMs) .statsPrintIntervalMs(tsLatestStatsPrintIntervalMs)
.statsNamePrefix("ts.latest") .statsNamePrefix("ts.latest")
.batchSortEnabled(false) .batchSortEnabled(batchSortEnabled)
.withResponse(true)
.build(); .build();
java.util.function.Function<TsKvLatestEntity, Integer> hashcodeFunction = entity -> entity.getEntityId().hashCode(); java.util.function.Function<TsKvLatestEntity, Integer> hashcodeFunction = entity -> entity.getEntityId().hashCode();
tsLatestQueue = new TbSqlBlockingQueueWrapper<>(tsLatestParams, hashcodeFunction, tsLatestBatchThreads, statsFactory); tsLatestQueue = new TbSqlBlockingQueueWrapper<>(tsLatestParams, hashcodeFunction, tsLatestBatchThreads, statsFactory);
tsLatestQueue.init(logExecutor, v -> { tsLatestQueue.init(logExecutor,
Map<TsKey, TsKvLatestEntity> trueLatest = new HashMap<>(); v -> insertLatestTsRepository.saveOrUpdate(v),
v.forEach(ts -> { Comparator.comparing((Function<TsKvLatestEntity, UUID>) AbstractTsKvEntity::getEntityId)
TsKey key = new TsKey(ts.getEntityId(), ts.getKey()); .thenComparingInt(AbstractTsKvEntity::getKey),
trueLatest.merge(key, ts, (oldTs, newTs) -> oldTs.getTs() <= newTs.getTs() ? newTs : oldTs); v -> {
}); Map<TsKey, TbSqlQueueElement<TsKvLatestEntity, Long>> trueLatest = new HashMap<>();
List<TsKvLatestEntity> latestEntities = new ArrayList<>(trueLatest.values()); v.forEach(element -> {
if (batchSortEnabled) { var entity = element.getEntity();
latestEntities.sort(Comparator.comparing((Function<TsKvLatestEntity, UUID>) AbstractTsKvEntity::getEntityId) TsKey key = new TsKey(entity.getEntityId(), entity.getKey());
.thenComparingInt(AbstractTsKvEntity::getKey)); trueLatest.merge(key, element, (oldElement, newElement) -> oldElement.getEntity().getTs() <= newElement.getEntity().getTs() ? newElement : oldElement);
} });
insertLatestTsRepository.saveOrUpdate(latestEntities); return new ArrayList<>(trueLatest.values());
}, (l, r) -> 0); });
} }
@PreDestroy @PreDestroy
@ -144,7 +146,7 @@ public class SqlTimeseriesLatestDao extends BaseAbstractSqlTimeseriesDao impleme
} }
@Override @Override
public ListenableFuture<Void> saveLatest(TenantId tenantId, EntityId entityId, TsKvEntry tsKvEntry) { public ListenableFuture<Long> saveLatest(TenantId tenantId, EntityId entityId, TsKvEntry tsKvEntry) {
return getSaveLatestFuture(entityId, tsKvEntry); return getSaveLatestFuture(entityId, tsKvEntry);
} }
@ -247,7 +249,7 @@ public class SqlTimeseriesLatestDao extends BaseAbstractSqlTimeseriesDao impleme
searchTsKvLatestRepository.findAllByEntityId(entityId.getId())))); searchTsKvLatestRepository.findAllByEntityId(entityId.getId()))));
} }
protected ListenableFuture<Void> getSaveLatestFuture(EntityId entityId, TsKvEntry tsKvEntry) { protected ListenableFuture<Long> getSaveLatestFuture(EntityId entityId, TsKvEntry tsKvEntry) {
TsKvLatestEntity latestEntity = new TsKvLatestEntity(); TsKvLatestEntity latestEntity = new TsKvLatestEntity();
latestEntity.setEntityId(entityId.getId()); latestEntity.setEntityId(entityId.getId());
latestEntity.setTs(tsKvEntry.getTs()); latestEntity.setTs(tsKvEntry.getTs());

224
dao/src/main/java/org/thingsboard/server/dao/sqlts/insert/latest/sql/SqlLatestInsertTsRepository.java

@ -17,33 +17,25 @@ package org.thingsboard.server.dao.sqlts.insert.latest.sql;
import jakarta.annotation.PostConstruct; import jakarta.annotation.PostConstruct;
import org.springframework.beans.factory.annotation.Value; import org.springframework.beans.factory.annotation.Value;
import org.springframework.jdbc.core.BatchPreparedStatementSetter;
import org.springframework.jdbc.core.PreparedStatementCreator;
import org.springframework.jdbc.core.SqlProvider;
import org.springframework.jdbc.support.GeneratedKeyHolder;
import org.springframework.jdbc.support.KeyHolder;
import org.springframework.stereotype.Repository; import org.springframework.stereotype.Repository;
import org.springframework.transaction.annotation.Transactional; import org.springframework.transaction.annotation.Transactional;
import org.thingsboard.server.dao.AbstractSequenceInsertRepository;
import org.thingsboard.server.dao.model.sqlts.latest.TsKvLatestEntity; import org.thingsboard.server.dao.model.sqlts.latest.TsKvLatestEntity;
import org.thingsboard.server.dao.sqlts.insert.AbstractInsertRepository;
import org.thingsboard.server.dao.sqlts.insert.latest.InsertLatestTsRepository; import org.thingsboard.server.dao.sqlts.insert.latest.InsertLatestTsRepository;
import org.thingsboard.server.dao.util.SqlDao; import org.thingsboard.server.dao.util.SqlDao;
import org.thingsboard.server.dao.util.SqlTsLatestAnyDao; import org.thingsboard.server.dao.util.SqlTsLatestAnyDao;
import java.sql.Connection;
import java.sql.PreparedStatement; import java.sql.PreparedStatement;
import java.sql.SQLException; import java.sql.SQLException;
import java.sql.Types; import java.sql.Types;
import java.util.ArrayList;
import java.util.List; import java.util.List;
import java.util.Map;
@SqlTsLatestAnyDao @SqlTsLatestAnyDao
@Repository @Repository
@Transactional @Transactional
@SqlDao @SqlDao
public class SqlLatestInsertTsRepository extends AbstractInsertRepository implements InsertLatestTsRepository { public class SqlLatestInsertTsRepository extends AbstractSequenceInsertRepository<TsKvLatestEntity> implements InsertLatestTsRepository {
@Value("${sql.ts_latest.update_by_latest_ts:true}") @Value("${sql.ts_latest.update_by_latest_ts:true}")
private Boolean updateByLatestTs; private Boolean updateByLatestTs;
@ -61,8 +53,6 @@ public class SqlLatestInsertTsRepository extends AbstractInsertRepository implem
private static final String RETURNING = " RETURNING seq_number"; private static final String RETURNING = " RETURNING seq_number";
private static final String SEQ_NUMBER = "seq_number";
private String batchUpdateQuery; private String batchUpdateQuery;
private String insertOrUpdateQuery; private String insertOrUpdateQuery;
@ -73,161 +63,89 @@ public class SqlLatestInsertTsRepository extends AbstractInsertRepository implem
} }
@Override @Override
public List<Long> saveOrUpdate(List<TsKvLatestEntity> entities) { protected void setOnBatchUpdateValues(PreparedStatement ps, int i, List<TsKvLatestEntity> entities) throws SQLException {
return transactionTemplate.execute(status -> { TsKvLatestEntity tsKvLatestEntity = entities.get(i);
List<Long> seqNumbers = new ArrayList<>(entities.size()); ps.setLong(1, tsKvLatestEntity.getTs());
if (tsKvLatestEntity.getBooleanValue() != null) {
ps.setBoolean(2, tsKvLatestEntity.getBooleanValue());
} else {
ps.setNull(2, Types.BOOLEAN);
}
KeyHolder keyHolder = new GeneratedKeyHolder(); ps.setString(3, replaceNullChars(tsKvLatestEntity.getStrValue()));
int[] updateResult = onBatchUpdate(entities, keyHolder); if (tsKvLatestEntity.getLongValue() != null) {
ps.setLong(4, tsKvLatestEntity.getLongValue());
} else {
ps.setNull(4, Types.BIGINT);
}
List<Map<String, Object>> seqNumbersList = keyHolder.getKeyList(); if (tsKvLatestEntity.getDoubleValue() != null) {
ps.setDouble(5, tsKvLatestEntity.getDoubleValue());
} else {
ps.setNull(5, Types.DOUBLE);
}
int notUpdatedCount = entities.size() - seqNumbersList.size(); ps.setString(6, replaceNullChars(tsKvLatestEntity.getJsonValue()));
List<Integer> toInsertIndexes = new ArrayList<>(notUpdatedCount); ps.setObject(7, tsKvLatestEntity.getEntityId());
List<TsKvLatestEntity> insertEntities = new ArrayList<>(notUpdatedCount); ps.setInt(8, tsKvLatestEntity.getKey());
int keyHolderIndex = 0; if (updateByLatestTs) {
for (int i = 0; i < updateResult.length; i++) { ps.setLong(9, tsKvLatestEntity.getTs());
if (updateResult[i] == 0) { }
insertEntities.add(entities.get(i)); }
seqNumbers.add(0L);
toInsertIndexes.add(i);
} else {
seqNumbers.add((Long) seqNumbersList.get(keyHolderIndex).get(SEQ_NUMBER));
keyHolderIndex++;
}
}
if (insertEntities.isEmpty()) { @Override
return seqNumbers; protected void setOnInsertOrUpdateValues(PreparedStatement ps, int i, List<TsKvLatestEntity> insertEntities) throws SQLException {
} TsKvLatestEntity tsKvLatestEntity = insertEntities.get(i);
ps.setObject(1, tsKvLatestEntity.getEntityId());
ps.setInt(2, tsKvLatestEntity.getKey());
ps.setLong(3, tsKvLatestEntity.getTs());
ps.setLong(9, tsKvLatestEntity.getTs());
if (updateByLatestTs) {
ps.setLong(15, tsKvLatestEntity.getTs());
}
onInsertOrUpdate(insertEntities, keyHolder); if (tsKvLatestEntity.getBooleanValue() != null) {
ps.setBoolean(4, tsKvLatestEntity.getBooleanValue());
ps.setBoolean(10, tsKvLatestEntity.getBooleanValue());
} else {
ps.setNull(4, Types.BOOLEAN);
ps.setNull(10, Types.BOOLEAN);
}
seqNumbersList = keyHolder.getKeyList(); ps.setString(5, replaceNullChars(tsKvLatestEntity.getStrValue()));
ps.setString(11, replaceNullChars(tsKvLatestEntity.getStrValue()));
for (int i = 0; i < seqNumbersList.size(); i++) { if (tsKvLatestEntity.getLongValue() != null) {
seqNumbers.set(toInsertIndexes.get(i), (Long) seqNumbersList.get(i).get(SEQ_NUMBER)); ps.setLong(6, tsKvLatestEntity.getLongValue());
} ps.setLong(12, tsKvLatestEntity.getLongValue());
} else {
ps.setNull(6, Types.BIGINT);
ps.setNull(12, Types.BIGINT);
}
return seqNumbers; if (tsKvLatestEntity.getDoubleValue() != null) {
}); ps.setDouble(7, tsKvLatestEntity.getDoubleValue());
} ps.setDouble(13, tsKvLatestEntity.getDoubleValue());
} else {
ps.setNull(7, Types.DOUBLE);
ps.setNull(13, Types.DOUBLE);
}
private int[] onBatchUpdate(List<TsKvLatestEntity> entities, KeyHolder keyHolder) { ps.setString(8, replaceNullChars(tsKvLatestEntity.getJsonValue()));
return jdbcTemplate.batchUpdate(new SimplePreparedStatementCreator(batchUpdateQuery), new BatchPreparedStatementSetter() { ps.setString(14, replaceNullChars(tsKvLatestEntity.getJsonValue()));
@Override
public void setValues(PreparedStatement ps, int i) throws SQLException {
TsKvLatestEntity tsKvLatestEntity = entities.get(i);
ps.setLong(1, tsKvLatestEntity.getTs());
if (tsKvLatestEntity.getBooleanValue() != null) {
ps.setBoolean(2, tsKvLatestEntity.getBooleanValue());
} else {
ps.setNull(2, Types.BOOLEAN);
}
ps.setString(3, replaceNullChars(tsKvLatestEntity.getStrValue()));
if (tsKvLatestEntity.getLongValue() != null) {
ps.setLong(4, tsKvLatestEntity.getLongValue());
} else {
ps.setNull(4, Types.BIGINT);
}
if (tsKvLatestEntity.getDoubleValue() != null) {
ps.setDouble(5, tsKvLatestEntity.getDoubleValue());
} else {
ps.setNull(5, Types.DOUBLE);
}
ps.setString(6, replaceNullChars(tsKvLatestEntity.getJsonValue()));
ps.setObject(7, tsKvLatestEntity.getEntityId());
ps.setInt(8, tsKvLatestEntity.getKey());
if (updateByLatestTs) {
ps.setLong(9, tsKvLatestEntity.getTs());
}
}
@Override
public int getBatchSize() {
return entities.size();
}
}, keyHolder);
} }
private void onInsertOrUpdate(List<TsKvLatestEntity> insertEntities, KeyHolder keyHolder) { @Override
jdbcTemplate.batchUpdate(new SimplePreparedStatementCreator(insertOrUpdateQuery), new BatchPreparedStatementSetter() { protected String getBatchUpdateQuery() {
@Override return batchUpdateQuery;
public void setValues(PreparedStatement ps, int i) throws SQLException {
TsKvLatestEntity tsKvLatestEntity = insertEntities.get(i);
ps.setObject(1, tsKvLatestEntity.getEntityId());
ps.setInt(2, tsKvLatestEntity.getKey());
ps.setLong(3, tsKvLatestEntity.getTs());
ps.setLong(9, tsKvLatestEntity.getTs());
if (updateByLatestTs) {
ps.setLong(15, tsKvLatestEntity.getTs());
}
if (tsKvLatestEntity.getBooleanValue() != null) {
ps.setBoolean(4, tsKvLatestEntity.getBooleanValue());
ps.setBoolean(10, tsKvLatestEntity.getBooleanValue());
} else {
ps.setNull(4, Types.BOOLEAN);
ps.setNull(10, Types.BOOLEAN);
}
ps.setString(5, replaceNullChars(tsKvLatestEntity.getStrValue()));
ps.setString(11, replaceNullChars(tsKvLatestEntity.getStrValue()));
if (tsKvLatestEntity.getLongValue() != null) {
ps.setLong(6, tsKvLatestEntity.getLongValue());
ps.setLong(12, tsKvLatestEntity.getLongValue());
} else {
ps.setNull(6, Types.BIGINT);
ps.setNull(12, Types.BIGINT);
}
if (tsKvLatestEntity.getDoubleValue() != null) {
ps.setDouble(7, tsKvLatestEntity.getDoubleValue());
ps.setDouble(13, tsKvLatestEntity.getDoubleValue());
} else {
ps.setNull(7, Types.DOUBLE);
ps.setNull(13, Types.DOUBLE);
}
ps.setString(8, replaceNullChars(tsKvLatestEntity.getJsonValue()));
ps.setString(14, replaceNullChars(tsKvLatestEntity.getJsonValue()));
}
@Override
public int getBatchSize() {
return insertEntities.size();
}
}, keyHolder);
} }
private static class SimplePreparedStatementCreator implements PreparedStatementCreator, SqlProvider { @Override
protected String getInsertOrUpdateQuery() {
private static final String[] COLUMNS = {SEQ_NUMBER}; return insertOrUpdateQuery;
private final String sql;
public SimplePreparedStatementCreator(String sql) {
this.sql = sql;
}
@Override
public PreparedStatement createPreparedStatement(Connection con) throws SQLException {
return con.prepareStatement(sql, COLUMNS);
}
@Override
public String getSql() {
return this.sql;
}
} }
} }

8
dao/src/main/java/org/thingsboard/server/dao/sqlts/timescale/TimescaleTimeseriesDao.java

@ -18,8 +18,9 @@ package org.thingsboard.server.dao.sqlts.timescale;
import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.ListenableFuture;
import com.google.common.util.concurrent.MoreExecutors; import com.google.common.util.concurrent.MoreExecutors;
import jakarta.annotation.PostConstruct;
import jakarta.annotation.PreDestroy;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.jetbrains.annotations.NotNull;
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.data.domain.PageRequest; import org.springframework.data.domain.PageRequest;
import org.springframework.data.domain.Sort; import org.springframework.data.domain.Sort;
@ -38,7 +39,6 @@ import org.thingsboard.server.dao.DaoUtil;
import org.thingsboard.server.dao.dictionary.KeyDictionaryDao; import org.thingsboard.server.dao.dictionary.KeyDictionaryDao;
import org.thingsboard.server.dao.model.sql.AbstractTsKvEntity; import org.thingsboard.server.dao.model.sql.AbstractTsKvEntity;
import org.thingsboard.server.dao.model.sqlts.timescale.ts.TimescaleTsKvEntity; import org.thingsboard.server.dao.model.sqlts.timescale.ts.TimescaleTsKvEntity;
import org.thingsboard.server.dao.model.sqlts.ts.TsKvEntity;
import org.thingsboard.server.dao.sql.TbSqlBlockingQueueParams; import org.thingsboard.server.dao.sql.TbSqlBlockingQueueParams;
import org.thingsboard.server.dao.sql.TbSqlBlockingQueueWrapper; import org.thingsboard.server.dao.sql.TbSqlBlockingQueueWrapper;
import org.thingsboard.server.dao.sqlts.AbstractSqlTimeseriesDao; import org.thingsboard.server.dao.sqlts.AbstractSqlTimeseriesDao;
@ -47,8 +47,6 @@ import org.thingsboard.server.dao.timeseries.TimeseriesDao;
import org.thingsboard.server.dao.util.TimeUtils; import org.thingsboard.server.dao.util.TimeUtils;
import org.thingsboard.server.dao.util.TimescaleDBTsDao; import org.thingsboard.server.dao.util.TimescaleDBTsDao;
import jakarta.annotation.PostConstruct;
import jakarta.annotation.PreDestroy;
import java.util.ArrayList; import java.util.ArrayList;
import java.util.Collections; import java.util.Collections;
import java.util.Comparator; import java.util.Comparator;
@ -77,7 +75,7 @@ public class TimescaleTimeseriesDao extends AbstractSqlTimeseriesDao implements
@Autowired @Autowired
protected KeyDictionaryDao keyDictionaryDao; protected KeyDictionaryDao keyDictionaryDao;
protected TbSqlBlockingQueueWrapper<TimescaleTsKvEntity> tsQueue; protected TbSqlBlockingQueueWrapper<TimescaleTsKvEntity, Void> tsQueue;
@PostConstruct @PostConstruct
protected void init() { protected void init() {

4
dao/src/main/java/org/thingsboard/server/dao/timeseries/BaseTimeseriesService.java

@ -187,8 +187,8 @@ public class BaseTimeseriesService implements TimeseriesService {
} }
@Override @Override
public ListenableFuture<List<Void>> saveLatest(TenantId tenantId, EntityId entityId, List<TsKvEntry> tsKvEntries) { public ListenableFuture<List<Long>> saveLatest(TenantId tenantId, EntityId entityId, List<TsKvEntry> tsKvEntries) {
List<ListenableFuture<Void>> futures = new ArrayList<>(tsKvEntries.size()); List<ListenableFuture<Long>> futures = new ArrayList<>(tsKvEntries.size());
for (TsKvEntry tsKvEntry : tsKvEntries) { for (TsKvEntry tsKvEntry : tsKvEntries) {
futures.add(timeseriesLatestDao.saveLatest(tenantId, entityId, tsKvEntry)); futures.add(timeseriesLatestDao.saveLatest(tenantId, entityId, tsKvEntry));
} }

2
dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesLatestDao.java

@ -100,7 +100,7 @@ public class CassandraBaseTimeseriesLatestDao extends AbstractCassandraBaseTimes
} }
@Override @Override
public ListenableFuture<Void> saveLatest(TenantId tenantId, EntityId entityId, TsKvEntry tsKvEntry) { public ListenableFuture<Long> saveLatest(TenantId tenantId, EntityId entityId, TsKvEntry tsKvEntry) {
BoundStatementBuilder stmtBuilder = new BoundStatementBuilder(getLatestStmt().bind()); BoundStatementBuilder stmtBuilder = new BoundStatementBuilder(getLatestStmt().bind());
stmtBuilder.setString(0, entityId.getEntityType().name()) stmtBuilder.setString(0, entityId.getEntityType().name())
.setUuid(1, entityId.getId()) .setUuid(1, entityId.getId())

2
dao/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesLatestDao.java

@ -42,7 +42,7 @@ public interface TimeseriesLatestDao {
ListenableFuture<List<TsKvEntry>> findAllLatest(TenantId tenantId, EntityId entityId); ListenableFuture<List<TsKvEntry>> findAllLatest(TenantId tenantId, EntityId entityId);
ListenableFuture<Void> saveLatest(TenantId tenantId, EntityId entityId, TsKvEntry tsKvEntry); ListenableFuture<Long> saveLatest(TenantId tenantId, EntityId entityId, TsKvEntry tsKvEntry);
ListenableFuture<TsKvLatestRemovingResult> removeLatest(TenantId tenantId, EntityId entityId, DeleteTsKvQuery query); ListenableFuture<TsKvLatestRemovingResult> removeLatest(TenantId tenantId, EntityId entityId, DeleteTsKvQuery query);

3
dao/src/main/resources/sql/schema-entities.sql

@ -102,6 +102,8 @@ CREATE TABLE IF NOT EXISTS audit_log (
action_failure_details varchar(1000000) action_failure_details varchar(1000000)
) PARTITION BY RANGE (created_time); ) PARTITION BY RANGE (created_time);
CREATE SEQUENCE IF NOT EXISTS attribute_kv_latest_seq cache 1000;
CREATE TABLE IF NOT EXISTS attribute_kv ( CREATE TABLE IF NOT EXISTS attribute_kv (
entity_id uuid, entity_id uuid,
attribute_type int, attribute_type int,
@ -112,6 +114,7 @@ CREATE TABLE IF NOT EXISTS attribute_kv (
dbl_v double precision, dbl_v double precision,
json_v json, json_v json,
last_update_ts bigint, last_update_ts bigint,
seq_number bigint,
CONSTRAINT attribute_kv_pkey PRIMARY KEY (entity_id, attribute_type, attribute_key) CONSTRAINT attribute_kv_pkey PRIMARY KEY (entity_id, attribute_type, attribute_key)
); );

Loading…
Cancel
Save