Browse Source

fix(rpc): skip rule-engine notify when async update matches no row

Async RPC persistence now returns a per-write Boolean: an INSERT-on-conflict always persists (true), while an UPDATE-by-id reports false when its WHERE id = ? matched no row (RPC deleted via TTL/manual delete). TbRpcService notifies the rule engine only when the write actually persisted, restoring the findById-null skip the async refactor had dropped.

Also: drop the unused tenantId param from createAsync/updateAsync; make TbRpcService's callback-thread count constructor-injectable (removes test reflection) with a comment on why callback striping exists alongside the queue partitioning; correct the RPC batch_threads yml comment; assert the Boolean contract in JpaRpcDaoTest and add a notification-suppression unit test.
pull/15853/head
dshvaika 3 months ago
parent
commit
ef571175f7
  1. 41
      application/src/main/java/org/thingsboard/server/service/rpc/TbRpcService.java
  2. 2
      application/src/main/resources/thingsboard.yml
  3. 24
      application/src/test/java/org/thingsboard/server/service/rpc/TbRpcServiceTest.java
  4. 4
      common/dao-api/src/main/java/org/thingsboard/server/dao/rpc/RpcService.java
  5. 8
      dao/src/main/java/org/thingsboard/server/dao/rpc/BaseRpcService.java
  6. 4
      dao/src/main/java/org/thingsboard/server/dao/rpc/RpcDao.java
  7. 14
      dao/src/main/java/org/thingsboard/server/dao/sql/rpc/JpaRpcDao.java
  8. 24
      dao/src/main/java/org/thingsboard/server/dao/sql/rpc/RpcInsertRepository.java
  9. 47
      dao/src/test/java/org/thingsboard/server/dao/sql/rpc/JpaRpcDaoTest.java

41
application/src/main/java/org/thingsboard/server/service/rpc/TbRpcService.java

