Browse Source

Cache transaction support for versioned entities

pull/10977/head
ViacheslavKlimov 2 years ago
parent
commit
31d2d14f60
  1. 4
      application/src/main/data/upgrade/3.7.0/schema_update.sql
  2. 5
      common/cache/src/main/java/org/thingsboard/server/cache/CaffeineTbTransactionalCache.java
  3. 2
      common/cache/src/main/java/org/thingsboard/server/cache/RedisTbCacheTransaction.java
  4. 13
      common/cache/src/main/java/org/thingsboard/server/cache/RedisTbTransactionalCache.java
  5. 6
      common/cache/src/main/java/org/thingsboard/server/cache/TbTransactionalCache.java
  6. 40
      common/cache/src/main/java/org/thingsboard/server/cache/VersionedRedisTbCache.java
  7. 8
      common/cache/src/main/java/org/thingsboard/server/cache/VersionedTbCache.java

4
application/src/main/data/upgrade/3.7.0/schema_update.sql

@ -14,10 +14,12 @@
-- limitations under the License.
--
-- UPDATE PUBLIC CUSTOMERS START
-- KV VERSIONING UPDATE START
CREATE SEQUENCE IF NOT EXISTS attribute_kv_version_seq cache 1000;
CREATE SEQUENCE IF NOT EXISTS ts_kv_latest_version_seq cache 1000;
ALTER TABLE attribute_kv ADD COLUMN version bigint default 0;
ALTER TABLE ts_kv_latest ADD COLUMN version bigint default 0;
-- KV VERSIONING UPDATE END

5
common/cache/src/main/java/org/thingsboard/server/cache/CaffeineTbTransactionalCache.java

