|
|
|
@ -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: |
|
|
|
|