@ -16,9 +16,7 @@
package org.thingsboard.server.service.rpc;
import com.google.common.util.concurrent.ListenableFuture;
import jakarta.annotation.PostConstruct;
import jakarta.annotation.PreDestroy;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.stereotype.Service;
@ -46,20 +44,24 @@ import java.util.concurrent.Executors;
@TbCoreComponent
@Service
@RequiredArgsConstructor
@Slf4j
public class TbRpcService {
private final RpcService rpcService;
private final TbClusterService tbClusterService;
@Value("${sql.rpc.callback_threads:3}")
private int callbackThreads;
// Post-persist rule-engine notifications run on these striped single-thread executors, keyed by
// rpcId, instead of inline on the SQL persist threads. A flush batch can carry up to batch_size
// entries, and each notification's pushMsgToRuleEngine may do a (potentially blocking) device-profile
// cache lookup; running them inline would serialize that work on the persist thread and stall the
// next DB flush. Striping by rpcId (rather than a shared multi-threaded pool) also preserves
// per-command notification order, e.g. RPC_QUEUED before RPC_DELIVERED.
private final ExecutorService[] callbackExecutors;
private ExecutorService[] callbackExecutors;
@PostConstruct
private void init() {
callbackExecutors = new ExecutorService[callbackThreads];
public TbRpcService(RpcService rpcService, TbClusterService tbClusterService,
@Value("${sql.rpc.callback_threads:3}") int callbackThreads) {
this.rpcService = rpcService;
this.tbClusterService = tbClusterService;
this.callbackExecutors = new ExecutorService[callbackThreads];
for (int i = 0; i < callbackThreads; i++) {
callbackExecutors[i] = Executors.newSingleThreadExecutor(
ThingsBoardThreadFactory.forName("rpc-persist-callback-" + i));
@ -68,10 +70,8 @@ public class TbRpcService {
@PreDestroy
private void destroy() {
if (callbackExecutors != null) {
for (ExecutorService executor : callbackExecutors) {
executor.shutdownNow();
}
for (ExecutorService executor : callbackExecutors) {
executor.shutdownNow();
}
}
@ -83,16 +83,23 @@ public class TbRpcService {
persist(tenantId, rpc, rpcService.updateAsync(rpc));
}
private void persist(TenantId tenantId, Rpc rpc, ListenableFuture<Void> future) {
private void persist(TenantId tenantId, Rpc rpc, ListenableFuture<Boolean> future) {
DonAsynchron.withCallback(future,
v -> pushRpcMsgToRuleEngine(tenantId, rpc),
persisted -> {
if (Boolean.TRUE.equals(persisted)) {
pushRpcMsgToRuleEngine(tenantId, rpc);
} else {
log.debug("[{}][{}][{}] Skipping rule engine notification for status [{}] - RPC row no longer exists",
tenantId, rpc.getDeviceId(), rpc.getId(), rpc.getStatus());
}
},
t -> log.error("[{}][{}][{}] Failed to persist RPC with status [{}]",
tenantId, rpc.getDeviceId(), rpc.getId(), rpc.getStatus(), t),
executorFor(rpc.getUuidId()));
}
private Executor executorFor(UUID rpcId) {
return callbackExecutors[(rpcId.hashCode() & 0x7FFFFFFF) % callbackThreads];
return callbackExecutors[(rpcId.hashCode() & 0x7FFFFFFF) % callbackExecutors.length];
}
private void pushRpcMsgToRuleEngine(TenantId tenantId, Rpc rpc) {

2
application/src/main/resources/thingsboard.yml

@ -428,7 +428,7 @@ sql:
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
batch_threads: "${SQL_RPC_BATCH_THREADS:3}" # number of queue partitions for batched RPC persistence; writes for a given RPC always map to the same partition (by rpcId hash) so they persist in submission order
callback_threads: "${SQL_RPC_CALLBACK_THREADS:3}" # striped threads for post-persist rule-engine notifications; each rpcId maps to a single stripe, preserving per-command notification order
events:
batch_size: "${SQL_EVENTS_BATCH_SIZE:10000}" # Batch size for persisting latest telemetry updates

24
application/src/test/java/org/thingsboard/server/service/rpc/TbRpcServiceTest.java

@ -18,7 +18,6 @@ package org.thingsboard.server.service.rpc;
import com.google.common.util.concurrent.Futures;
import org.junit.Before;
import org.junit.Test;
import org.springframework.test.util.ReflectionTestUtils;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.server.cluster.TbClusterService;
import org.thingsboard.server.common.data.id.DeviceId;
@ -34,6 +33,7 @@ import java.util.UUID;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.ArgumentMatchers.isNull;
import static org.mockito.Mockito.after;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.timeout;
import static org.mockito.Mockito.verify;
@ -47,15 +47,13 @@ public class TbRpcServiceTest {
@Before
public void setUp() {
tbRpcService = new TbRpcService(rpcService, clusterService);
ReflectionTestUtils.setField(tbRpcService, "callbackThreads", 1);
ReflectionTestUtils.invokeMethod(tbRpcService, "init");
tbRpcService = new TbRpcService(rpcService, clusterService, 1);
}
@Test
public void savePersistsViaUpdateAsyncThenPushesToRuleEngine() {
Rpc rpc = newRpc();
when(rpcService.updateAsync(rpc)).thenReturn(Futures.immediateFuture(null));
when(rpcService.updateAsync(rpc)).thenReturn(Futures.immediateFuture(true));
tbRpcService.save(rpc.getTenantId(), rpc);
@ -67,7 +65,7 @@ public class TbRpcServiceTest {
@Test
public void createPersistsViaCreateAsyncThenPushesToRuleEngine() {
Rpc rpc = newRpc();
when(rpcService.createAsync(rpc)).thenReturn(Futures.immediateFuture(null));
when(rpcService.createAsync(rpc)).thenReturn(Futures.immediateFuture(true));
tbRpcService.create(rpc.getTenantId(), rpc);
@ -76,6 +74,20 @@ public class TbRpcServiceTest {
.pushMsgToRuleEngine(eq(rpc.getTenantId()), eq(rpc.getDeviceId()), any(TbMsg.class), isNull());
}
@Test
public void saveDoesNotNotifyRuleEngineWhenRowMissing() {
Rpc rpc = newRpc();
// updateAsync resolves false: the UPDATE matched no row (RPC was deleted), so the rule engine
// must not be notified for a status change that never persisted.
when(rpcService.updateAsync(rpc)).thenReturn(Futures.immediateFuture(false));
tbRpcService.save(rpc.getTenantId(), rpc);
verify(rpcService).updateAsync(rpc);
verify(clusterService, after(500).never())
.pushMsgToRuleEngine(any(TenantId.class), any(DeviceId.class), any(TbMsg.class), isNull());
}
private Rpc newRpc() {
Rpc rpc = new Rpc(new RpcId(UUID.randomUUID()));
rpc.setTenantId(TenantId.SYS_TENANT_ID);

4
common/dao-api/src/main/java/org/thingsboard/server/dao/rpc/RpcService.java

@ -29,9 +29,9 @@ public interface RpcService extends EntityDaoService {
Rpc save(Rpc rpc);
ListenableFuture<Void> createAsync(Rpc rpc);
ListenableFuture<Boolean> createAsync(Rpc rpc);
ListenableFuture<Void> updateAsync(Rpc rpc);
ListenableFuture<Boolean> updateAsync(Rpc rpc);
void deleteRpc(TenantId tenantId, RpcId id);

8
dao/src/main/java/org/thingsboard/server/dao/rpc/BaseRpcService.java

@ -55,15 +55,15 @@ public class BaseRpcService implements RpcService {
}
@Override
public ListenableFuture<Void> createAsync(Rpc rpc) {
public ListenableFuture<Boolean> createAsync(Rpc rpc) {
log.trace("Executing createAsync, [{}]", rpc);
return rpcDao.createAsync(rpc.getTenantId(), rpc);
return rpcDao.createAsync(rpc);
}
@Override
public ListenableFuture<Void> updateAsync(Rpc rpc) {
public ListenableFuture<Boolean> updateAsync(Rpc rpc) {
log.trace("Executing updateAsync, [{}]", rpc);
return rpcDao.updateAsync(rpc.getTenantId(), rpc);
return rpcDao.updateAsync(rpc);
}
@Override

4
dao/src/main/java/org/thingsboard/server/dao/rpc/RpcDao.java

@ -26,9 +26,9 @@ import org.thingsboard.server.dao.Dao;
public interface RpcDao extends Dao<Rpc> {
ListenableFuture<Void> createAsync(TenantId tenantId, Rpc rpc);
ListenableFuture<Boolean> createAsync(Rpc rpc);
ListenableFuture<Void> updateAsync(TenantId tenantId, Rpc rpc);
ListenableFuture<Boolean> updateAsync(Rpc rpc);
PageData<Rpc> findAllByDeviceId(TenantId tenantId, DeviceId deviceId, PageLink pageLink);

14
dao/src/main/java/org/thingsboard/server/dao/sql/rpc/JpaRpcDao.java

@ -71,7 +71,10 @@ public class JpaRpcDao extends JpaAbstractDao<RpcEntity, Rpc> implements RpcDao,
@Value("${sql.batch_sort:true}")
private boolean batchSortEnabled;
private TbSqlBlockingQueueWrapper<RpcQueueEntry, Void> queue;
// Boolean response per queued write: an INSERT-on-conflict always persists (true), while an
// UPDATE-by-id reports false when it matched no row (the RPC was deleted in the meantime), so the
// service layer can skip the rule-engine notification for a row that no longer exists.
private TbSqlBlockingQueueWrapper<RpcQueueEntry, Boolean> queue;
@PostConstruct
private void init() {
@ -82,12 +85,13 @@ public class JpaRpcDao extends JpaAbstractDao<RpcEntity, Rpc> implements RpcDao,
.statsPrintIntervalMs(statsPrintIntervalMs)
.statsNamePrefix("rpc")
.batchSortEnabled(batchSortEnabled)
.withResponse(false)
.withResponse(true)
.build();
Function<RpcQueueEntry, Integer> hashcodeFunction = entry -> entry.entity().getUuid().hashCode();
queue = new TbSqlBlockingQueueWrapper<>(params, hashcodeFunction, batchThreads, statsFactory);
queue.init(logExecutor, entries -> rpcInsertRepository.saveOrUpdate(entries),
Comparator.comparing((RpcQueueEntry entry) -> entry.entity().getUuid()));
Comparator.comparing((RpcQueueEntry entry) -> entry.entity().getUuid()),
Function.identity());
}
@PreDestroy
@ -98,12 +102,12 @@ public class JpaRpcDao extends JpaAbstractDao<RpcEntity, Rpc> implements RpcDao,
}
@Override
public ListenableFuture<Void> createAsync(TenantId tenantId, Rpc rpc) {
public ListenableFuture<Boolean> createAsync(Rpc rpc) {
return queue.add(RpcQueueEntry.forInsert(new RpcEntity(rpc)));
}
@Override
public ListenableFuture<Void> updateAsync(TenantId tenantId, Rpc rpc) {
public ListenableFuture<Boolean> updateAsync(Rpc rpc) {
return queue.add(RpcQueueEntry.forUpdate(new RpcEntity(rpc)));
}

24
dao/src/main/java/org/thingsboard/server/dao/sql/rpc/RpcInsertRepository.java

@ -24,6 +24,7 @@ import org.thingsboard.server.dao.sqlts.insert.AbstractInsertRepository;
import java.sql.PreparedStatement;
import java.sql.SQLException;
import java.util.ArrayList;
import java.util.List;
@Repository
@ -48,10 +49,10 @@ public class RpcInsertRepository extends AbstractInsertRepository {
void bind(PreparedStatement ps, RpcEntity rpc) throws SQLException;
}
public void saveOrUpdate(List<RpcQueueEntry> entries) {
public List<Boolean> saveOrUpdate(List<RpcQueueEntry> entries) {
List<RpcEntity> inserts = entries.stream().filter(RpcQueueEntry::insert).map(RpcQueueEntry::entity).toList();
List<RpcEntity> updates = entries.stream().filter(entry -> !entry.insert()).map(RpcQueueEntry::entity).toList();
transactionTemplate.execute(status -> {
int[] updateCounts = transactionTemplate.execute(status -> {
// Inserts run first so a create and a status update for the same rpcId coalesced into one
// batch still apply in create -> update order.
if (!inserts.isEmpty()) {
@ -68,18 +69,29 @@ public class RpcInsertRepository extends AbstractInsertRepository {
});
}
if (!updates.isEmpty()) {
batch(UPDATE, updates, (ps, rpc) -> {
return batch(UPDATE, updates, (ps, rpc) -> {
ps.setString(1, rpc.getStatus().name());
ps.setString(2, toJsonStr(rpc.getResponse()));
ps.setObject(3, rpc.getUuid());
});
}
return null;
return new int[0];
});
// Result is aligned to the submission order of entries: an insert always persists
// (INSERT ... ON CONFLICT), an update persists only if its WHERE id = ? matched a still-existing
// row. updateCounts keeps the same relative order as the filtered updates, so a single cursor
// walks it as we encounter update entries; insert entries short-circuit and don't advance it.
List<Boolean> persisted = new ArrayList<>(entries.size());
int updateIdx = 0;
for (RpcQueueEntry entry : entries) {
persisted.add(entry.insert() || updateCounts[updateIdx++] > 0);
}
return persisted;
}
private void batch(String sql, List<RpcEntity> entities, ColumnBinder binder) {
jdbcTemplate.batchUpdate(sql, new BatchPreparedStatementSetter() {
private int[] batch(String sql, List<RpcEntity> entities, ColumnBinder binder) {
return jdbcTemplate.batchUpdate(sql, new BatchPreparedStatementSetter() {
@Override
public void setValues(PreparedStatement ps, int i) throws SQLException {
binder.bind(ps, entities.get(i));

47
dao/src/test/java/org/thingsboard/server/dao/sql/rpc/JpaRpcDaoTest.java

@ -63,31 +63,19 @@ public class JpaRpcDaoTest extends AbstractJpaDaoTest {
@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);
DeviceId deviceId = new DeviceId(UUID.randomUUID());
rpcDao.createAsync(rpc.getTenantId(), rpc).get(5, TimeUnit.SECONDS);
// A create always persists (INSERT ... ON CONFLICT), so the future resolves true.
assertThat(rpcDao.createAsync(rpc(id, deviceId, RpcStatus.QUEUED, null)).get(5, TimeUnit.SECONDS)).isTrue();
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.updateAsync(update.getTenantId(), update).get(5, TimeUnit.SECONDS);
// The update matches the existing row, so the future resolves true.
assertThat(rpcDao.updateAsync(rpc(id, deviceId, RpcStatus.DELIVERED, JacksonUtil.toJsonNode("{\"ok\":true}")))
.get(5, TimeUnit.SECONDS)).isTrue();
Rpc afterUpdate = rpcDao.findById(TenantId.SYS_TENANT_ID, id);
assertThat(afterUpdate.getStatus()).isEqualTo(RpcStatus.DELIVERED);
@ -121,8 +109,8 @@ public class JpaRpcDaoTest extends AbstractJpaDaoTest {
// Enqueue both writes for the same rpcId back-to-back so they coalesce into one flush batch.
// Same rpcId -> same partition; the queue's stable sort must keep submission order
// (QUEUED before DELIVERED), so the final persisted row must be DELIVERED, never QUEUED.
var queuedFuture = rpcDao.createAsync(TenantId.SYS_TENANT_ID, queued);
var deliveredFuture = rpcDao.updateAsync(TenantId.SYS_TENANT_ID, delivered);
var queuedFuture = rpcDao.createAsync(queued);
var deliveredFuture = rpcDao.updateAsync(delivered);
queuedFuture.get(5, TimeUnit.SECONDS);
deliveredFuture.get(5, TimeUnit.SECONDS);
@ -138,12 +126,12 @@ public class JpaRpcDaoTest extends AbstractJpaDaoTest {
DeviceId deviceId = new DeviceId(UUID.randomUUID());
// Initial create.
rpcDao.createAsync(TenantId.SYS_TENANT_ID, rpc(id, deviceId, RpcStatus.QUEUED, null)).get(5, TimeUnit.SECONDS);
rpcDao.createAsync(rpc(id, deviceId, RpcStatus.QUEUED, null)).get(5, TimeUnit.SECONDS);
// Delivery timeout with closeTransportSessionOnRpcDeliveryTimeout=true re-queues the RPC: the
// device actor persists status=QUEUED again as a status update so init() can re-pick it up.
// The update must land on the existing row, never be dropped.
rpcDao.updateAsync(TenantId.SYS_TENANT_ID, rpc(id, deviceId, RpcStatus.QUEUED, null)).get(5, TimeUnit.SECONDS);
rpcDao.updateAsync(rpc(id, deviceId, RpcStatus.QUEUED, null)).get(5, TimeUnit.SECONDS);
Rpc stored = rpcDao.findById(TenantId.SYS_TENANT_ID, id);
assertThat(stored).isNotNull();
@ -155,14 +143,14 @@ public class JpaRpcDaoTest extends AbstractJpaDaoTest {
UUID id = UUID.randomUUID();
DeviceId deviceId = new DeviceId(UUID.randomUUID());
rpcDao.createAsync(TenantId.SYS_TENANT_ID, rpc(id, deviceId, RpcStatus.QUEUED, null)).get(5, TimeUnit.SECONDS);
rpcDao.createAsync(rpc(id, deviceId, RpcStatus.QUEUED, null)).get(5, TimeUnit.SECONDS);
// A successful response is stored.
rpcDao.updateAsync(TenantId.SYS_TENANT_ID, rpc(id, deviceId, RpcStatus.SUCCESSFUL, JacksonUtil.toJsonNode("{\"ok\":true}")))
rpcDao.updateAsync(rpc(id, deviceId, RpcStatus.SUCCESSFUL, JacksonUtil.toJsonNode("{\"ok\":true}")))
.get(5, TimeUnit.SECONDS);
// A later status update carries no response - it must NOT clobber the stored one.
rpcDao.updateAsync(TenantId.SYS_TENANT_ID, rpc(id, deviceId, RpcStatus.EXPIRED, null)).get(5, TimeUnit.SECONDS);
rpcDao.updateAsync(rpc(id, deviceId, RpcStatus.EXPIRED, null)).get(5, TimeUnit.SECONDS);
Rpc stored = rpcDao.findById(TenantId.SYS_TENANT_ID, id);
assertThat(stored.getStatus()).isEqualTo(RpcStatus.EXPIRED);
@ -174,15 +162,16 @@ public class JpaRpcDaoTest extends AbstractJpaDaoTest {
UUID id = UUID.randomUUID();
DeviceId deviceId = new DeviceId(UUID.randomUUID());
rpcDao.createAsync(TenantId.SYS_TENANT_ID, rpc(id, deviceId, RpcStatus.QUEUED, null)).get(5, TimeUnit.SECONDS);
rpcDao.createAsync(rpc(id, deviceId, RpcStatus.QUEUED, null)).get(5, TimeUnit.SECONDS);
// RPC is removed (TTL cleanup / manual delete) while a response is still in flight.
rpcDao.removeById(TenantId.SYS_TENANT_ID, id);
assertThat(rpcDao.findById(TenantId.SYS_TENANT_ID, id)).isNull();
// A late status update must not re-create the deleted row.
rpcDao.updateAsync(TenantId.SYS_TENANT_ID, rpc(id, deviceId, RpcStatus.SUCCESSFUL, JacksonUtil.toJsonNode("{\"ok\":true}")))
.get(5, TimeUnit.SECONDS);
// A late status update must not re-create the deleted row, and the future must resolve false
// (no row matched) so the service layer can skip the rule-engine notification.
assertThat(rpcDao.updateAsync(rpc(id, deviceId, RpcStatus.SUCCESSFUL, JacksonUtil.toJsonNode("{\"ok\":true}")))
.get(5, TimeUnit.SECONDS)).isFalse();
assertThat(rpcDao.findById(TenantId.SYS_TENANT_ID, id)).isNull();
}

Loading…
Cancel
Save