From 123457f8eb5313770618615cea3570fd08e37c27 Mon Sep 17 00:00:00 2001 From: Andrii Shvaika Date: Wed, 14 Apr 2021 13:49:09 +0300 Subject: [PATCH] Prepared Statement initialization lock --- .../CassandraBaseTimeseriesDao.java | 151 ++++++++++++------ 1 file changed, 105 insertions(+), 46 deletions(-) diff --git a/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java b/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java index 7d09578dd1..240d5a0b88 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java @@ -62,6 +62,8 @@ import java.util.Collections; import java.util.List; import java.util.Optional; import java.util.concurrent.TimeUnit; +import java.util.concurrent.locks.Lock; +import java.util.concurrent.locks.ReentrantLock; import java.util.stream.Collectors; import static com.datastax.oss.driver.api.querybuilder.QueryBuilder.literal; @@ -107,6 +109,7 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD private PreparedStatement[] fetchStmtsDesc; private PreparedStatement deleteStmt; private PreparedStatement deletePartitionStmt; + private final Lock stmtCreationLock = new ReentrantLock(); private boolean isInstall() { return environment.acceptsProfiles(Profiles.of("install")); @@ -545,13 +548,20 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD private PreparedStatement getDeleteStmt() { if (deleteStmt == null) { - deleteStmt = prepare("DELETE FROM " + ModelConstants.TS_KV_CF + - " WHERE " + ModelConstants.ENTITY_TYPE_COLUMN + EQUALS_PARAM - + "AND " + ModelConstants.ENTITY_ID_COLUMN + EQUALS_PARAM - + "AND " + ModelConstants.KEY_COLUMN + EQUALS_PARAM - + "AND " + ModelConstants.PARTITION_COLUMN + EQUALS_PARAM - + "AND " + ModelConstants.TS_COLUMN + " >= ? " - + "AND " + ModelConstants.TS_COLUMN + " < ?"); + stmtCreationLock.lock(); + try { + if (deleteStmt == null) { + deleteStmt = prepare("DELETE FROM " + ModelConstants.TS_KV_CF + + " WHERE " + ModelConstants.ENTITY_TYPE_COLUMN + EQUALS_PARAM + + "AND " + ModelConstants.ENTITY_ID_COLUMN + EQUALS_PARAM + + "AND " + ModelConstants.KEY_COLUMN + EQUALS_PARAM + + "AND " + ModelConstants.PARTITION_COLUMN + EQUALS_PARAM + + "AND " + ModelConstants.TS_COLUMN + " >= ? " + + "AND " + ModelConstants.TS_COLUMN + " < ?"); + } + } finally { + stmtCreationLock.unlock(); + } } return deleteStmt; } @@ -585,27 +595,41 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD private PreparedStatement getDeletePartitionStmt() { if (deletePartitionStmt == null) { - deletePartitionStmt = prepare("DELETE FROM " + ModelConstants.TS_KV_PARTITIONS_CF + - " WHERE " + ModelConstants.ENTITY_TYPE_COLUMN + EQUALS_PARAM - + "AND " + ModelConstants.ENTITY_ID_COLUMN + EQUALS_PARAM - + "AND " + ModelConstants.PARTITION_COLUMN + EQUALS_PARAM - + "AND " + ModelConstants.KEY_COLUMN + EQUALS_PARAM); + stmtCreationLock.lock(); + try { + if (deletePartitionStmt == null) { + deletePartitionStmt = prepare("DELETE FROM " + ModelConstants.TS_KV_PARTITIONS_CF + + " WHERE " + ModelConstants.ENTITY_TYPE_COLUMN + EQUALS_PARAM + + "AND " + ModelConstants.ENTITY_ID_COLUMN + EQUALS_PARAM + + "AND " + ModelConstants.PARTITION_COLUMN + EQUALS_PARAM + + "AND " + ModelConstants.KEY_COLUMN + EQUALS_PARAM); + } + } finally { + stmtCreationLock.unlock(); + } } return deletePartitionStmt; } private PreparedStatement getSaveStmt(DataType dataType) { if (saveStmts == null) { - saveStmts = new PreparedStatement[DataType.values().length]; - for (DataType type : DataType.values()) { - saveStmts[type.ordinal()] = prepare(INSERT_INTO + ModelConstants.TS_KV_CF + - "(" + ModelConstants.ENTITY_TYPE_COLUMN + - "," + ModelConstants.ENTITY_ID_COLUMN + - "," + ModelConstants.KEY_COLUMN + - "," + ModelConstants.PARTITION_COLUMN + - "," + ModelConstants.TS_COLUMN + - "," + getColumnName(type) + ")" + - " VALUES(?, ?, ?, ?, ?, ?)"); + stmtCreationLock.lock(); + try { + if (saveStmts == null) { + saveStmts = new PreparedStatement[DataType.values().length]; + for (DataType type : DataType.values()) { + saveStmts[type.ordinal()] = prepare(INSERT_INTO + ModelConstants.TS_KV_CF + + "(" + ModelConstants.ENTITY_TYPE_COLUMN + + "," + ModelConstants.ENTITY_ID_COLUMN + + "," + ModelConstants.KEY_COLUMN + + "," + ModelConstants.PARTITION_COLUMN + + "," + ModelConstants.TS_COLUMN + + "," + getColumnName(type) + ")" + + " VALUES(?, ?, ?, ?, ?, ?)"); + } + } + } finally { + stmtCreationLock.unlock(); } } return saveStmts[dataType.ordinal()]; @@ -613,16 +637,23 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD private PreparedStatement getSaveTtlStmt(DataType dataType) { if (saveTtlStmts == null) { - saveTtlStmts = new PreparedStatement[DataType.values().length]; - for (DataType type : DataType.values()) { - saveTtlStmts[type.ordinal()] = prepare(INSERT_INTO + ModelConstants.TS_KV_CF + - "(" + ModelConstants.ENTITY_TYPE_COLUMN + - "," + ModelConstants.ENTITY_ID_COLUMN + - "," + ModelConstants.KEY_COLUMN + - "," + ModelConstants.PARTITION_COLUMN + - "," + ModelConstants.TS_COLUMN + - "," + getColumnName(type) + ")" + - " VALUES(?, ?, ?, ?, ?, ?) USING TTL ?"); + stmtCreationLock.lock(); + try { + if (saveTtlStmts == null) { + saveTtlStmts = new PreparedStatement[DataType.values().length]; + for (DataType type : DataType.values()) { + saveTtlStmts[type.ordinal()] = prepare(INSERT_INTO + ModelConstants.TS_KV_CF + + "(" + ModelConstants.ENTITY_TYPE_COLUMN + + "," + ModelConstants.ENTITY_ID_COLUMN + + "," + ModelConstants.KEY_COLUMN + + "," + ModelConstants.PARTITION_COLUMN + + "," + ModelConstants.TS_COLUMN + + "," + getColumnName(type) + ")" + + " VALUES(?, ?, ?, ?, ?, ?) USING TTL ?"); + } + } + } finally { + stmtCreationLock.unlock(); } } return saveTtlStmts[dataType.ordinal()]; @@ -630,24 +661,38 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD private PreparedStatement getPartitionInsertStmt() { if (partitionInsertStmt == null) { - partitionInsertStmt = prepare(INSERT_INTO + ModelConstants.TS_KV_PARTITIONS_CF + - "(" + ModelConstants.ENTITY_TYPE_COLUMN + - "," + ModelConstants.ENTITY_ID_COLUMN + - "," + ModelConstants.PARTITION_COLUMN + - "," + ModelConstants.KEY_COLUMN + ")" + - " VALUES(?, ?, ?, ?)"); + stmtCreationLock.lock(); + try { + if (partitionInsertStmt == null) { + partitionInsertStmt = prepare(INSERT_INTO + ModelConstants.TS_KV_PARTITIONS_CF + + "(" + ModelConstants.ENTITY_TYPE_COLUMN + + "," + ModelConstants.ENTITY_ID_COLUMN + + "," + ModelConstants.PARTITION_COLUMN + + "," + ModelConstants.KEY_COLUMN + ")" + + " VALUES(?, ?, ?, ?)"); + } + } finally { + stmtCreationLock.unlock(); + } } return partitionInsertStmt; } private PreparedStatement getPartitionInsertTtlStmt() { if (partitionInsertTtlStmt == null) { - partitionInsertTtlStmt = prepare(INSERT_INTO + ModelConstants.TS_KV_PARTITIONS_CF + - "(" + ModelConstants.ENTITY_TYPE_COLUMN + - "," + ModelConstants.ENTITY_ID_COLUMN + - "," + ModelConstants.PARTITION_COLUMN + - "," + ModelConstants.KEY_COLUMN + ")" + - " VALUES(?, ?, ?, ?) USING TTL ?"); + stmtCreationLock.lock(); + try { + if (partitionInsertTtlStmt == null) { + partitionInsertTtlStmt = prepare(INSERT_INTO + ModelConstants.TS_KV_PARTITIONS_CF + + "(" + ModelConstants.ENTITY_TYPE_COLUMN + + "," + ModelConstants.ENTITY_ID_COLUMN + + "," + ModelConstants.PARTITION_COLUMN + + "," + ModelConstants.KEY_COLUMN + ")" + + " VALUES(?, ?, ?, ?) USING TTL ?"); + } + } finally { + stmtCreationLock.unlock(); + } } return partitionInsertTtlStmt; } @@ -713,12 +758,26 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD switch (orderBy) { case ASC_ORDER: if (fetchStmtsAsc == null) { - fetchStmtsAsc = initFetchStmt(orderBy); + stmtCreationLock.lock(); + try { + if (fetchStmtsAsc == null) { + fetchStmtsAsc = initFetchStmt(orderBy); + } + } finally { + stmtCreationLock.unlock(); + } } return fetchStmtsAsc[aggType.ordinal()]; case DESC_ORDER: if (fetchStmtsDesc == null) { - fetchStmtsDesc = initFetchStmt(orderBy); + stmtCreationLock.lock(); + try { + if (fetchStmtsDesc == null) { + fetchStmtsDesc = initFetchStmt(orderBy); + } + } finally { + stmtCreationLock.unlock(); + } } return fetchStmtsDesc[aggType.ordinal()]; default: