From 09eb511599225583100fe2eb2ab7c1f79bc98f48 Mon Sep 17 00:00:00 2001 From: dshvaika Date: Mon, 8 Dec 2025 16:51:40 +0200 Subject: [PATCH 1/5] Fix key_dictionary race causing cached keyId 0 --- .../dictionary/KeyDictionaryCompositeKey.java | 2 +- .../sqlts/dictionary/KeyDictionaryEntry.java | 5 +- .../sqlts/dictionary/JpaKeyDictionaryDao.java | 59 ++++------ .../dictionary/KeyDictionaryRepository.java | 4 + .../dictionary/KeyDictionaryDaoTest.java | 111 ++++++++++++++++++ 5 files changed, 141 insertions(+), 40 deletions(-) create mode 100644 dao/src/test/java/org/thingsboard/server/dao/sqlts/dictionary/KeyDictionaryDaoTest.java diff --git a/dao/src/main/java/org/thingsboard/server/dao/model/sqlts/dictionary/KeyDictionaryCompositeKey.java b/dao/src/main/java/org/thingsboard/server/dao/model/sqlts/dictionary/KeyDictionaryCompositeKey.java index 4f3285b9bf..00e49ea703 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/model/sqlts/dictionary/KeyDictionaryCompositeKey.java +++ b/dao/src/main/java/org/thingsboard/server/dao/model/sqlts/dictionary/KeyDictionaryCompositeKey.java @@ -25,7 +25,7 @@ import java.io.Serializable; @Data @NoArgsConstructor @AllArgsConstructor -public class KeyDictionaryCompositeKey implements Serializable{ +public class KeyDictionaryCompositeKey implements Serializable { @Transient private static final long serialVersionUID = -4089175869616037523L; diff --git a/dao/src/main/java/org/thingsboard/server/dao/model/sqlts/dictionary/KeyDictionaryEntry.java b/dao/src/main/java/org/thingsboard/server/dao/model/sqlts/dictionary/KeyDictionaryEntry.java index a95c7a2bc6..d98105c0bc 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/model/sqlts/dictionary/KeyDictionaryEntry.java +++ b/dao/src/main/java/org/thingsboard/server/dao/model/sqlts/dictionary/KeyDictionaryEntry.java @@ -36,8 +36,7 @@ public final class KeyDictionaryEntry { @Column(name = KEY_COLUMN) private String key; - @Column(name = KEY_ID_COLUMN, unique = true, columnDefinition = "int") - @Generated - private int keyId; + @Column(name = KEY_ID_COLUMN, unique = true, columnDefinition = "int", insertable = false, updatable = false) + private Integer keyId; } \ No newline at end of file diff --git a/dao/src/main/java/org/thingsboard/server/dao/sqlts/dictionary/JpaKeyDictionaryDao.java b/dao/src/main/java/org/thingsboard/server/dao/sqlts/dictionary/JpaKeyDictionaryDao.java index c14c069f23..1ef2c2d6aa 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sqlts/dictionary/JpaKeyDictionaryDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sqlts/dictionary/JpaKeyDictionaryDao.java @@ -17,8 +17,6 @@ package org.thingsboard.server.dao.sqlts.dictionary; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; -import org.hibernate.exception.ConstraintViolationException; -import org.springframework.dao.DataIntegrityViolationException; import org.springframework.stereotype.Component; import org.springframework.transaction.annotation.Propagation; import org.springframework.transaction.annotation.Transactional; @@ -48,43 +46,32 @@ public class JpaKeyDictionaryDao extends JpaAbstractDaoListeningExecutorService @Transactional(propagation = Propagation.NOT_SUPPORTED) @Override public Integer getOrSaveKeyId(String strKey) { - Integer keyId = keyDictionaryMap.get(strKey); - if (keyId == null) { - Optional tsKvDictionaryOptional; - tsKvDictionaryOptional = keyDictionaryRepository.findById(new KeyDictionaryCompositeKey(strKey)); - if (tsKvDictionaryOptional.isEmpty()) { - creationLock.lock(); - try { - keyId = keyDictionaryMap.get(strKey); - if (keyId != null) { - return keyId; - } - tsKvDictionaryOptional = keyDictionaryRepository.findById(new KeyDictionaryCompositeKey(strKey)); - if (tsKvDictionaryOptional.isEmpty()) { - KeyDictionaryEntry keyDictionaryEntry = new KeyDictionaryEntry(); - keyDictionaryEntry.setKey(strKey); - try { - KeyDictionaryEntry saved = keyDictionaryRepository.save(keyDictionaryEntry); - keyDictionaryMap.put(saved.getKey(), saved.getKeyId()); - keyId = saved.getKeyId(); - } catch (DataIntegrityViolationException | ConstraintViolationException e) { - tsKvDictionaryOptional = keyDictionaryRepository.findById(new KeyDictionaryCompositeKey(strKey)); - KeyDictionaryEntry dictionary = tsKvDictionaryOptional.orElseThrow(() -> new RuntimeException("Failed to get KeyDictionaryEntry entity from DB!")); - keyDictionaryMap.put(dictionary.getKey(), dictionary.getKeyId()); - keyId = dictionary.getKeyId(); - } - } else { - keyId = tsKvDictionaryOptional.get().getKeyId(); - } - } finally { - creationLock.unlock(); + Integer cached = keyDictionaryMap.get(strKey); + if (cached != null) { + return cached; + } + creationLock.lock(); + try { + Integer keyId = keyDictionaryMap.get(strKey); + if (keyId != null) { + return keyId; + } + keyId = keyDictionaryRepository.upsertAndGetKeyId(strKey); + if (keyId == null || keyId == 0) { + log.warn("upsertAndGetKeyId returned: [{}] for key: [{}], falling back to findById", keyId, strKey); + KeyDictionaryCompositeKey id = new KeyDictionaryCompositeKey(strKey); + Optional entryOpt = keyDictionaryRepository.findById(id); + if (entryOpt.isEmpty() || + entryOpt.get().getKeyId() == null || + entryOpt.get().getKeyId() == 0) { + throw new IllegalStateException("Failed to resolve keyId for string key: " + strKey + " after fallback. keyId: " + keyId); } - } else { - keyId = tsKvDictionaryOptional.get().getKeyId(); - keyDictionaryMap.put(strKey, keyId); } + keyDictionaryMap.put(strKey, keyId); + return keyId; + } finally { + creationLock.unlock(); } - return keyId; } @Override diff --git a/dao/src/main/java/org/thingsboard/server/dao/sqlts/dictionary/KeyDictionaryRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sqlts/dictionary/KeyDictionaryRepository.java index d264cd9966..a1141b62ff 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sqlts/dictionary/KeyDictionaryRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sqlts/dictionary/KeyDictionaryRepository.java @@ -19,6 +19,7 @@ import org.springframework.data.domain.Page; import org.springframework.data.domain.Pageable; import org.springframework.data.jpa.repository.JpaRepository; import org.springframework.data.jpa.repository.Query; +import org.springframework.data.repository.query.Param; import org.thingsboard.server.dao.model.sqlts.dictionary.KeyDictionaryCompositeKey; import org.thingsboard.server.dao.model.sqlts.dictionary.KeyDictionaryEntry; @@ -31,4 +32,7 @@ public interface KeyDictionaryRepository extends JpaRepository findAll(Pageable pageable); + @Query(value = "INSERT INTO key_dictionary (key) VALUES (:key) ON CONFLICT (key) DO UPDATE SET key = EXCLUDED.key RETURNING key_id", nativeQuery = true) + Integer upsertAndGetKeyId(@Param("key") String key); + } \ No newline at end of file diff --git a/dao/src/test/java/org/thingsboard/server/dao/sqlts/dictionary/KeyDictionaryDaoTest.java b/dao/src/test/java/org/thingsboard/server/dao/sqlts/dictionary/KeyDictionaryDaoTest.java new file mode 100644 index 0000000000..11c14c99f4 --- /dev/null +++ b/dao/src/test/java/org/thingsboard/server/dao/sqlts/dictionary/KeyDictionaryDaoTest.java @@ -0,0 +1,111 @@ +/** + * Copyright © 2016-2025 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.sqlts.dictionary; + +import org.junit.Test; +import org.springframework.beans.factory.annotation.Autowired; +import org.thingsboard.server.dao.dictionary.KeyDictionaryDao; +import org.thingsboard.server.dao.model.sqlts.dictionary.KeyDictionaryCompositeKey; +import org.thingsboard.server.dao.model.sqlts.dictionary.KeyDictionaryEntry; +import org.thingsboard.server.dao.service.AbstractServiceTest; +import org.thingsboard.server.dao.service.DaoSqlTest; + +import java.util.Arrays; +import java.util.Optional; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.TimeUnit; + +import static org.assertj.core.api.Assertions.assertThat; + +@DaoSqlTest +public class KeyDictionaryDaoTest extends AbstractServiceTest { + + @Autowired + private KeyDictionaryDao keyDictionaryDao; + + @Autowired + private KeyDictionaryRepository keyDictionaryRepository; + + private static final String KEY = "testKeyDictionaryDaoTestKey"; + + @Test + public void testGetOrSaveKeyId_concurrent() throws Exception { + int threads = 8; + ExecutorService executor = Executors.newFixedThreadPool(threads); + + CountDownLatch allReady = new CountDownLatch(threads); + CountDownLatch start = new CountDownLatch(1); + CountDownLatch allDone = new CountDownLatch(threads); + + Integer[] keyIds = new Integer[threads]; + + try { + for (int i = 0; i < threads; i++) { + final int idx = i; + executor.submit(() -> { + allReady.countDown(); + try { + // wait until all threads are ready + start.await(); + // concurrent call + Integer id = keyDictionaryDao.getOrSaveKeyId(KEY); + keyIds[idx] = id; + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } finally { + allDone.countDown(); + } + }); + } + + // ensure all threads are queued + allReady.await(5, TimeUnit.SECONDS); + // fire the start gun + start.countDown(); + // wait for all to finish + allDone.await(10, TimeUnit.SECONDS); + } finally { + executor.shutdownNow(); + } + + // basic sanity + for (int i = 0; i < threads; i++) { + assertThat(keyIds[i]) + .as("keyId[%s]", i) + .isNotNull() + .isGreaterThan(0); + } + + // all threads must see the same keyId + int first = keyIds[0]; + assertThat(first).isGreaterThan(0); + assertThat(Arrays.stream(keyIds).distinct().count()) + .as("all threads should get the same keyId") + .isEqualTo(1); + + // DB must have exactly one row for this key and the same id + KeyDictionaryCompositeKey id = new KeyDictionaryCompositeKey(KEY); + Optional entry = keyDictionaryRepository.findById(id); + + assertThat(entry.isPresent()).isTrue(); + assertThat(entry.get().getKeyId()).isEqualTo(first); + + keyDictionaryRepository.deleteById(id); + } + +} From d66e9ecf74b430ac6ab9bee4fa6b863157fbb419 Mon Sep 17 00:00:00 2001 From: dshvaika Date: Mon, 8 Dec 2025 17:45:00 +0200 Subject: [PATCH 2/5] Added new lines to the end of files --- .../server/dao/model/sqlts/dictionary/KeyDictionaryEntry.java | 2 +- .../server/dao/sqlts/dictionary/KeyDictionaryRepository.java | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/dao/src/main/java/org/thingsboard/server/dao/model/sqlts/dictionary/KeyDictionaryEntry.java b/dao/src/main/java/org/thingsboard/server/dao/model/sqlts/dictionary/KeyDictionaryEntry.java index d98105c0bc..8365e8facc 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/model/sqlts/dictionary/KeyDictionaryEntry.java +++ b/dao/src/main/java/org/thingsboard/server/dao/model/sqlts/dictionary/KeyDictionaryEntry.java @@ -39,4 +39,4 @@ public final class KeyDictionaryEntry { @Column(name = KEY_ID_COLUMN, unique = true, columnDefinition = "int", insertable = false, updatable = false) private Integer keyId; -} \ No newline at end of file +} diff --git a/dao/src/main/java/org/thingsboard/server/dao/sqlts/dictionary/KeyDictionaryRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sqlts/dictionary/KeyDictionaryRepository.java index a1141b62ff..e836cedb19 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sqlts/dictionary/KeyDictionaryRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sqlts/dictionary/KeyDictionaryRepository.java @@ -35,4 +35,4 @@ public interface KeyDictionaryRepository extends JpaRepository Date: Mon, 8 Dec 2025 18:36:01 +0200 Subject: [PATCH 3/5] Added find before lock for the startup of application --- .../dao/sqlts/dictionary/JpaKeyDictionaryDao.java | 12 ++++++++++-- 1 file changed, 10 insertions(+), 2 deletions(-) diff --git a/dao/src/main/java/org/thingsboard/server/dao/sqlts/dictionary/JpaKeyDictionaryDao.java b/dao/src/main/java/org/thingsboard/server/dao/sqlts/dictionary/JpaKeyDictionaryDao.java index 1ef2c2d6aa..bcb9284371 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sqlts/dictionary/JpaKeyDictionaryDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sqlts/dictionary/JpaKeyDictionaryDao.java @@ -50,6 +50,15 @@ public class JpaKeyDictionaryDao extends JpaAbstractDaoListeningExecutorService if (cached != null) { return cached; } + var compositeKey = new KeyDictionaryCompositeKey(strKey); + Optional entryOpt = keyDictionaryRepository.findById(compositeKey); + if (entryOpt.isPresent()) { + Integer keyId = entryOpt.get().getKeyId(); + if (keyId != null) { + keyDictionaryMap.put(strKey, keyId); + return keyId; + } + } creationLock.lock(); try { Integer keyId = keyDictionaryMap.get(strKey); @@ -59,8 +68,7 @@ public class JpaKeyDictionaryDao extends JpaAbstractDaoListeningExecutorService keyId = keyDictionaryRepository.upsertAndGetKeyId(strKey); if (keyId == null || keyId == 0) { log.warn("upsertAndGetKeyId returned: [{}] for key: [{}], falling back to findById", keyId, strKey); - KeyDictionaryCompositeKey id = new KeyDictionaryCompositeKey(strKey); - Optional entryOpt = keyDictionaryRepository.findById(id); + entryOpt = keyDictionaryRepository.findById(compositeKey); if (entryOpt.isEmpty() || entryOpt.get().getKeyId() == null || entryOpt.get().getKeyId() == 0) { From c0af057590273148da4e65fdb147ba3a9c38a23b Mon Sep 17 00:00:00 2001 From: dshvaika Date: Tue, 9 Dec 2025 11:31:02 +0200 Subject: [PATCH 4/5] refactoring logic in getOrSaveKeyId method --- .../sqlts/dictionary/JpaKeyDictionaryDao.java | 44 ++++++++++--------- 1 file changed, 23 insertions(+), 21 deletions(-) diff --git a/dao/src/main/java/org/thingsboard/server/dao/sqlts/dictionary/JpaKeyDictionaryDao.java b/dao/src/main/java/org/thingsboard/server/dao/sqlts/dictionary/JpaKeyDictionaryDao.java index bcb9284371..9d23bdbcb5 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sqlts/dictionary/JpaKeyDictionaryDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sqlts/dictionary/JpaKeyDictionaryDao.java @@ -51,32 +51,29 @@ public class JpaKeyDictionaryDao extends JpaAbstractDaoListeningExecutorService return cached; } var compositeKey = new KeyDictionaryCompositeKey(strKey); - Optional entryOpt = keyDictionaryRepository.findById(compositeKey); - if (entryOpt.isPresent()) { - Integer keyId = entryOpt.get().getKeyId(); - if (keyId != null) { - keyDictionaryMap.put(strKey, keyId); - return keyId; - } + Optional existingId = keyDictionaryRepository.findById(compositeKey) + .map(KeyDictionaryEntry::getKeyId) + .filter(id -> id != 0); + if (existingId.isPresent()) { + return cacheAndReturn(strKey, existingId.get()); } creationLock.lock(); try { - Integer keyId = keyDictionaryMap.get(strKey); - if (keyId != null) { - return keyId; + Integer fromCache = keyDictionaryMap.get(strKey); + if (fromCache != null) { + return fromCache; } - keyId = keyDictionaryRepository.upsertAndGetKeyId(strKey); - if (keyId == null || keyId == 0) { - log.warn("upsertAndGetKeyId returned: [{}] for key: [{}], falling back to findById", keyId, strKey); - entryOpt = keyDictionaryRepository.findById(compositeKey); - if (entryOpt.isEmpty() || - entryOpt.get().getKeyId() == null || - entryOpt.get().getKeyId() == 0) { - throw new IllegalStateException("Failed to resolve keyId for string key: " + strKey + " after fallback. keyId: " + keyId); - } + Integer keyId = keyDictionaryRepository.upsertAndGetKeyId(strKey); + if (keyId != null && keyId != 0) { + return cacheAndReturn(strKey, keyId); } - keyDictionaryMap.put(strKey, keyId); - return keyId; + log.warn("upsertAndGetKeyId returned: [{}] for key: [{}], falling back to findById", keyId, strKey); + keyId = keyDictionaryRepository.findById(compositeKey) + .map(KeyDictionaryEntry::getKeyId) + .filter(id -> id != 0) + .orElseThrow(() -> new IllegalStateException( + "Failed to resolve keyId for string key: " + strKey + " after fallback.")); + return cacheAndReturn(strKey, keyId); } finally { creationLock.unlock(); } @@ -93,4 +90,9 @@ public class JpaKeyDictionaryDao extends JpaAbstractDaoListeningExecutorService return DaoUtil.pageToPageData(keyDictionaryRepository.findAll(DaoUtil.toPageable(pageLink))); } + private Integer cacheAndReturn(String key, Integer keyId) { + keyDictionaryMap.put(key, keyId); + return keyId; + } + } From 2c04943984abebdd7fe07772046e20d2173e8531 Mon Sep 17 00:00:00 2001 From: dshvaika Date: Tue, 9 Dec 2025 12:11:14 +0200 Subject: [PATCH 5/5] removed checks for 0 value --- .../server/dao/sqlts/dictionary/JpaKeyDictionaryDao.java | 7 ++----- 1 file changed, 2 insertions(+), 5 deletions(-) diff --git a/dao/src/main/java/org/thingsboard/server/dao/sqlts/dictionary/JpaKeyDictionaryDao.java b/dao/src/main/java/org/thingsboard/server/dao/sqlts/dictionary/JpaKeyDictionaryDao.java index 9d23bdbcb5..46bc10010e 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sqlts/dictionary/JpaKeyDictionaryDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sqlts/dictionary/JpaKeyDictionaryDao.java @@ -51,9 +51,7 @@ public class JpaKeyDictionaryDao extends JpaAbstractDaoListeningExecutorService return cached; } var compositeKey = new KeyDictionaryCompositeKey(strKey); - Optional existingId = keyDictionaryRepository.findById(compositeKey) - .map(KeyDictionaryEntry::getKeyId) - .filter(id -> id != 0); + Optional existingId = keyDictionaryRepository.findById(compositeKey).map(KeyDictionaryEntry::getKeyId); if (existingId.isPresent()) { return cacheAndReturn(strKey, existingId.get()); } @@ -64,13 +62,12 @@ public class JpaKeyDictionaryDao extends JpaAbstractDaoListeningExecutorService return fromCache; } Integer keyId = keyDictionaryRepository.upsertAndGetKeyId(strKey); - if (keyId != null && keyId != 0) { + if (keyId != null) { return cacheAndReturn(strKey, keyId); } log.warn("upsertAndGetKeyId returned: [{}] for key: [{}], falling back to findById", keyId, strKey); keyId = keyDictionaryRepository.findById(compositeKey) .map(KeyDictionaryEntry::getKeyId) - .filter(id -> id != 0) .orElseThrow(() -> new IllegalStateException( "Failed to resolve keyId for string key: " + strKey + " after fallback.")); return cacheAndReturn(strKey, keyId);