Browse Source

Implement events batch insert. Remove 'region' filter from tenant find queries.

pull/6147/head
Igor Kulikov 5 years ago
parent
commit
c3dde26c69
  1. 14
      application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java
  2. 2
      application/src/main/java/org/thingsboard/server/actors/stats/StatsActor.java
  3. 5
      application/src/main/resources/thingsboard.yml
  4. 6
      application/src/test/java/org/thingsboard/server/actors/stats/StatsActorTest.java
  5. 6
      common/dao-api/src/main/java/org/thingsboard/server/dao/event/EventService.java
  6. 18
      dao/src/main/java/org/thingsboard/server/dao/event/BaseEventService.java
  7. 18
      dao/src/main/java/org/thingsboard/server/dao/event/EventDao.java
  8. 91
      dao/src/main/java/org/thingsboard/server/dao/sql/event/AbstractEventInsertRepository.java
  9. 72
      dao/src/main/java/org/thingsboard/server/dao/sql/event/EventInsertRepository.java
  10. 63
      dao/src/main/java/org/thingsboard/server/dao/sql/event/HsqlEventInsertRepository.java
  11. 114
      dao/src/main/java/org/thingsboard/server/dao/sql/event/JpaBaseEventDao.java
  12. 53
      dao/src/main/java/org/thingsboard/server/dao/sql/event/PsqlEventInsertRepository.java
  13. 10
      dao/src/main/java/org/thingsboard/server/dao/sql/tenant/JpaTenantDao.java
  14. 14
      dao/src/main/java/org/thingsboard/server/dao/sql/tenant/TenantRepository.java
  15. 7
      dao/src/main/java/org/thingsboard/server/dao/tenant/TenantDao.java
  16. 15
      dao/src/main/java/org/thingsboard/server/dao/tenant/TenantServiceImpl.java
  17. 33
      dao/src/test/java/org/thingsboard/server/dao/service/event/BaseEventServiceTest.java
  18. 30
      dao/src/test/java/org/thingsboard/server/dao/sql/event/JpaBaseEventDaoTest.java
  19. 16
      dao/src/test/java/org/thingsboard/server/dao/sql/tenant/JpaTenantDaoTest.java

14
application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java

