diff --git a/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java b/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java index 889bd616fc..1be741a8c1 100644 --- a/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java +++ b/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 future = eventService.saveAsync(event); - Futures.addCallback(future, new FutureCallback() { + ListenableFuture future = eventService.saveAsync(event); + Futures.addCallback(future, new FutureCallback() { @Override - public void onSuccess(@Nullable Event event) { + public void onSuccess(@Nullable Void event) { } @@ -605,10 +605,10 @@ public class ActorSystemContext { } event.setBody(node); - ListenableFuture future = eventService.saveAsync(event); - Futures.addCallback(future, new FutureCallback() { + ListenableFuture future = eventService.saveAsync(event); + Futures.addCallback(future, new FutureCallback() { @Override - public void onSuccess(@Nullable Event event) { + public void onSuccess(@Nullable Void event) { } diff --git a/application/src/main/java/org/thingsboard/server/actors/stats/StatsActor.java b/application/src/main/java/org/thingsboard/server/actors/stats/StatsActor.java index b547741453..8cdf5cec02 100644 --- a/application/src/main/java/org/thingsboard/server/actors/stats/StatsActor.java +++ b/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) { diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index eaa76d4273..5a37c9c3aa 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/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 diff --git a/application/src/test/java/org/thingsboard/server/actors/stats/StatsActorTest.java b/application/src/test/java/org/thingsboard/server/actors/stats/StatsActorTest.java index d617ed830b..0e928b7dd6 100644 --- a/application/src/test/java/org/thingsboard/server/actors/stats/StatsActorTest.java +++ b/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)); } } diff --git a/common/dao-api/src/main/java/org/thingsboard/server/dao/event/EventService.java b/common/dao-api/src/main/java/org/thingsboard/server/dao/event/EventService.java index f7ddfbc8de..83624df8e3 100644 --- a/common/dao-api/src/main/java/org/thingsboard/server/dao/event/EventService.java +++ b/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 saveAsync(Event event); - - Optional saveIfNotExists(Event event); + ListenableFuture saveAsync(Event event); Optional findEvent(TenantId tenantId, EntityId entityId, String eventType, String eventUid); diff --git a/dao/src/main/java/org/thingsboard/server/dao/event/BaseEventService.java b/dao/src/main/java/org/thingsboard/server/dao/event/BaseEventService.java index b80de363f1..714eefd0b1 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/event/BaseEventService.java +++ b/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 saveAsync(Event event) { + public ListenableFuture saveAsync(Event event) { eventValidator.validate(event, Event::getTenantId); checkAndTruncateDebugEvent(event); return eventDao.saveAsync(event); } - @Override - public Optional 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(); diff --git a/dao/src/main/java/org/thingsboard/server/dao/event/EventDao.java b/dao/src/main/java/org/thingsboard/server/dao/event/EventDao.java index e0f1160db1..f3574fb419 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/event/EventDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/event/EventDao.java @@ -33,29 +33,13 @@ import java.util.UUID; */ public interface EventDao extends Dao { - /** - * 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 saveAsync(Event event); - - /** - * Save event object if it is not yet saved - * - * @param event the event object - * @return saved event object - */ - Optional saveIfNotExists(Event event); + ListenableFuture saveAsync(Event event); /** * Find event by tenantId, entityId and eventUid. diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/event/AbstractEventInsertRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sql/event/AbstractEventInsertRepository.java deleted file mode 100644 index d3e631e887..0000000000 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/event/AbstractEventInsertRepository.java +++ /dev/null @@ -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); - } -} \ No newline at end of file diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/event/EventInsertRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sql/event/EventInsertRepository.java index 866a9b5483..1f2d4c01be 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/event/EventInsertRepository.java +++ b/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_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 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; + } } diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/event/HsqlEventInsertRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sql/event/HsqlEventInsertRepository.java deleted file mode 100644 index a85a3eaa12..0000000000 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/event/HsqlEventInsertRepository.java +++ /dev/null @@ -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)"; - } -} \ No newline at end of file diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/event/JpaBaseEventDao.java b/dao/src/main/java/org/thingsboard/server/dao/sql/event/JpaBaseEventDao.java index 8398e54dde..6f42857a96 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/event/JpaBaseEventDao.java +++ b/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 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 queue; + + @PostConstruct + private void init() { + TbSqlBlockingQueueParams params = TbSqlBlockingQueueParams.builder() + .logName("Events") + .batchSize(batchSize) + .maxDelay(maxDelay) + .statsPrintIntervalMs(statsPrintIntervalMs) + .statsNamePrefix("events") + .batchSortEnabled(batchSortEnabled) + .build(); + Function 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 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 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 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 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 saveIfNotExists(Event event) { - return save(new EventEntity(event), true); + private ListenableFuture addToQueue(EventEntity entity) { + return queue.add(entity); } @Override @@ -264,25 +315,6 @@ public class JpaBaseEventDao extends JpaAbstractDao implemen eventCleanupRepository.cleanupEvents(regularEventStartTs, regularEventEndTs, debugEventStartTs, debugEventEndTs); } - public Optional 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; } diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/event/PsqlEventInsertRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sql/event/PsqlEventInsertRepository.java deleted file mode 100644 index 3e7baa571c..0000000000 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/event/PsqlEventInsertRepository.java +++ /dev/null @@ -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 *"; - } -} \ No newline at end of file diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/tenant/JpaTenantDao.java b/dao/src/main/java/org/thingsboard/server/dao/sql/tenant/JpaTenantDao.java index e7a1961fab..2a3ce2bcab 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/tenant/JpaTenantDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/tenant/JpaTenantDao.java @@ -58,19 +58,17 @@ public class JpaTenantDao extends JpaAbstractSearchTextDao } @Override - public PageData findTenantsByRegion(TenantId tenantId, String region, PageLink pageLink) { + public PageData findTenants(TenantId tenantId, PageLink pageLink) { return DaoUtil.toPageData(tenantRepository - .findByRegionNextPage( - region, + .findTenantsNextPage( Objects.toString(pageLink.getTextSearch(), ""), DaoUtil.toPageable(pageLink))); } @Override - public PageData findTenantInfosByRegion(TenantId tenantId, String region, PageLink pageLink) { + public PageData findTenantInfos(TenantId tenantId, PageLink pageLink) { return DaoUtil.toPageData(tenantRepository - .findTenantInfoByRegionNextPage( - region, + .findTenantInfosNextPage( Objects.toString(pageLink.getTextSearch(), ""), DaoUtil.toPageable(pageLink, TenantInfoEntity.tenantInfoColumnMap))); } diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/tenant/TenantRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sql/tenant/TenantRepository.java index 4ae1278bec..8f0bb05021 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/tenant/TenantRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/tenant/TenantRepository.java @@ -36,19 +36,15 @@ public interface TenantRepository extends PagingAndSortingRepository 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 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 findTenantInfoByRegionNextPage(@Param("region") String region, - @Param("textSearch") String textSearch, + "WHERE LOWER(t.searchText) LIKE LOWER(CONCAT('%', :textSearch, '%'))") + Page findTenantInfosNextPage(@Param("textSearch") String textSearch, Pageable pageable); @Query("SELECT t.id FROM TenantEntity t") diff --git a/dao/src/main/java/org/thingsboard/server/dao/tenant/TenantDao.java b/dao/src/main/java/org/thingsboard/server/dao/tenant/TenantDao.java index 94ca9c8c1e..576d102f00 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/tenant/TenantDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/tenant/TenantDao.java @@ -37,15 +37,14 @@ public interface TenantDao extends Dao { 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 findTenantsByRegion(TenantId tenantId, String region, PageLink pageLink); + PageData findTenants(TenantId tenantId, PageLink pageLink); - PageData findTenantInfosByRegion(TenantId tenantId, String region, PageLink pageLink); + PageData findTenantInfos(TenantId tenantId, PageLink pageLink); PageData findTenantsIds(PageLink pageLink); diff --git a/dao/src/main/java/org/thingsboard/server/dao/tenant/TenantServiceImpl.java b/dao/src/main/java/org/thingsboard/server/dao/tenant/TenantServiceImpl.java index 056f5987d8..a8a1dc6e5b 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/tenant/TenantServiceImpl.java +++ b/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 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 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 tenantValidator = @@ -214,12 +213,12 @@ public class TenantServiceImpl extends AbstractEntityService implements TenantSe } }; - private PaginatedRemover tenantsRemover = - new PaginatedRemover() { + private PaginatedRemover tenantsRemover = + new PaginatedRemover<>() { @Override - protected PageData findEntities(TenantId tenantId, String region, PageLink pageLink) { - return tenantDao.findTenantsByRegion(tenantId, region, pageLink); + protected PageData findEntities(TenantId tenantId, TenantId id, PageLink pageLink) { + return tenantDao.findTenants(tenantId, pageLink); } @Override diff --git a/dao/src/test/java/org/thingsboard/server/dao/service/event/BaseEventServiceTest.java b/dao/src/test/java/org/thingsboard/server/dao/service/event/BaseEventServiceTest.java index 155db106ff..7d81e5afc6 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/service/event/BaseEventServiceTest.java +++ b/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 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 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; } } diff --git a/dao/src/test/java/org/thingsboard/server/dao/sql/event/JpaBaseEventDaoTest.java b/dao/src/test/java/org/thingsboard/server/dao/sql/event/JpaBaseEventDaoTest.java index 4b392c6030..07f0a52338 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/sql/event/JpaBaseEventDaoTest.java +++ b/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 optEvent1 = eventDao.saveIfNotExists(event); - assertTrue("Optional is expected to be non-empty", optEvent1.isPresent()); - assertEquals(event, optEvent1.get()); - Optional 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(); } diff --git a/dao/src/test/java/org/thingsboard/server/dao/sql/tenant/JpaTenantDaoTest.java b/dao/src/test/java/org/thingsboard/server/dao/sql/tenant/JpaTenantDaoTest.java index 03d1eb01b1..916b572fba 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/sql/tenant/JpaTenantDaoTest.java +++ b/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 tenants1 = tenantDao.findTenantsByRegion(AbstractServiceTest.SYSTEM_TENANT_ID, "REGION_1", pageLink); + PageData tenants1 = tenantDao.findTenants(AbstractServiceTest.SYSTEM_TENANT_ID, pageLink); assertEquals(20, tenants1.getData().size()); pageLink = pageLink.nextPageLink(); - PageData tenants2 = tenantDao.findTenantsByRegion(AbstractServiceTest.SYSTEM_TENANT_ID, "REGION_1", + PageData tenants2 = tenantDao.findTenants(AbstractServiceTest.SYSTEM_TENANT_ID, pageLink); assertEquals(10, tenants2.getData().size()); pageLink = pageLink.nextPageLink(); - PageData tenants3 = tenantDao.findTenantsByRegion(AbstractServiceTest.SYSTEM_TENANT_ID, "REGION_1", + PageData 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));