@ -54,6 +54,11 @@ public abstract class CaffeineTbTransactionalCache<K extends Serializable, V ext
return SimpleTbCacheValueWrapper.wrap(cache.get(key));
}
@Override
public TbCacheValueWrapper<V> get(K key, boolean transactionMode) {
return get(key);
}
@Override
public void put(K key, V value) {
lock.lock();

2
common/cache/src/main/java/org/thingsboard/server/cache/RedisTbCacheTransaction.java

@ -31,7 +31,7 @@ public class RedisTbCacheTransaction<K extends Serializable, V extends Serializa
@Override
public void put(K key, V value) {
cache.put(key, value, connection);
cache.put(key, value, connection, true);
}
@Override

13
common/cache/src/main/java/org/thingsboard/server/cache/RedisTbTransactionalCache.java

@ -77,9 +77,14 @@ public abstract class RedisTbTransactionalCache<K extends Serializable, V extend
@Override
public TbCacheValueWrapper<V> get(K key) {
return get(key, false);
}
@Override
public TbCacheValueWrapper<V> get(K key, boolean transactionMode) {
try (var connection = connectionFactory.getConnection()) {
byte[] rawKey = getRawKey(key);
byte[] rawValue = doGet(connection, rawKey);
byte[] rawValue = doGet(connection, rawKey, transactionMode);
if (rawValue == null || rawValue.length == 0) {
return null;
} else if (Arrays.equals(rawValue, BINARY_NULL_VALUE)) {
@ -96,18 +101,18 @@ public abstract class RedisTbTransactionalCache<K extends Serializable, V extend
}
}
protected byte[] doGet(RedisConnection connection, byte[] rawKey) {
protected byte[] doGet(RedisConnection connection, byte[] rawKey, boolean transactionMode) {
return connection.stringCommands().get(rawKey);
}
@Override
public void put(K key, V value) {
try (var connection = connectionFactory.getConnection()) {
put(key, value, connection);
put(key, value, connection, false);
}
}
public void put(K key, V value, RedisConnection connection) {
public void put(K key, V value, RedisConnection connection, boolean transactionMode) {
put(connection, key, value, RedisStringCommands.SetOption.UPSERT);
}

6
common/cache/src/main/java/org/thingsboard/server/cache/TbTransactionalCache.java

@ -27,6 +27,8 @@ public interface TbTransactionalCache<K extends Serializable, V extends Serializ
TbCacheValueWrapper<V> get(K key);
TbCacheValueWrapper<V> get(K key, boolean transactionMode);
void put(K key, V value);
void putIfAbsent(K key, V value);
@ -60,7 +62,7 @@ public interface TbTransactionalCache<K extends Serializable, V extends Serializ
}
default V getAndPutInTransaction(K key, Supplier<V> dbCall, boolean cacheNullValue) {
TbCacheValueWrapper<V> cacheValueWrapper = get(key);
TbCacheValueWrapper<V> cacheValueWrapper = get(key, true);
if (cacheValueWrapper != null) {
return cacheValueWrapper.get();
}
@ -95,7 +97,7 @@ public interface TbTransactionalCache<K extends Serializable, V extends Serializ
}
default <R> R getAndPutInTransaction(K key, Supplier<R> dbCall, Function<V, R> cacheValueToResult, Function<R, V> dbValueToCacheValue, boolean cacheNullValue) {
TbCacheValueWrapper<V> cacheValueWrapper = get(key);
TbCacheValueWrapper<V> cacheValueWrapper = get(key, true);
if (cacheValueWrapper != null) {
var cacheValue = cacheValueWrapper.get();
return cacheValue == null ? null : cacheValueToResult.apply(cacheValue);

40
common/cache/src/main/java/org/thingsboard/server/cache/VersionedRedisTbCache.java

@ -88,31 +88,30 @@ public abstract class VersionedRedisTbCache<K extends Serializable, V extends Se
}
@Override
protected byte[] doGet(RedisConnection connection, byte[] rawKey) {
protected byte[] doGet(RedisConnection connection, byte[] rawKey, boolean transactionMode) {
if (transactionMode) {
return super.doGet(connection, rawKey, true);
}
return connection.stringCommands().getRange(rawKey, VERSION_SIZE, VALUE_END_OFFSET);
}
@Override
public void put(K key, V value) {
Long version;
if (value == null) {
version = 0L;
} else if (value.getVersion() != null) {
version = value.getVersion();
} else {
Long version = getVersion(value);
if (version == null) {
return;
}
doPut(key, value, version, cacheTtl);
}
@Override
public void put(K key, V value, RedisConnection connection) {
Long version;
if (value == null) {
version = 0L;
} else if (value.getVersion() != null) {
version = value.getVersion();
} else {
public void put(K key, V value, RedisConnection connection, boolean transactionMode) {
if (transactionMode) {
super.put(key, value, connection, true); // because scripting commands are not supported in transaction mode
return;
}
Long version = getVersion(value);
if (version == null) {
return;
}
byte[] rawKey = getRawKey(key);
@ -121,9 +120,6 @@ public abstract class VersionedRedisTbCache<K extends Serializable, V extends Se
private void doPut(K key, V value, Long version, Expiration expiration) {
log.trace("put [{}][{}][{}]", key, value, version);
if (version == null) {
return;
}
final byte[] rawKey = getRawKey(key);
try (var connection = getConnection(rawKey)) {
doPut(rawKey, value, version, expiration, connection);
@ -169,4 +165,14 @@ public abstract class VersionedRedisTbCache<K extends Serializable, V extends Se
throw new NotImplementedException("evictOrPut is not supported by versioned cache");
}
private Long getVersion(V value) {
if (value == null) {
return 0L;
} else if (value.getVersion() != null) {
return value.getVersion();
} else {
return null;
}
}
}

8
common/cache/src/main/java/org/thingsboard/server/cache/VersionedTbCache.java

@ -27,11 +27,17 @@ public interface VersionedTbCache<K extends Serializable, V extends Serializable
TbCacheValueWrapper<V> get(K key);
default V get(K key, Supplier<V> supplier) {
return get(key, supplier, true);
}
default V get(K key, Supplier<V> supplier, boolean putToCache) {
return Optional.ofNullable(get(key))
.map(TbCacheValueWrapper::get)
.orElseGet(() -> {
V value = supplier.get();
put(key, value);
if (putToCache) {
put(key, value);
}
return value;
});
}

Loading…
Cancel
Save