@ -467,7 +467,7 @@ public class ActorSystemContext {
}
private void persistEvent(Event event) {
eventService.save(event);
eventService.saveAsync(event);
}
private String toString(Throwable e) {
@ -552,10 +552,10 @@ public class ActorSystemContext {
}
event.setBody(node);
ListenableFuture<Event> future = eventService.saveAsync(event);
Futures.addCallback(future, new FutureCallback<Event>() {
ListenableFuture<Void> future = eventService.saveAsync(event);
Futures.addCallback(future, new FutureCallback<Void>() {
@Override
public void onSuccess(@Nullable Event event) {
public void onSuccess(@Nullable Void event) {
}
@ -605,10 +605,10 @@ public class ActorSystemContext {
}
event.setBody(node);
ListenableFuture<Event> future = eventService.saveAsync(event);
Futures.addCallback(future, new FutureCallback<Event>() {
ListenableFuture<Void> future = eventService.saveAsync(event);
Futures.addCallback(future, new FutureCallback<Void>() {
@Override
public void onSuccess(@Nullable Event event) {
public void onSuccess(@Nullable Void event) {
}

2
application/src/main/java/org/thingsboard/server/actors/stats/StatsActor.java

@ -59,7 +59,7 @@ public class StatsActor extends ContextAwareActor {
event.setTenantId(msg.getTenantId());
event.setType(DataConstants.STATS);
event.setBody(toBodyJson(systemContext.getServiceInfoProvider().getServiceId(), msg.getMessagesProcessed(), msg.getErrorsOccurred()));
systemContext.getEventService().save(event);
systemContext.getEventService().saveAsync(event);
}
private JsonNode toBodyJson(String serviceId, long messagesProcessed, long errorsOccurred) {

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

@ -278,6 +278,11 @@ sql:
stats_print_interval_ms: "${SQL_TS_LATEST_BATCH_STATS_PRINT_MS:10000}"
batch_threads: "${SQL_TS_LATEST_BATCH_THREADS:3}" # batch thread count have 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}"
events:
batch_size: "${SQL_EVENTS_BATCH_SIZE:10000}"
batch_max_delay: "${SQL_EVENTS_BATCH_MAX_DELAY_MS:100}"
stats_print_interval_ms: "${SQL_EVENTS_BATCH_STATS_PRINT_MS:10000}"
batch_threads: "${SQL_EVENTS_BATCH_THREADS:3}" # batch thread count have to be a prime number like 3 or 5 to gain perfect hash distribution
# Specify whether to sort entities before batch update. Should be enabled for cluster mode to avoid deadlocks
batch_sort: "${SQL_BATCH_SORT:false}"
# Specify whether to remove null characters from strValue of attributes and timeseries before insert

6
application/src/test/java/org/thingsboard/server/actors/stats/StatsActorTest.java

@ -59,11 +59,11 @@ class StatsActorTest {
@Test
void givenNonEmptyStatMessage_whenOnStatsPersistMsg_thenNoAction() {
statsActor.onStatsPersistMsg(new StatsPersistMsg(0, 1, TenantId.SYS_TENANT_ID, TenantId.SYS_TENANT_ID));
verify(eventService, times(1)).save(any(Event.class));
verify(eventService, times(1)).saveAsync(any(Event.class));
statsActor.onStatsPersistMsg(new StatsPersistMsg(1, 0, TenantId.SYS_TENANT_ID, TenantId.SYS_TENANT_ID));
verify(eventService, times(2)).save(any(Event.class));
verify(eventService, times(2)).saveAsync(any(Event.class));
statsActor.onStatsPersistMsg(new StatsPersistMsg(1, 1, TenantId.SYS_TENANT_ID, TenantId.SYS_TENANT_ID));
verify(eventService, times(3)).save(any(Event.class));
verify(eventService, times(3)).saveAsync(any(Event.class));
}
}

6
common/dao-api/src/main/java/org/thingsboard/server/dao/event/EventService.java

@ -28,11 +28,7 @@ import java.util.Optional;
public interface EventService {
Event save(Event event);
ListenableFuture<Event> saveAsync(Event event);
Optional<Event> saveIfNotExists(Event event);
ListenableFuture<Void> saveAsync(Event event);
Optional<Event> findEvent(TenantId tenantId, EntityId entityId, String eventType, String eventUid);

18
dao/src/main/java/org/thingsboard/server/dao/event/BaseEventService.java

@ -47,28 +47,12 @@ public class BaseEventService implements EventService {
public EventDao eventDao;
@Override
public Event save(Event event) {
eventValidator.validate(event, Event::getTenantId);
return eventDao.save(event.getTenantId(), event);
}
@Override
public ListenableFuture<Event> saveAsync(Event event) {
public ListenableFuture<Void> saveAsync(Event event) {
eventValidator.validate(event, Event::getTenantId);
checkAndTruncateDebugEvent(event);
return eventDao.saveAsync(event);
}
@Override
public Optional<Event> saveIfNotExists(Event event) {
eventValidator.validate(event, Event::getTenantId);
if (StringUtils.isEmpty(event.getUid())) {
throw new DataValidationException("Event uid should be specified!.");
}
checkAndTruncateDebugEvent(event);
return eventDao.saveIfNotExists(event);
}
private void checkAndTruncateDebugEvent(Event event) {
if (event.getType().startsWith("DEBUG") && event.getBody() != null && event.getBody().has("data")) {
String dataStr = event.getBody().get("data").asText();

18
dao/src/main/java/org/thingsboard/server/dao/event/EventDao.java

@ -33,29 +33,13 @@ import java.util.UUID;
*/
public interface EventDao extends Dao<Event> {
/**
* Save or update event object
*
* @param event the event object
* @return saved event object
*/
Event save(TenantId tenantId, Event event);
/**
* Save or update event object async
*
* @param event the event object
* @return saved event object future
*/
ListenableFuture<Event> saveAsync(Event event);
/**
* Save event object if it is not yet saved
*
* @param event the event object
* @return saved event object
*/
Optional<Event> saveIfNotExists(Event event);
ListenableFuture<Void> saveAsync(Event event);
/**
* Find event by tenantId, entityId and eventUid.

91
dao/src/main/java/org/thingsboard/server/dao/sql/event/AbstractEventInsertRepository.java

@ -1,91 +0,0 @@
/**
* Copyright © 2016-2022 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.event;
import lombok.extern.slf4j.Slf4j;
import org.hibernate.exception.ConstraintViolationException;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.data.jpa.repository.Modifying;
import org.springframework.transaction.PlatformTransactionManager;
import org.springframework.transaction.TransactionDefinition;
import org.springframework.transaction.TransactionStatus;
import org.springframework.transaction.support.DefaultTransactionDefinition;
import org.thingsboard.server.dao.model.sql.EventEntity;
import javax.persistence.EntityManager;
import javax.persistence.PersistenceContext;
import javax.persistence.Query;
@Slf4j
public abstract class AbstractEventInsertRepository implements EventInsertRepository {
@PersistenceContext
protected EntityManager entityManager;
@Autowired
protected PlatformTransactionManager transactionManager;
protected EventEntity saveAndGet(EventEntity entity, String insertOrUpdateOnPrimaryKeyConflict, String insertOrUpdateOnUniqueKeyConflict) {
EventEntity eventEntity = null;
TransactionStatus insertTransaction = getTransactionStatus(TransactionDefinition.PROPAGATION_REQUIRED);
try {
eventEntity = processSaveOrUpdate(entity, insertOrUpdateOnPrimaryKeyConflict);
transactionManager.commit(insertTransaction);
} catch (Throwable throwable) {
transactionManager.rollback(insertTransaction);
if (throwable.getCause() instanceof ConstraintViolationException) {
log.trace("Insert request leaded in a violation of a defined integrity constraint {} for Entity with entityId {} and entityType {}", throwable.getMessage(), entity.getEventUid(), entity.getEventType());
TransactionStatus transaction = getTransactionStatus(TransactionDefinition.PROPAGATION_REQUIRES_NEW);
try {
eventEntity = processSaveOrUpdate(entity, insertOrUpdateOnUniqueKeyConflict);
transactionManager.commit(transaction);
} catch (Throwable th) {
log.trace("Could not execute the update statement for Entity with entityId {} and entityType {}", entity.getEventUid(), entity.getEventType());
transactionManager.rollback(transaction);
}
} else {
log.trace("Could not execute the insert statement for Entity with entityId {} and entityType {}", entity.getEventUid(), entity.getEventType());
}
}
return eventEntity;
}
@Modifying
protected abstract EventEntity doProcessSaveOrUpdate(EventEntity entity, String query);
protected Query getQuery(EventEntity entity, String query) {
return entityManager.createNativeQuery(query, EventEntity.class)
.setParameter("id", entity.getUuid())
.setParameter("created_time", entity.getCreatedTime())
.setParameter("body", entity.getBody().toString())
.setParameter("entity_id", entity.getEntityId())
.setParameter("entity_type", entity.getEntityType().name())
.setParameter("event_type", entity.getEventType())
.setParameter("event_uid", entity.getEventUid())
.setParameter("tenant_id", entity.getTenantId())
.setParameter("ts", entity.getTs());
}
private EventEntity processSaveOrUpdate(EventEntity entity, String query) {
return doProcessSaveOrUpdate(entity, query);
}
private TransactionStatus getTransactionStatus(int propagationRequired) {
DefaultTransactionDefinition insertDefinition = new DefaultTransactionDefinition();
insertDefinition.setPropagationBehavior(propagationRequired);
return transactionManager.getTransaction(insertDefinition);
}
}

72
dao/src/main/java/org/thingsboard/server/dao/sql/event/EventInsertRepository.java

@ -15,10 +15,78 @@
*/
package org.thingsboard.server.dao.sql.event;
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.transaction.TransactionStatus;
import org.springframework.transaction.annotation.Transactional;
import org.springframework.transaction.support.TransactionCallbackWithoutResult;
import org.springframework.transaction.support.TransactionTemplate;
import org.thingsboard.server.dao.model.sql.EventEntity;
import org.thingsboard.server.dao.util.PsqlDao;
public interface EventInsertRepository {
import java.sql.PreparedStatement;
import java.sql.SQLException;
import java.util.List;
import java.util.regex.Pattern;
EventEntity saveOrUpdate(EventEntity entity);
@PsqlDao
@Repository
@Transactional
public class EventInsertRepository {
private static final ThreadLocal<Pattern> PATTERN_THREAD_LOCAL = ThreadLocal.withInitial(() -> Pattern.compile(String.valueOf(Character.MIN_VALUE)));
private static final String EMPTY_STR = "";
private static final String INSERT =
"INSERT INTO event (id, created_time, body, entity_id, entity_type, event_type, event_uid, tenant_id, ts) " +
"VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?) " +
"ON CONFLICT DO NOTHING;";
@Autowired
protected JdbcTemplate jdbcTemplate;
@Autowired
private TransactionTemplate transactionTemplate;
@Value("${sql.remove_null_chars:true}")
private boolean removeNullChars;
protected void save(List<EventEntity> entities) {
transactionTemplate.execute(new TransactionCallbackWithoutResult() {
@Override
protected void doInTransactionWithoutResult(TransactionStatus status) {
jdbcTemplate.batchUpdate(INSERT, new BatchPreparedStatementSetter() {
@Override
public void setValues(PreparedStatement ps, int i) throws SQLException {
EventEntity event = entities.get(i);
ps.setObject(1, event.getId());
ps.setLong(2, event.getCreatedTime());
ps.setString(3, replaceNullChars(event.getBody().toString()));
ps.setObject(4, event.getEntityId());
ps.setString(5, event.getEntityType().name());
ps.setString(6, event.getEventType());
ps.setString(7, event.getEventUid());
ps.setObject(8, event.getTenantId());
ps.setLong(9, event.getTs());
}
@Override
public int getBatchSize() {
return entities.size();
}
});
}
});
}
private String replaceNullChars(String strValue) {
if (removeNullChars && strValue != null) {
return PATTERN_THREAD_LOCAL.get().matcher(strValue).replaceAll(EMPTY_STR);
}
return strValue;
}
}

63
dao/src/main/java/org/thingsboard/server/dao/sql/event/HsqlEventInsertRepository.java

@ -1,63 +0,0 @@
/**
* Copyright © 2016-2022 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.event;
import org.springframework.stereotype.Repository;
import org.thingsboard.server.dao.model.sql.EventEntity;
import org.thingsboard.server.dao.util.HsqlDao;
import javax.persistence.Query;
@HsqlDao
@Repository
public class HsqlEventInsertRepository extends AbstractEventInsertRepository {
private static final String P_KEY_CONFLICT_STATEMENT = "(event.id=I.id)";
private static final String UNQ_KEY_CONFLICT_STATEMENT = "(event.tenant_id=I.tenant_id AND event.entity_type=I.entity_type AND event.entity_id=I.entity_id AND event.event_type=I.event_type AND event.event_uid=I.event_uid)";
private static final String INSERT_OR_UPDATE_ON_P_KEY_CONFLICT = getInsertString(P_KEY_CONFLICT_STATEMENT);
private static final String INSERT_OR_UPDATE_ON_UNQ_KEY_CONFLICT = getInsertString(UNQ_KEY_CONFLICT_STATEMENT);
@Override
public EventEntity saveOrUpdate(EventEntity entity) {
return saveAndGet(entity, INSERT_OR_UPDATE_ON_P_KEY_CONFLICT, INSERT_OR_UPDATE_ON_UNQ_KEY_CONFLICT);
}
@Override
protected EventEntity doProcessSaveOrUpdate(EventEntity entity, String query) {
getQuery(entity, query).executeUpdate();
return entityManager.find(EventEntity.class, entity.getUuid());
}
protected Query getQuery(EventEntity entity, String query) {
return entityManager.createNativeQuery(query, EventEntity.class)
.setParameter("id", entity.getUuid().toString())
.setParameter("created_time", entity.getCreatedTime())
.setParameter("body", entity.getBody().toString())
.setParameter("entity_id", entity.getEntityId().toString())
.setParameter("entity_type", entity.getEntityType().name())
.setParameter("event_type", entity.getEventType())
.setParameter("event_uid", entity.getEventUid())
.setParameter("tenant_id", entity.getTenantId().toString())
.setParameter("ts", entity.getTs());
}
private static String getInsertString(String conflictStatement) {
return "MERGE INTO event USING (VALUES UUID(:id), :created_time, :body, UUID(:entity_id), :entity_type, :event_type, :event_uid, UUID(:tenant_id), :ts) I (id, created_time, body, entity_id, entity_type, event_type, event_uid, tenant_id, ts) ON " + conflictStatement
+ " WHEN MATCHED THEN UPDATE SET event.id = I.id, event.created_time = I.created_time, event.body = I.body, event.entity_id = I.entity_id, event.entity_type = I.entity_type, event.event_type = I.event_type, event.event_uid = I.event_uid, event.tenant_id = I.tenant_id, event.ts = I.ts" +
" WHEN NOT MATCHED THEN INSERT (id, created_time, body, entity_id, entity_type, event_type, event_uid, tenant_id, ts) VALUES (I.id, I.created_time, I.body, I.entity_id, I.entity_type, I.event_type, I.event_uid, I.tenant_id, I.ts)";
}
}

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

@ -20,6 +20,7 @@ import com.google.common.util.concurrent.ListenableFuture;
import lombok.extern.slf4j.Slf4j;
import org.apache.commons.lang3.StringUtils;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.data.domain.PageRequest;
import org.springframework.data.repository.CrudRepository;
import org.springframework.stereotype.Component;
@ -34,15 +35,24 @@ import org.thingsboard.server.common.data.id.EventId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.page.PageData;
import org.thingsboard.server.common.data.page.TimePageLink;
import org.thingsboard.server.common.stats.StatsFactory;
import org.thingsboard.server.dao.DaoUtil;
import org.thingsboard.server.dao.event.EventDao;
import org.thingsboard.server.dao.model.sql.AttributeKvEntity;
import org.thingsboard.server.dao.model.sql.EventEntity;
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 javax.annotation.PostConstruct;
import javax.annotation.PreDestroy;
import java.util.Comparator;
import java.util.List;
import java.util.Objects;
import java.util.Optional;
import java.util.UUID;
import java.util.function.Function;
import static org.thingsboard.server.dao.model.ModelConstants.NULL_UUID;
@ -74,8 +84,55 @@ public class JpaBaseEventDao extends JpaAbstractDao<EventEntity, Event> implemen
return eventRepository;
}
@Autowired
ScheduledLogExecutorComponent logExecutor;
@Autowired
private StatsFactory statsFactory;
@Value("${sql.events.batch_size:10000}")
private int batchSize;
@Value("${sql.events.batch_max_delay:100}")
private long maxDelay;
@Value("${sql.events.stats_print_interval_ms:10000}")
private long statsPrintIntervalMs;
@Value("${sql.events.batch_threads:3}")
private int batchThreads;
@Value("${sql.batch_sort:false}")
private boolean batchSortEnabled;
private TbSqlBlockingQueueWrapper<EventEntity> queue;
@PostConstruct
private void init() {
TbSqlBlockingQueueParams params = TbSqlBlockingQueueParams.builder()
.logName("Events")
.batchSize(batchSize)
.maxDelay(maxDelay)
.statsPrintIntervalMs(statsPrintIntervalMs)
.statsNamePrefix("events")
.batchSortEnabled(batchSortEnabled)
.build();
Function<EventEntity, Integer> hashcodeFunction = entity -> entity.getEntityId().hashCode();
queue = new TbSqlBlockingQueueWrapper<>(params, hashcodeFunction, batchThreads, statsFactory);
queue.init(logExecutor, v -> eventInsertRepository.save(v),
Comparator.comparing((EventEntity eventEntity) -> eventEntity.getTs())
);
}
@PreDestroy
private void destroy() {
if (queue != null) {
queue.destroy();
}
}
@Override
public Event save(TenantId tenantId, Event event) {
public ListenableFuture<Void> saveAsync(Event event) {
log.debug("Save event [{}] ", event);
if (event.getId() == null) {
UUID timeBased = Uuids.timeBased();
@ -92,33 +149,27 @@ public class JpaBaseEventDao extends JpaAbstractDao<EventEntity, Event> implemen
if (StringUtils.isEmpty(event.getUid())) {
event.setUid(event.getId().toString());
}
return save(new EventEntity(event), false).orElse(null);
return save(new EventEntity(event));
}
@Override
public ListenableFuture<Event> saveAsync(Event event) {
log.debug("Save event [{}] ", event);
if (event.getId() == null) {
UUID timeBased = Uuids.timeBased();
event.setId(new EventId(timeBased));
event.setCreatedTime(Uuids.unixTimestamp(timeBased));
} else if (event.getCreatedTime() == 0L) {
UUID eventId = event.getId().getId();
if (eventId.version() == 1) {
event.setCreatedTime(Uuids.unixTimestamp(eventId));
} else {
event.setCreatedTime(System.currentTimeMillis());
}
private ListenableFuture<Void> save(EventEntity entity) {
log.debug("Save event [{}] ", entity);
if (entity.getTenantId() == null) {
log.trace("Save system event with predefined id {}", systemTenantId);
entity.setTenantId(systemTenantId);
}
if (StringUtils.isEmpty(event.getUid())) {
event.setUid(event.getId().toString());
if (entity.getUuid() == null) {
entity.setUuid(Uuids.timeBased());
}
return service.submit(() -> save(new EventEntity(event), false).orElse(null));
if (StringUtils.isEmpty(entity.getEventUid())) {
entity.setEventUid(entity.getUuid().toString());
}
return addToQueue(entity);
}
@Override
public Optional<Event> saveIfNotExists(Event event) {
return save(new EventEntity(event), true);
private ListenableFuture<Void> addToQueue(EventEntity entity) {
return queue.add(entity);
}
@Override
@ -264,25 +315,6 @@ public class JpaBaseEventDao extends JpaAbstractDao<EventEntity, Event> implemen
eventCleanupRepository.cleanupEvents(regularEventStartTs, regularEventEndTs, debugEventStartTs, debugEventEndTs);
}
public Optional<Event> save(EventEntity entity, boolean ifNotExists) {
log.debug("Save event [{}] ", entity);
if (entity.getTenantId() == null) {
log.trace("Save system event with predefined id {}", systemTenantId);
entity.setTenantId(systemTenantId);
}
if (entity.getUuid() == null) {
entity.setUuid(Uuids.timeBased());
}
if (StringUtils.isEmpty(entity.getEventUid())) {
entity.setEventUid(entity.getUuid().toString());
}
if (ifNotExists &&
eventRepository.findByTenantIdAndEntityTypeAndEntityId(entity.getTenantId(), entity.getEntityType(), entity.getEntityId()) != null) {
return Optional.empty();
}
return Optional.of(DaoUtil.getData(eventInsertRepository.saveOrUpdate(entity)));
}
private long notNull(Long value) {
return value != null ? value : 0;
}

53
dao/src/main/java/org/thingsboard/server/dao/sql/event/PsqlEventInsertRepository.java

@ -1,53 +0,0 @@
/**
* Copyright © 2016-2022 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.event;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Repository;
import org.thingsboard.server.dao.model.sql.EventEntity;
import org.thingsboard.server.dao.util.PsqlDao;
@Slf4j
@PsqlDao
@Repository
public class PsqlEventInsertRepository extends AbstractEventInsertRepository {
private static final String P_KEY_CONFLICT_STATEMENT = "(id)";
private static final String UNQ_KEY_CONFLICT_STATEMENT = "(tenant_id, created_time, entity_type, entity_id, event_type, event_uid)";
private static final String UPDATE_P_KEY_STATEMENT = "id = :id";
private static final String UPDATE_UNQ_KEY_STATEMENT = "created_time = :created_time, tenant_id = :tenant_id, entity_type = :entity_type, entity_id = :entity_id, event_type = :event_type, event_uid = :event_uid";
private static final String INSERT_OR_UPDATE_ON_P_KEY_CONFLICT = getInsertOrUpdateString(P_KEY_CONFLICT_STATEMENT, UPDATE_UNQ_KEY_STATEMENT);
private static final String INSERT_OR_UPDATE_ON_UNQ_KEY_CONFLICT = getInsertOrUpdateString(UNQ_KEY_CONFLICT_STATEMENT, UPDATE_P_KEY_STATEMENT);
@Override
public EventEntity saveOrUpdate(EventEntity entity) {
return saveAndGet(entity, INSERT_OR_UPDATE_ON_P_KEY_CONFLICT, INSERT_OR_UPDATE_ON_UNQ_KEY_CONFLICT);
}
@Override
protected EventEntity doProcessSaveOrUpdate(EventEntity entity, String query) {
return (EventEntity) getQuery(entity, query).getSingleResult();
}
private static String getInsertOrUpdateString(String eventKeyStatement, String updateKeyStatement) {
return "INSERT INTO event (id, created_time, body, entity_id, entity_type, event_type, event_uid, tenant_id, ts) " +
"VALUES (:id, :created_time, :body, :entity_id, :entity_type, :event_type, :event_uid, :tenant_id, :ts) " +
"ON CONFLICT " + eventKeyStatement + " DO UPDATE SET body = :body, ts = :ts," + updateKeyStatement + " returning *";
}
}

10
dao/src/main/java/org/thingsboard/server/dao/sql/tenant/JpaTenantDao.java

@ -58,19 +58,17 @@ public class JpaTenantDao extends JpaAbstractSearchTextDao<TenantEntity, Tenant>
}
@Override
public PageData<Tenant> findTenantsByRegion(TenantId tenantId, String region, PageLink pageLink) {
public PageData<Tenant> findTenants(TenantId tenantId, PageLink pageLink) {
return DaoUtil.toPageData(tenantRepository
.findByRegionNextPage(
region,
.findTenantsNextPage(
Objects.toString(pageLink.getTextSearch(), ""),
DaoUtil.toPageable(pageLink)));
}
@Override
public PageData<TenantInfo> findTenantInfosByRegion(TenantId tenantId, String region, PageLink pageLink) {
public PageData<TenantInfo> findTenantInfos(TenantId tenantId, PageLink pageLink) {
return DaoUtil.toPageData(tenantRepository
.findTenantInfoByRegionNextPage(
region,
.findTenantInfosNextPage(
Objects.toString(pageLink.getTextSearch(), ""),
DaoUtil.toPageable(pageLink, TenantInfoEntity.tenantInfoColumnMap)));
}

14
dao/src/main/java/org/thingsboard/server/dao/sql/tenant/TenantRepository.java

@ -36,19 +36,15 @@ public interface TenantRepository extends PagingAndSortingRepository<TenantEntit
"WHERE t.id = :tenantId")
TenantInfoEntity findTenantInfoById(@Param("tenantId") UUID tenantId);
@Query("SELECT t FROM TenantEntity t WHERE t.region = :region " +
"AND LOWER(t.searchText) LIKE LOWER(CONCAT('%', :textSearch, '%'))")
Page<TenantEntity> findByRegionNextPage(@Param("region") String region,
@Param("textSearch") String textSearch,
Pageable pageable);
@Query("SELECT t FROM TenantEntity t WHERE LOWER(t.searchText) LIKE LOWER(CONCAT('%', :textSearch, '%'))")
Page<TenantEntity> findTenantsNextPage(@Param("textSearch") String textSearch,
Pageable pageable);
@Query("SELECT new org.thingsboard.server.dao.model.sql.TenantInfoEntity(t, p.name) " +
"FROM TenantEntity t " +
"LEFT JOIN TenantProfileEntity p on p.id = t.tenantProfileId " +
"WHERE t.region = :region " +
"AND LOWER(t.searchText) LIKE LOWER(CONCAT('%', :textSearch, '%'))")
Page<TenantInfoEntity> findTenantInfoByRegionNextPage(@Param("region") String region,
@Param("textSearch") String textSearch,
"WHERE LOWER(t.searchText) LIKE LOWER(CONCAT('%', :textSearch, '%'))")
Page<TenantInfoEntity> findTenantInfosNextPage(@Param("textSearch") String textSearch,
Pageable pageable);
@Query("SELECT t.id FROM TenantEntity t")

7
dao/src/main/java/org/thingsboard/server/dao/tenant/TenantDao.java

@ -37,15 +37,14 @@ public interface TenantDao extends Dao<Tenant> {
Tenant save(TenantId tenantId, Tenant tenant);
/**
* Find tenants by region and page link.
* Find tenants by page link.
*
* @param region the region
* @param pageLink the page link
* @return the list of tenant objects
*/
PageData<Tenant> findTenantsByRegion(TenantId tenantId, String region, PageLink pageLink);
PageData<Tenant> findTenants(TenantId tenantId, PageLink pageLink);
PageData<TenantInfo> findTenantInfosByRegion(TenantId tenantId, String region, PageLink pageLink);
PageData<TenantInfo> findTenantInfos(TenantId tenantId, PageLink pageLink);
PageData<TenantId> findTenantsIds(PageLink pageLink);

15
dao/src/main/java/org/thingsboard/server/dao/tenant/TenantServiceImpl.java

@ -24,7 +24,6 @@ import org.springframework.stereotype.Service;
import org.thingsboard.server.common.data.Tenant;
import org.thingsboard.server.common.data.TenantInfo;
import org.thingsboard.server.common.data.TenantProfile;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.page.PageData;
import org.thingsboard.server.common.data.page.PageLink;
@ -168,20 +167,20 @@ public class TenantServiceImpl extends AbstractEntityService implements TenantSe
public PageData<Tenant> findTenants(PageLink pageLink) {
log.trace("Executing findTenants pageLink [{}]", pageLink);
Validator.validatePageLink(pageLink);
return tenantDao.findTenantsByRegion(TenantId.SYS_TENANT_ID, DEFAULT_TENANT_REGION, pageLink);
return tenantDao.findTenants(TenantId.SYS_TENANT_ID, pageLink);
}
@Override
public PageData<TenantInfo> findTenantInfos(PageLink pageLink) {
log.trace("Executing findTenantInfos pageLink [{}]", pageLink);
Validator.validatePageLink(pageLink);
return tenantDao.findTenantInfosByRegion(TenantId.SYS_TENANT_ID, DEFAULT_TENANT_REGION, pageLink);
return tenantDao.findTenantInfos(TenantId.SYS_TENANT_ID, pageLink);
}
@Override
public void deleteTenants() {
log.trace("Executing deleteTenants");
tenantsRemover.removeEntities(TenantId.SYS_TENANT_ID, DEFAULT_TENANT_REGION);
tenantsRemover.removeEntities(TenantId.SYS_TENANT_ID, TenantId.SYS_TENANT_ID);
}
private DataValidator<Tenant> tenantValidator =
@ -214,12 +213,12 @@ public class TenantServiceImpl extends AbstractEntityService implements TenantSe
}
};
private PaginatedRemover<String, Tenant> tenantsRemover =
new PaginatedRemover<String, Tenant>() {
private PaginatedRemover<TenantId, Tenant> tenantsRemover =
new PaginatedRemover<>() {
@Override
protected PageData<Tenant> findEntities(TenantId tenantId, String region, PageLink pageLink) {
return tenantDao.findTenantsByRegion(tenantId, region, pageLink);
protected PageData<Tenant> findEntities(TenantId tenantId, TenantId id, PageLink pageLink) {
return tenantDao.findTenants(tenantId, pageLink);
}
@Override

33
dao/src/test/java/org/thingsboard/server/dao/service/event/BaseEventServiceTest.java

@ -17,6 +17,7 @@ package org.thingsboard.server.dao.service.event;
import com.datastax.oss.driver.api.core.uuid.Uuids;
import org.junit.Assert;
import org.junit.Ignore;
import org.junit.Test;
import org.thingsboard.server.common.data.DataConstants;
import org.thingsboard.server.common.data.Event;
@ -35,6 +36,7 @@ import java.time.LocalDateTime;
import java.time.Month;
import java.time.ZoneOffset;
import java.util.Optional;
import java.util.concurrent.ExecutionException;
public abstract class BaseEventServiceTest extends AbstractServiceTest {
@ -42,21 +44,13 @@ public abstract class BaseEventServiceTest extends AbstractServiceTest {
public void saveEvent() throws Exception {
DeviceId devId = new DeviceId(Uuids.timeBased());
Event event = generateEvent(null, devId, "ALARM", Uuids.timeBased().toString());
Event saved = eventService.save(event);
eventService.saveAsync(event).get();
Optional<Event> loaded = eventService.findEvent(event.getTenantId(), event.getEntityId(), event.getType(), event.getUid());
Assert.assertTrue(loaded.isPresent());
Assert.assertNotNull(loaded.get());
Assert.assertEquals(saved, loaded.get());
}
@Test
public void saveEventIfNotExists() throws Exception {
DeviceId devId = new DeviceId(Uuids.timeBased());
Event event = generateEvent(null, devId, "ALARM", Uuids.timeBased().toString());
Optional<Event> saved = eventService.saveIfNotExists(event);
Assert.assertTrue(saved.isPresent());
saved = eventService.saveIfNotExists(event);
Assert.assertFalse(saved.isPresent());
Assert.assertEquals(event.getEntityId(), loaded.get().getEntityId());
Assert.assertEquals(event.getType(), loaded.get().getType());
Assert.assertEquals(event.getBody(), loaded.get().getBody());
}
@Test
@ -96,11 +90,11 @@ public abstract class BaseEventServiceTest extends AbstractServiceTest {
@Test
public void findEventsByTypeAndTimeDescOrder() throws Exception {
long timeBeforeStartTime = LocalDateTime.of(2016, Month.NOVEMBER, 1, 11, 30).toEpochSecond(ZoneOffset.UTC);
long startTime = LocalDateTime.of(2016, Month.NOVEMBER, 1, 12, 0).toEpochSecond(ZoneOffset.UTC);
long eventTime = LocalDateTime.of(2016, Month.NOVEMBER, 1, 12, 30).toEpochSecond(ZoneOffset.UTC);
long endTime = LocalDateTime.of(2016, Month.NOVEMBER, 1, 13, 0).toEpochSecond(ZoneOffset.UTC);
long timeAfterEndTime = LocalDateTime.of(2016, Month.NOVEMBER, 1, 13, 30).toEpochSecond(ZoneOffset.UTC);
long timeBeforeStartTime = LocalDateTime.of(2017, Month.NOVEMBER, 1, 11, 30).toEpochSecond(ZoneOffset.UTC);
long startTime = LocalDateTime.of(2017, Month.NOVEMBER, 1, 12, 0).toEpochSecond(ZoneOffset.UTC);
long eventTime = LocalDateTime.of(2017, Month.NOVEMBER, 1, 12, 30).toEpochSecond(ZoneOffset.UTC);
long endTime = LocalDateTime.of(2017, Month.NOVEMBER, 1, 13, 0).toEpochSecond(ZoneOffset.UTC);
long timeAfterEndTime = LocalDateTime.of(2017, Month.NOVEMBER, 1, 13, 30).toEpochSecond(ZoneOffset.UTC);
CustomerId customerId = new CustomerId(Uuids.timeBased());
TenantId tenantId = TenantId.fromUUID(Uuids.timeBased());
@ -129,9 +123,10 @@ public abstract class BaseEventServiceTest extends AbstractServiceTest {
Assert.assertFalse(events.hasNext());
}
private Event saveEventWithProvidedTime(long time, EntityId entityId, TenantId tenantId) throws IOException {
private Event saveEventWithProvidedTime(long time, EntityId entityId, TenantId tenantId) throws Exception {
Event event = generateEvent(tenantId, entityId, DataConstants.STATS, null);
event.setId(new EventId(Uuids.startOf(time)));
return eventService.save(event);
eventService.saveAsync(event).get();
return event;
}
}

30
dao/src/test/java/org/thingsboard/server/dao/sql/event/JpaBaseEventDaoTest.java

@ -37,6 +37,7 @@ import java.io.IOException;
import java.util.List;
import java.util.Optional;
import java.util.UUID;
import java.util.concurrent.ExecutionException;
import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertFalse;
@ -55,19 +56,6 @@ public class JpaBaseEventDaoTest extends AbstractJpaDaoTest {
@Autowired
private EventDao eventDao;
@Test
public void testSaveIfNotExists() {
UUID eventId = Uuids.timeBased();
UUID tenantId = Uuids.timeBased();
UUID entityId = Uuids.timeBased();
Event event = getEvent(eventId, tenantId, entityId);
Optional<Event> optEvent1 = eventDao.saveIfNotExists(event);
assertTrue("Optional is expected to be non-empty", optEvent1.isPresent());
assertEquals(event, optEvent1.get());
Optional<Event> optEvent2 = eventDao.saveIfNotExists(event);
assertFalse("Optional is expected to be empty", optEvent2.isPresent());
}
@Test
@DatabaseSetup("classpath:dbunit/event.xml")
public void findEvent() {
@ -82,7 +70,7 @@ public class JpaBaseEventDaoTest extends AbstractJpaDaoTest {
}
@Test
public void findEventsByEntityIdAndPageLink() {
public void findEventsByEntityIdAndPageLink() throws Exception {
UUID tenantId = Uuids.timeBased();
UUID entityId1 = Uuids.timeBased();
UUID entityId2 = Uuids.timeBased();
@ -116,7 +104,7 @@ public class JpaBaseEventDaoTest extends AbstractJpaDaoTest {
}
@Test
public void findEventsByEntityIdAndEventTypeAndPageLink() {
public void findEventsByEntityIdAndEventTypeAndPageLink() throws Exception {
UUID tenantId = Uuids.timeBased();
UUID entityId1 = Uuids.timeBased();
UUID entityId2 = Uuids.timeBased();
@ -144,27 +132,27 @@ public class JpaBaseEventDaoTest extends AbstractJpaDaoTest {
assertEquals(2, events5.getData().size());
}
private long createEventsTwoEntitiesTwoTypes(UUID tenantId, UUID entityId1, UUID entityId2, long startTime, int count) {
private long createEventsTwoEntitiesTwoTypes(UUID tenantId, UUID entityId1, UUID entityId2, long startTime, int count) throws Exception {
for (int i = 0; i < count / 2; i++) {
String type = i % 2 == 0 ? STATS : ALARM;
UUID eventId1 = Uuids.timeBased();
Event event1 = getEvent(eventId1, tenantId, entityId1, type);
eventDao.save(TenantId.fromUUID(tenantId), event1);
eventDao.saveAsync(event1).get();
UUID eventId2 = Uuids.timeBased();
Event event2 = getEvent(eventId2, tenantId, entityId2, type);
eventDao.save(TenantId.fromUUID(tenantId), event2);
eventDao.saveAsync(event2).get();
}
return System.currentTimeMillis();
}
private long createEventsTwoEntities(UUID tenantId, UUID entityId1, UUID entityId2, long startTime, int count) {
private long createEventsTwoEntities(UUID tenantId, UUID entityId1, UUID entityId2, long startTime, int count) throws Exception {
for (int i = 0; i < count / 2; i++) {
UUID eventId1 = Uuids.timeBased();
Event event1 = getEvent(eventId1, tenantId, entityId1);
eventDao.save(TenantId.fromUUID(tenantId), event1);
eventDao.saveAsync(event1).get();
UUID eventId2 = Uuids.timeBased();
Event event2 = getEvent(eventId2, tenantId, entityId2);
eventDao.save(TenantId.fromUUID(tenantId), event2);
eventDao.saveAsync(event2).get();
}
return System.currentTimeMillis();
}

16
dao/src/test/java/org/thingsboard/server/dao/sql/tenant/JpaTenantDaoTest.java

@ -66,36 +66,34 @@ public class JpaTenantDaoTest extends AbstractJpaDaoTest {
@Test
//@DatabaseSetup("classpath:dbunit/empty_dataset.xml")
public void testFindTenantsByRegion() {
public void testFindTenants() {
createTenants();
assertEquals(60, tenantDao.find(AbstractServiceTest.SYSTEM_TENANT_ID).size());
assertEquals(30, tenantDao.find(AbstractServiceTest.SYSTEM_TENANT_ID).size());
PageLink pageLink = new PageLink(20, 0, "title");
PageData<Tenant> tenants1 = tenantDao.findTenantsByRegion(AbstractServiceTest.SYSTEM_TENANT_ID, "REGION_1", pageLink);
PageData<Tenant> tenants1 = tenantDao.findTenants(AbstractServiceTest.SYSTEM_TENANT_ID, pageLink);
assertEquals(20, tenants1.getData().size());
pageLink = pageLink.nextPageLink();
PageData<Tenant> tenants2 = tenantDao.findTenantsByRegion(AbstractServiceTest.SYSTEM_TENANT_ID, "REGION_1",
PageData<Tenant> tenants2 = tenantDao.findTenants(AbstractServiceTest.SYSTEM_TENANT_ID,
pageLink);
assertEquals(10, tenants2.getData().size());
pageLink = pageLink.nextPageLink();
PageData<Tenant> tenants3 = tenantDao.findTenantsByRegion(AbstractServiceTest.SYSTEM_TENANT_ID, "REGION_1",
PageData<Tenant> tenants3 = tenantDao.findTenants(AbstractServiceTest.SYSTEM_TENANT_ID,
pageLink);
assertEquals(0, tenants3.getData().size());
}
private void createTenants() {
for (int i = 0; i < 30; i++) {
createTenant("REGION_1", "TITLE", i);
createTenant("REGION_2", "TITLE", i);
createTenant("TITLE", i);
}
}
void createTenant(String region, String title, int index) {
void createTenant(String title, int index) {
Tenant tenant = new Tenant();
tenant.setId(TenantId.fromUUID(Uuids.timeBased()));
tenant.setRegion(region);
tenant.setTitle(title + "_" + index);
tenant.setTenantProfileId(tenantProfile.getId());
createdTenants.add(tenantDao.save(TenantId.SYS_TENANT_ID, tenant));

Loading…
Cancel
Save