diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index f678076c57..3209c54746 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -424,6 +424,12 @@ sql: stats_print_interval_ms: "${SQL_TS_LATEST_BATCH_STATS_PRINT_MS:10000}" # Interval in milliseconds for printing latest telemetry updates statistic batch_threads: "${SQL_TS_LATEST_BATCH_THREADS:3}" # batch thread count has to be a prime number like 3 or 5 to gain perfect hash distribution update_by_latest_ts: "${SQL_TS_UPDATE_BY_LATEST_TIMESTAMP:true}" # Update latest values only if the timestamp of the new record is greater or equals the timestamp of the previously saved latest value. The latest values are stored separately from historical values for fast lookup from DB. Insert of historical value happens in any case + rpc: + batch_size: "${SQL_RPC_BATCH_SIZE:1000}" # Batch size for persisting RPC inserts/updates + batch_max_delay: "${SQL_RPC_BATCH_MAX_DELAY_MS:50}" # Max timeout for RPC entries queue polling, in milliseconds + stats_print_interval_ms: "${SQL_RPC_BATCH_STATS_PRINT_MS:10000}" # Interval in milliseconds for printing RPC persistence statistic + batch_threads: "${SQL_RPC_BATCH_THREADS:3}" # batch thread count has to be a prime number like 3 or 5 to gain perfect hash distribution + callback_threads: "${SQL_RPC_CALLBACK_THREADS:3}" # striped threads for post-persist rule-engine notifications; keep equal to batch_threads for aligned ordering events: batch_size: "${SQL_EVENTS_BATCH_SIZE:10000}" # Batch size for persisting latest telemetry updates batch_max_delay: "${SQL_EVENTS_BATCH_MAX_DELAY_MS:100}" # Max timeout for latest telemetry entries queue polling. The value set in milliseconds diff --git a/common/dao-api/src/main/java/org/thingsboard/server/dao/rpc/RpcService.java b/common/dao-api/src/main/java/org/thingsboard/server/dao/rpc/RpcService.java index 855063d657..deb3bf8cf6 100644 --- a/common/dao-api/src/main/java/org/thingsboard/server/dao/rpc/RpcService.java +++ b/common/dao-api/src/main/java/org/thingsboard/server/dao/rpc/RpcService.java @@ -29,6 +29,8 @@ public interface RpcService extends EntityDaoService { Rpc save(Rpc rpc); + ListenableFuture saveAsync(Rpc rpc); + void deleteRpc(TenantId tenantId, RpcId id); void deleteAllRpcByTenantId(TenantId tenantId); diff --git a/dao/src/main/java/org/thingsboard/server/dao/rpc/BaseRpcService.java b/dao/src/main/java/org/thingsboard/server/dao/rpc/BaseRpcService.java index b096c69d4b..4162422e69 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/rpc/BaseRpcService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/rpc/BaseRpcService.java @@ -54,6 +54,12 @@ public class BaseRpcService implements RpcService { return rpcDao.save(rpc.getTenantId(), rpc); } + @Override + public ListenableFuture saveAsync(Rpc rpc) { + log.trace("Executing saveAsync, [{}]", rpc); + return rpcDao.saveAsync(rpc.getTenantId(), rpc); + } + @Override public void deleteRpc(TenantId tenantId, RpcId rpcId) { log.trace("Executing deleteRpc, tenantId [{}], rpcId [{}]", tenantId, rpcId); diff --git a/dao/src/main/java/org/thingsboard/server/dao/rpc/RpcDao.java b/dao/src/main/java/org/thingsboard/server/dao/rpc/RpcDao.java index 37fe2950b4..a823b05507 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/rpc/RpcDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/rpc/RpcDao.java @@ -15,6 +15,7 @@ */ package org.thingsboard.server.dao.rpc; +import com.google.common.util.concurrent.ListenableFuture; import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.page.PageData; @@ -25,6 +26,8 @@ import org.thingsboard.server.dao.Dao; public interface RpcDao extends Dao { + ListenableFuture saveAsync(TenantId tenantId, Rpc rpc); + PageData findAllByDeviceId(TenantId tenantId, DeviceId deviceId, PageLink pageLink); PageData findAllByDeviceIdAndStatus(TenantId tenantId, DeviceId deviceId, RpcStatus rpcStatus, PageLink pageLink); diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/rpc/JpaRpcDao.java b/dao/src/main/java/org/thingsboard/server/dao/sql/rpc/JpaRpcDao.java index a929b6d1e4..e46dc2a704 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/rpc/JpaRpcDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/rpc/JpaRpcDao.java @@ -15,8 +15,12 @@ */ package org.thingsboard.server.dao.sql.rpc; -import lombok.AllArgsConstructor; +import com.google.common.util.concurrent.ListenableFuture; +import jakarta.annotation.PostConstruct; +import jakarta.annotation.PreDestroy; import lombok.extern.slf4j.Slf4j; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.annotation.Value; import org.springframework.data.jpa.repository.JpaRepository; import org.springframework.stereotype.Component; import org.springframework.transaction.annotation.Transactional; @@ -27,22 +31,76 @@ import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.page.PageLink; import org.thingsboard.server.common.data.rpc.Rpc; import org.thingsboard.server.common.data.rpc.RpcStatus; +import org.thingsboard.server.common.stats.StatsFactory; import org.thingsboard.server.dao.DaoUtil; import org.thingsboard.server.dao.TenantEntityDao; import org.thingsboard.server.dao.model.sql.RpcEntity; import org.thingsboard.server.dao.rpc.RpcDao; import org.thingsboard.server.dao.sql.JpaAbstractDao; +import org.thingsboard.server.dao.sql.ScheduledLogExecutorComponent; +import org.thingsboard.server.dao.sql.TbSqlBlockingQueueParams; +import org.thingsboard.server.dao.sql.TbSqlBlockingQueueWrapper; import org.thingsboard.server.dao.util.SqlDao; +import java.util.Comparator; import java.util.UUID; +import java.util.function.Function; @Slf4j @Component -@AllArgsConstructor @SqlDao public class JpaRpcDao extends JpaAbstractDao implements RpcDao, TenantEntityDao { - private final RpcRepository rpcRepository; + @Autowired + private RpcRepository rpcRepository; + @Autowired + private RpcInsertRepository rpcInsertRepository; + @Autowired + private ScheduledLogExecutorComponent logExecutor; + @Autowired + private StatsFactory statsFactory; + + @Value("${sql.rpc.batch_size:1000}") + private int batchSize; + @Value("${sql.rpc.batch_max_delay:50}") + private long maxDelay; + @Value("${sql.rpc.stats_print_interval_ms:10000}") + private long statsPrintIntervalMs; + @Value("${sql.rpc.batch_threads:3}") + private int batchThreads; + @Value("${sql.batch_sort:true}") + private boolean batchSortEnabled; + + private TbSqlBlockingQueueWrapper queue; + + @PostConstruct + private void init() { + TbSqlBlockingQueueParams params = TbSqlBlockingQueueParams.builder() + .logName("RPC") + .batchSize(batchSize) + .maxDelay(maxDelay) + .statsPrintIntervalMs(statsPrintIntervalMs) + .statsNamePrefix("rpc") + .batchSortEnabled(batchSortEnabled) + .withResponse(false) + .build(); + Function hashcodeFunction = entity -> entity.getUuid().hashCode(); + queue = new TbSqlBlockingQueueWrapper<>(params, hashcodeFunction, batchThreads, statsFactory); + queue.init(logExecutor, entities -> rpcInsertRepository.saveOrUpdate(entities), + Comparator.comparing(RpcEntity::getUuid)); + } + + @PreDestroy + private void destroy() { + if (queue != null) { + queue.destroy(); + } + } + + @Override + public ListenableFuture saveAsync(TenantId tenantId, Rpc rpc) { + return queue.add(new RpcEntity(rpc)); + } @Override protected Class getEntityClass() { diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/rpc/RpcInsertRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sql/rpc/RpcInsertRepository.java new file mode 100644 index 0000000000..5073702ef4 --- /dev/null +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/rpc/RpcInsertRepository.java @@ -0,0 +1,66 @@ +/** + * Copyright © 2016-2026 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.rpc; + +import com.fasterxml.jackson.databind.JsonNode; +import org.springframework.jdbc.core.BatchPreparedStatementSetter; +import org.springframework.stereotype.Repository; +import org.thingsboard.common.util.JacksonUtil; +import org.thingsboard.server.dao.model.sql.RpcEntity; +import org.thingsboard.server.dao.sqlts.insert.AbstractInsertRepository; + +import java.sql.PreparedStatement; +import java.sql.SQLException; +import java.util.List; + +@Repository +public class RpcInsertRepository extends AbstractInsertRepository { + + private static final String INSERT_OR_UPDATE = + "INSERT INTO rpc (id, created_time, tenant_id, device_id, expiration_time, request, response, additional_info, status) " + + "VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?) " + + "ON CONFLICT (id) DO UPDATE SET status = EXCLUDED.status, response = EXCLUDED.response;"; + + public void saveOrUpdate(List entities) { + transactionTemplate.execute(status -> { + jdbcTemplate.batchUpdate(INSERT_OR_UPDATE, new BatchPreparedStatementSetter() { + @Override + public void setValues(PreparedStatement ps, int i) throws SQLException { + RpcEntity rpc = entities.get(i); + ps.setObject(1, rpc.getUuid()); + ps.setLong(2, rpc.getCreatedTime()); + ps.setObject(3, rpc.getTenantId()); + ps.setObject(4, rpc.getDeviceId()); + ps.setLong(5, rpc.getExpirationTime()); + ps.setString(6, toJsonStr(rpc.getRequest())); + ps.setString(7, toJsonStr(rpc.getResponse())); + ps.setString(8, toJsonStr(rpc.getAdditionalInfo())); + ps.setString(9, rpc.getStatus().name()); + } + + @Override + public int getBatchSize() { + return entities.size(); + } + }); + return null; + }); + } + + private String toJsonStr(JsonNode node) { + return node == null ? null : replaceNullChars(JacksonUtil.toString(node)); + } +} diff --git a/dao/src/test/java/org/thingsboard/server/dao/sql/rpc/JpaRpcDaoTest.java b/dao/src/test/java/org/thingsboard/server/dao/sql/rpc/JpaRpcDaoTest.java index 921339a92b..8b95aa7c3b 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/sql/rpc/JpaRpcDaoTest.java +++ b/dao/src/test/java/org/thingsboard/server/dao/sql/rpc/JpaRpcDaoTest.java @@ -19,12 +19,14 @@ import org.junit.Test; import org.springframework.beans.factory.annotation.Autowired; import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.server.common.data.id.DeviceId; +import org.thingsboard.server.common.data.id.RpcId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.rpc.Rpc; import org.thingsboard.server.common.data.rpc.RpcStatus; import org.thingsboard.server.dao.AbstractJpaDaoTest; import java.util.UUID; +import java.util.concurrent.TimeUnit; import static org.assertj.core.api.Assertions.assertThat; @@ -57,4 +59,38 @@ public class JpaRpcDaoTest extends AbstractJpaDaoTest { assertThat(rpcDao.deleteOutdatedRpcByTenantIdBatch(tenantId, System.currentTimeMillis() + 1, batchSize)).isEqualTo(1); } + @Test + public void saveAsyncInsertThenUpsert() throws Exception { + UUID id = UUID.randomUUID(); + Rpc rpc = new Rpc(new RpcId(id)); + rpc.setCreatedTime(System.currentTimeMillis()); + rpc.setTenantId(TenantId.SYS_TENANT_ID); + rpc.setDeviceId(new DeviceId(UUID.randomUUID())); + rpc.setExpirationTime(System.currentTimeMillis() + 60_000); + rpc.setRequest(JacksonUtil.toJsonNode("{\"method\":\"x\"}")); + rpc.setStatus(RpcStatus.QUEUED); + + rpcDao.saveAsync(rpc.getTenantId(), rpc).get(5, TimeUnit.SECONDS); + + Rpc stored = rpcDao.findById(TenantId.SYS_TENANT_ID, id); + assertThat(stored).isNotNull(); + assertThat(stored.getStatus()).isEqualTo(RpcStatus.QUEUED); + assertThat(stored.getResponse()).isNull(); + + Rpc update = new Rpc(new RpcId(id)); + update.setCreatedTime(rpc.getCreatedTime()); + update.setTenantId(TenantId.SYS_TENANT_ID); + update.setDeviceId(rpc.getDeviceId()); + update.setExpirationTime(rpc.getExpirationTime()); + update.setRequest(rpc.getRequest()); + update.setStatus(RpcStatus.DELIVERED); + update.setResponse(JacksonUtil.toJsonNode("{\"ok\":true}")); + + rpcDao.saveAsync(update.getTenantId(), update).get(5, TimeUnit.SECONDS); + + Rpc afterUpdate = rpcDao.findById(TenantId.SYS_TENANT_ID, id); + assertThat(afterUpdate.getStatus()).isEqualTo(RpcStatus.DELIVERED); + assertThat(afterUpdate.getResponse()).isEqualTo(JacksonUtil.toJsonNode("{\"ok\":true}")); + } + }