Browse Source

feat(rpc): batched async upsert persistence in RPC DAO

pull/15853/head
dshvaika 4 weeks ago
parent
commit
0f41970cc7
  1. 6
      application/src/main/resources/thingsboard.yml
  2. 2
      common/dao-api/src/main/java/org/thingsboard/server/dao/rpc/RpcService.java
  3. 6
      dao/src/main/java/org/thingsboard/server/dao/rpc/BaseRpcService.java
  4. 3
      dao/src/main/java/org/thingsboard/server/dao/rpc/RpcDao.java
  5. 64
      dao/src/main/java/org/thingsboard/server/dao/sql/rpc/JpaRpcDao.java
  6. 66
      dao/src/main/java/org/thingsboard/server/dao/sql/rpc/RpcInsertRepository.java
  7. 36
      dao/src/test/java/org/thingsboard/server/dao/sql/rpc/JpaRpcDaoTest.java

6
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

2
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<Void> saveAsync(Rpc rpc);
void deleteRpc(TenantId tenantId, RpcId id);
void deleteAllRpcByTenantId(TenantId tenantId);

6
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<Void> 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);

3
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<Rpc> {
ListenableFuture<Void> saveAsync(TenantId tenantId, Rpc rpc);
PageData<Rpc> findAllByDeviceId(TenantId tenantId, DeviceId deviceId, PageLink pageLink);
PageData<Rpc> findAllByDeviceIdAndStatus(TenantId tenantId, DeviceId deviceId, RpcStatus rpcStatus, PageLink pageLink);

64
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<RpcEntity, Rpc> implements RpcDao, TenantEntityDao<Rpc> {
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<RpcEntity, Void> 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<RpcEntity, Integer> 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<Void> saveAsync(TenantId tenantId, Rpc rpc) {
return queue.add(new RpcEntity(rpc));
}
@Override
protected Class<RpcEntity> getEntityClass() {

66
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<RpcEntity> 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));
}
}

36
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}"));
}
}

Loading…
Cancel
Save