From bc3c06e51d9cca7509116b113f332b1c93175f0c Mon Sep 17 00:00:00 2001 From: ShvaykaD Date: Tue, 10 Sep 2019 15:42:34 +0300 Subject: [PATCH] Feature/timescale db + fixed violation constraint exception on primary key for sql time-series. (#1975) * init commit * update license * init ts-upgrade * update aggregation queries * aggregation update * revert upgrade init * merge with master * refactoring * revert thingsboard.yml * fix typo * fix typo * change packages * update packages * code update * add ts dao configs * fixed violation exception on primary key for sql timeseries * fix typo * fix typo * fix typo --- .../install/ThingsboardInstallService.java | 3 +- .../SqlTimescaleDatabaseSchemaService.java | 31 ++ .../src/main/resources/thingsboard.yml | 9 +- .../thingsboard/server/dao/JpaDaoConfig.java | 1 + .../server/dao/SqlTsDaoConfig.java | 35 +++ .../server/dao/TimescaleDaoConfig.java | 35 +++ .../server/dao/model/ModelConstants.java | 1 + .../dao/model/sql/AbsractTsKvEntity.java | 98 ++++++ .../timescale/TimescaleTsKvCompositeKey.java | 37 +++ .../sqlts/timescale/TimescaleTsKvEntity.java | 184 +++++++++++ .../{sql => sqlts/ts}/TsKvCompositeKey.java | 2 +- .../model/{sql => sqlts/ts}/TsKvEntity.java | 81 +---- .../ts}/TsKvLatestCompositeKey.java | 2 +- .../{sql => sqlts/ts}/TsKvLatestEntity.java | 2 +- .../dao/sqlts/AbstractSqlTimeseriesDao.java | 130 ++++++++ .../AbstractTimeseriesInsertRepository.java | 65 ++++ .../timescale/AggregationRepository.java | 101 ++++++ .../timescale/TimescaleInsertRepository.java | 93 ++++++ .../timescale/TimescaleTimeseriesDao.java | 295 ++++++++++++++++++ .../timescale/TsKvTimescaleRepository.java | 66 ++++ .../ts/HsqlTimeseriesInsertRepository.java | 93 ++++++ .../ts}/JpaTimeseriesDao.java | 284 +++++++---------- .../ts/PsqlTimeseriesInsertRepository.java | 93 ++++++ .../ts}/TsKvLatestRepository.java | 6 +- .../ts}/TsKvRepository.java | 6 +- .../server/dao/util/TimescaleDBTsDao.java | 22 ++ .../resources/sql/schema-timescale-idx.sql | 17 + .../main/resources/sql/schema-timescale.sql | 31 ++ .../server/dao/AbstractJpaDaoTest.java | 2 +- .../timeseries/BaseTimeseriesServiceTest.java | 20 -- dao/src/test/resources/sql-test.properties | 3 +- .../sql/timescale/drop-all-tables.sql | 21 ++ 32 files changed, 1586 insertions(+), 283 deletions(-) create mode 100644 application/src/main/java/org/thingsboard/server/service/install/SqlTimescaleDatabaseSchemaService.java create mode 100644 dao/src/main/java/org/thingsboard/server/dao/SqlTsDaoConfig.java create mode 100644 dao/src/main/java/org/thingsboard/server/dao/TimescaleDaoConfig.java create mode 100644 dao/src/main/java/org/thingsboard/server/dao/model/sql/AbsractTsKvEntity.java create mode 100644 dao/src/main/java/org/thingsboard/server/dao/model/sqlts/timescale/TimescaleTsKvCompositeKey.java create mode 100644 dao/src/main/java/org/thingsboard/server/dao/model/sqlts/timescale/TimescaleTsKvEntity.java rename dao/src/main/java/org/thingsboard/server/dao/model/{sql => sqlts/ts}/TsKvCompositeKey.java (95%) rename dao/src/main/java/org/thingsboard/server/dao/model/{sql => sqlts/ts}/TsKvEntity.java (58%) rename dao/src/main/java/org/thingsboard/server/dao/model/{sql => sqlts/ts}/TsKvLatestCompositeKey.java (95%) rename dao/src/main/java/org/thingsboard/server/dao/model/{sql => sqlts/ts}/TsKvLatestEntity.java (98%) create mode 100644 dao/src/main/java/org/thingsboard/server/dao/sqlts/AbstractSqlTimeseriesDao.java create mode 100644 dao/src/main/java/org/thingsboard/server/dao/sqlts/AbstractTimeseriesInsertRepository.java create mode 100644 dao/src/main/java/org/thingsboard/server/dao/sqlts/timescale/AggregationRepository.java create mode 100644 dao/src/main/java/org/thingsboard/server/dao/sqlts/timescale/TimescaleInsertRepository.java create mode 100644 dao/src/main/java/org/thingsboard/server/dao/sqlts/timescale/TimescaleTimeseriesDao.java create mode 100644 dao/src/main/java/org/thingsboard/server/dao/sqlts/timescale/TsKvTimescaleRepository.java create mode 100644 dao/src/main/java/org/thingsboard/server/dao/sqlts/ts/HsqlTimeseriesInsertRepository.java rename dao/src/main/java/org/thingsboard/server/dao/{sql/timeseries => sqlts/ts}/JpaTimeseriesDao.java (62%) create mode 100644 dao/src/main/java/org/thingsboard/server/dao/sqlts/ts/PsqlTimeseriesInsertRepository.java rename dao/src/main/java/org/thingsboard/server/dao/{sql/timeseries => sqlts/ts}/TsKvLatestRepository.java (84%) rename dao/src/main/java/org/thingsboard/server/dao/{sql/timeseries => sqlts/ts}/TsKvRepository.java (97%) create mode 100644 dao/src/main/java/org/thingsboard/server/dao/util/TimescaleDBTsDao.java create mode 100644 dao/src/main/resources/sql/schema-timescale-idx.sql create mode 100644 dao/src/main/resources/sql/schema-timescale.sql create mode 100644 dao/src/test/resources/sql/timescale/drop-all-tables.sql diff --git a/application/src/main/java/org/thingsboard/server/install/ThingsboardInstallService.java b/application/src/main/java/org/thingsboard/server/install/ThingsboardInstallService.java index 7b0c7cb780..f384817465 100644 --- a/application/src/main/java/org/thingsboard/server/install/ThingsboardInstallService.java +++ b/application/src/main/java/org/thingsboard/server/install/ThingsboardInstallService.java @@ -23,11 +23,11 @@ import org.springframework.context.ApplicationContext; import org.springframework.context.annotation.Profile; import org.springframework.stereotype.Service; import org.thingsboard.server.service.component.ComponentDiscoveryService; -import org.thingsboard.server.service.install.update.DataUpdateService; import org.thingsboard.server.service.install.DatabaseUpgradeService; import org.thingsboard.server.service.install.EntityDatabaseSchemaService; import org.thingsboard.server.service.install.SystemDataLoaderService; import org.thingsboard.server.service.install.TsDatabaseSchemaService; +import org.thingsboard.server.service.install.update.DataUpdateService; @Service @Profile("install") @@ -135,7 +135,6 @@ public class ThingsboardInstallService { systemDataLoaderService.deleteSystemWidgetBundle("date"); systemDataLoaderService.loadSystemWidgets(); - break; default: throw new RuntimeException("Unable to upgrade ThingsBoard, unsupported fromVersion: " + upgradeFromVersion); diff --git a/application/src/main/java/org/thingsboard/server/service/install/SqlTimescaleDatabaseSchemaService.java b/application/src/main/java/org/thingsboard/server/service/install/SqlTimescaleDatabaseSchemaService.java new file mode 100644 index 0000000000..23a335fe35 --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/install/SqlTimescaleDatabaseSchemaService.java @@ -0,0 +1,31 @@ +/** + * Copyright © 2016-2019 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.service.install; + +import org.springframework.context.annotation.Profile; +import org.springframework.stereotype.Service; +import org.thingsboard.server.dao.util.SqlDao; +import org.thingsboard.server.dao.util.TimescaleDBTsDao; + +@Service +@TimescaleDBTsDao +@Profile("install") +public class SqlTimescaleDatabaseSchemaService extends SqlAbstractDatabaseSchemaService + implements TsDatabaseSchemaService { + public SqlTimescaleDatabaseSchemaService() { + super("schema-timescale.sql", "schema-timescale-idx.sql"); + } +} \ No newline at end of file diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index 1aba5bbb33..94cee48053 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -119,8 +119,9 @@ database: entities: type: "${DATABASE_ENTITIES_TYPE:sql}" # cassandra OR sql ts: - type: "${DATABASE_TS_TYPE:sql}" # cassandra OR sql (for hybrid mode, only this value should be cassandra) + type: "${DATABASE_TS_TYPE:sql}" # cassandra, sql, or timescale (for hybrid mode, DATABASE_TS_TYPE value should be cassandra, or timescale) +# note: timescale works only with postgreSQL database for DATABASE_ENTITIES_TYPE. # Cassandra driver configuration parameters cassandra: @@ -189,10 +190,10 @@ cassandra: # SQL configuration parameters sql: - # Specify executor service type used to perform timeseries insert tasks: SINGLE FIXED CACHED + # Specify executor service type used to perform timeseries insert tasks: SINGLE or FIXED ts_inserts_executor_type: "${SQL_TS_INSERTS_EXECUTOR_TYPE:fixed}" # Specify thread pool size for FIXED executor service type - ts_inserts_fixed_thread_pool_size: "${SQL_TS_INSERTS_FIXED_THREAD_POOL_SIZE:10}" + ts_inserts_fixed_thread_pool_size: "${SQL_TS_INSERTS_FIXED_THREAD_POOL_SIZE:200}" # Actor system parameters actors: @@ -327,6 +328,8 @@ spring: url: "${SPRING_DATASOURCE_URL:jdbc:postgresql://localhost:5432/thingsboard}" username: "${SPRING_DATASOURCE_USERNAME:postgres}" password: "${SPRING_DATASOURCE_PASSWORD:postgres}" + hikari: + maximumPoolSize: "${SPRING_DATASOURCE_MAXIMUM_POOL_SIZE:50}" # Audit log parameters audit-log: diff --git a/dao/src/main/java/org/thingsboard/server/dao/JpaDaoConfig.java b/dao/src/main/java/org/thingsboard/server/dao/JpaDaoConfig.java index 5805eeaccf..58019f9208 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/JpaDaoConfig.java +++ b/dao/src/main/java/org/thingsboard/server/dao/JpaDaoConfig.java @@ -22,6 +22,7 @@ import org.springframework.context.annotation.Configuration; import org.springframework.data.jpa.repository.config.EnableJpaRepositories; import org.springframework.transaction.annotation.EnableTransactionManagement; import org.thingsboard.server.dao.util.SqlDao; +import org.thingsboard.server.dao.util.TimescaleDBTsDao; /** * @author Valerii Sosliuk diff --git a/dao/src/main/java/org/thingsboard/server/dao/SqlTsDaoConfig.java b/dao/src/main/java/org/thingsboard/server/dao/SqlTsDaoConfig.java new file mode 100644 index 0000000000..bc37316dda --- /dev/null +++ b/dao/src/main/java/org/thingsboard/server/dao/SqlTsDaoConfig.java @@ -0,0 +1,35 @@ +/** + * Copyright © 2016-2019 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; + +import org.springframework.boot.autoconfigure.EnableAutoConfiguration; +import org.springframework.boot.autoconfigure.domain.EntityScan; +import org.springframework.context.annotation.ComponentScan; +import org.springframework.context.annotation.Configuration; +import org.springframework.data.jpa.repository.config.EnableJpaRepositories; +import org.springframework.transaction.annotation.EnableTransactionManagement; +import org.thingsboard.server.dao.util.SqlTsDao; + +@Configuration +@EnableAutoConfiguration +@ComponentScan("org.thingsboard.server.dao.sqlts.ts") +@EnableJpaRepositories("org.thingsboard.server.dao.sqlts.ts") +@EntityScan("org.thingsboard.server.dao.model.sqlts.ts") +@EnableTransactionManagement +@SqlTsDao +public class SqlTsDaoConfig { + +} \ No newline at end of file diff --git a/dao/src/main/java/org/thingsboard/server/dao/TimescaleDaoConfig.java b/dao/src/main/java/org/thingsboard/server/dao/TimescaleDaoConfig.java new file mode 100644 index 0000000000..87982203a4 --- /dev/null +++ b/dao/src/main/java/org/thingsboard/server/dao/TimescaleDaoConfig.java @@ -0,0 +1,35 @@ +/** + * Copyright © 2016-2019 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; + +import org.springframework.boot.autoconfigure.EnableAutoConfiguration; +import org.springframework.boot.autoconfigure.domain.EntityScan; +import org.springframework.context.annotation.ComponentScan; +import org.springframework.context.annotation.Configuration; +import org.springframework.data.jpa.repository.config.EnableJpaRepositories; +import org.springframework.transaction.annotation.EnableTransactionManagement; +import org.thingsboard.server.dao.util.TimescaleDBTsDao; + +@Configuration +@EnableAutoConfiguration +@ComponentScan("org.thingsboard.server.dao.sqlts.timescale") +@EnableJpaRepositories("org.thingsboard.server.dao.sqlts.timescale") +@EntityScan("org.thingsboard.server.dao.model.sqlts.timescale") +@EnableTransactionManagement +@TimescaleDBTsDao +public class TimescaleDaoConfig { + +} \ No newline at end of file diff --git a/dao/src/main/java/org/thingsboard/server/dao/model/ModelConstants.java b/dao/src/main/java/org/thingsboard/server/dao/model/ModelConstants.java index d9cb365a14..f065443489 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/model/ModelConstants.java +++ b/dao/src/main/java/org/thingsboard/server/dao/model/ModelConstants.java @@ -47,6 +47,7 @@ public class ModelConstants { public static final String ENTITY_TYPE_PROPERTY = "entity_type"; public static final String ENTITY_TYPE_COLUMN = ENTITY_TYPE_PROPERTY; + public static final String TENANT_ID_COLUMN = "tenant_id"; public static final String ENTITY_ID_COLUMN = "entity_id"; public static final String ATTRIBUTE_TYPE_COLUMN = "attribute_type"; public static final String ATTRIBUTE_KEY_COLUMN = "attribute_key"; diff --git a/dao/src/main/java/org/thingsboard/server/dao/model/sql/AbsractTsKvEntity.java b/dao/src/main/java/org/thingsboard/server/dao/model/sql/AbsractTsKvEntity.java new file mode 100644 index 0000000000..2773e5e4ab --- /dev/null +++ b/dao/src/main/java/org/thingsboard/server/dao/model/sql/AbsractTsKvEntity.java @@ -0,0 +1,98 @@ +/** + * Copyright © 2016-2019 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.model.sql; + +import lombok.Data; +import org.thingsboard.server.common.data.kv.BasicTsKvEntry; +import org.thingsboard.server.common.data.kv.BooleanDataEntry; +import org.thingsboard.server.common.data.kv.DoubleDataEntry; +import org.thingsboard.server.common.data.kv.KvEntry; +import org.thingsboard.server.common.data.kv.LongDataEntry; +import org.thingsboard.server.common.data.kv.StringDataEntry; +import org.thingsboard.server.common.data.kv.TsKvEntry; +import org.thingsboard.server.dao.model.ToData; + +import javax.persistence.Column; +import javax.persistence.Id; +import javax.persistence.MappedSuperclass; + +import static org.thingsboard.server.dao.model.ModelConstants.BOOLEAN_VALUE_COLUMN; +import static org.thingsboard.server.dao.model.ModelConstants.DOUBLE_VALUE_COLUMN; +import static org.thingsboard.server.dao.model.ModelConstants.ENTITY_ID_COLUMN; +import static org.thingsboard.server.dao.model.ModelConstants.KEY_COLUMN; +import static org.thingsboard.server.dao.model.ModelConstants.LONG_VALUE_COLUMN; +import static org.thingsboard.server.dao.model.ModelConstants.STRING_VALUE_COLUMN; +import static org.thingsboard.server.dao.model.ModelConstants.TS_COLUMN; + +@Data +@MappedSuperclass +public abstract class AbsractTsKvEntity implements ToData { + + protected static final String SUM = "SUM"; + protected static final String AVG = "AVG"; + protected static final String MIN = "MIN"; + protected static final String MAX = "MAX"; + + @Id + @Column(name = ENTITY_ID_COLUMN) + protected String entityId; + + @Id + @Column(name = TS_COLUMN) + protected Long ts; + + @Id + @Column(name = KEY_COLUMN) + protected String key; + + @Column(name = BOOLEAN_VALUE_COLUMN) + protected Boolean booleanValue; + + @Column(name = STRING_VALUE_COLUMN) + protected String strValue; + + @Column(name = LONG_VALUE_COLUMN) + protected Long longValue; + + @Column(name = DOUBLE_VALUE_COLUMN) + protected Double doubleValue; + + @Override + public TsKvEntry toData() { + KvEntry kvEntry = null; + if (strValue != null) { + kvEntry = new StringDataEntry(key, strValue); + } else if (longValue != null) { + kvEntry = new LongDataEntry(key, longValue); + } else if (doubleValue != null) { + kvEntry = new DoubleDataEntry(key, doubleValue); + } else if (booleanValue != null) { + kvEntry = new BooleanDataEntry(key, booleanValue); + } + return new BasicTsKvEntry(ts, kvEntry); + } + + public abstract boolean isNotEmpty(); + + protected static boolean isAllNull(Object... args) { + for (Object arg : args) { + if(arg != null) { + return false; + } + } + return true; + } +} \ No newline at end of file diff --git a/dao/src/main/java/org/thingsboard/server/dao/model/sqlts/timescale/TimescaleTsKvCompositeKey.java b/dao/src/main/java/org/thingsboard/server/dao/model/sqlts/timescale/TimescaleTsKvCompositeKey.java new file mode 100644 index 0000000000..a8e7494627 --- /dev/null +++ b/dao/src/main/java/org/thingsboard/server/dao/model/sqlts/timescale/TimescaleTsKvCompositeKey.java @@ -0,0 +1,37 @@ +/** + * Copyright © 2016-2019 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.model.sqlts.timescale; + +import lombok.AllArgsConstructor; +import lombok.Data; +import lombok.NoArgsConstructor; + +import javax.persistence.Transient; +import java.io.Serializable; + +@Data +@AllArgsConstructor +@NoArgsConstructor +public class TimescaleTsKvCompositeKey implements Serializable { + + @Transient + private static final long serialVersionUID = -4089175869616037523L; + + private String tenantId; + private String entityId; + private String key; + private long ts; +} \ No newline at end of file diff --git a/dao/src/main/java/org/thingsboard/server/dao/model/sqlts/timescale/TimescaleTsKvEntity.java b/dao/src/main/java/org/thingsboard/server/dao/model/sqlts/timescale/TimescaleTsKvEntity.java new file mode 100644 index 0000000000..fa212dd2f9 --- /dev/null +++ b/dao/src/main/java/org/thingsboard/server/dao/model/sqlts/timescale/TimescaleTsKvEntity.java @@ -0,0 +1,184 @@ +/** + * Copyright © 2016-2019 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.model.sqlts.timescale; + +import lombok.Data; +import lombok.EqualsAndHashCode; +import org.springframework.util.StringUtils; +import org.thingsboard.server.common.data.kv.TsKvEntry; +import org.thingsboard.server.dao.model.ToData; +import org.thingsboard.server.dao.model.sql.AbsractTsKvEntity; + +import javax.persistence.Column; +import javax.persistence.ColumnResult; +import javax.persistence.ConstructorResult; +import javax.persistence.Entity; +import javax.persistence.Id; +import javax.persistence.IdClass; +import javax.persistence.NamedNativeQueries; +import javax.persistence.NamedNativeQuery; +import javax.persistence.SqlResultSetMapping; +import javax.persistence.SqlResultSetMappings; +import javax.persistence.Table; + +import static org.thingsboard.server.dao.model.ModelConstants.TENANT_ID_COLUMN; +import static org.thingsboard.server.dao.sqlts.timescale.AggregationRepository.FIND_AVG; +import static org.thingsboard.server.dao.sqlts.timescale.AggregationRepository.FIND_AVG_QUERY; +import static org.thingsboard.server.dao.sqlts.timescale.AggregationRepository.FIND_COUNT; +import static org.thingsboard.server.dao.sqlts.timescale.AggregationRepository.FIND_COUNT_QUERY; +import static org.thingsboard.server.dao.sqlts.timescale.AggregationRepository.FIND_MAX; +import static org.thingsboard.server.dao.sqlts.timescale.AggregationRepository.FIND_MAX_QUERY; +import static org.thingsboard.server.dao.sqlts.timescale.AggregationRepository.FIND_MIN; +import static org.thingsboard.server.dao.sqlts.timescale.AggregationRepository.FIND_MIN_QUERY; +import static org.thingsboard.server.dao.sqlts.timescale.AggregationRepository.FIND_SUM; +import static org.thingsboard.server.dao.sqlts.timescale.AggregationRepository.FIND_SUM_QUERY; +import static org.thingsboard.server.dao.sqlts.timescale.AggregationRepository.FROM_WHERE_CLAUSE; + +@Data +@EqualsAndHashCode(callSuper = true) +@Entity +@Table(name = "tenant_ts_kv") +@IdClass(TimescaleTsKvCompositeKey.class) +@SqlResultSetMappings({ + @SqlResultSetMapping( + name = "timescaleAggregationMapping", + classes = { + @ConstructorResult( + targetClass = TimescaleTsKvEntity.class, + columns = { + @ColumnResult(name = "tsBucket", type = Long.class), + @ColumnResult(name = "interval", type = Long.class), + @ColumnResult(name = "longValue", type = Long.class), + @ColumnResult(name = "doubleValue", type = Double.class), + @ColumnResult(name = "longCountValue", type = Long.class), + @ColumnResult(name = "doubleCountValue", type = Long.class), + @ColumnResult(name = "strValue", type = String.class), + @ColumnResult(name = "aggType", type = String.class), + } + ), + }), + @SqlResultSetMapping( + name = "timescaleCountMapping", + classes = { + @ConstructorResult( + targetClass = TimescaleTsKvEntity.class, + columns = { + @ColumnResult(name = "tsBucket", type = Long.class), + @ColumnResult(name = "interval", type = Long.class), + @ColumnResult(name = "booleanValueCount", type = Long.class), + @ColumnResult(name = "strValueCount", type = Long.class), + @ColumnResult(name = "longValueCount", type = Long.class), + @ColumnResult(name = "doubleValueCount", type = Long.class), + } + ) + }), +}) +@NamedNativeQueries({ + @NamedNativeQuery( + name = FIND_AVG, + query = FIND_AVG_QUERY + FROM_WHERE_CLAUSE, + resultSetMapping = "timescaleAggregationMapping" + ), + @NamedNativeQuery( + name = FIND_MAX, + query = FIND_MAX_QUERY + FROM_WHERE_CLAUSE, + resultSetMapping = "timescaleAggregationMapping" + ), + @NamedNativeQuery( + name = FIND_MIN, + query = FIND_MIN_QUERY + FROM_WHERE_CLAUSE, + resultSetMapping = "timescaleAggregationMapping" + ), + @NamedNativeQuery( + name = FIND_SUM, + query = FIND_SUM_QUERY + FROM_WHERE_CLAUSE, + resultSetMapping = "timescaleAggregationMapping" + ), + @NamedNativeQuery( + name = FIND_COUNT, + query = FIND_COUNT_QUERY + FROM_WHERE_CLAUSE, + resultSetMapping = "timescaleCountMapping" + ) +}) +public final class TimescaleTsKvEntity extends AbsractTsKvEntity implements ToData { + + @Id + @Column(name = TENANT_ID_COLUMN) + private String tenantId; + + public TimescaleTsKvEntity() { } + + public TimescaleTsKvEntity(Long tsBucket, Long interval, Long longValue, Double doubleValue, Long longCountValue, Long doubleCountValue, String strValue, String aggType) { + if (!StringUtils.isEmpty(strValue)) { + this.strValue = strValue; + } + if (!isAllNull(tsBucket, interval, longValue, doubleValue, longCountValue, doubleCountValue)) { + this.ts = tsBucket + interval/2; + switch (aggType) { + case AVG: + double sum = 0.0; + if (longValue != null) { + sum += longValue; + } + if (doubleValue != null) { + sum += doubleValue; + } + long totalCount = longCountValue + doubleCountValue; + if (totalCount > 0) { + this.doubleValue = sum / (longCountValue + doubleCountValue); + } else { + this.doubleValue = 0.0; + } + break; + case SUM: + if (doubleCountValue > 0) { + this.doubleValue = doubleValue + (longValue != null ? longValue.doubleValue() : 0.0); + } else { + this.longValue = longValue; + } + break; + case MIN: + case MAX: + if (longCountValue > 0 && doubleCountValue > 0) { + this.doubleValue = MAX.equals(aggType) ? Math.max(doubleValue, longValue.doubleValue()) : Math.min(doubleValue, longValue.doubleValue()); + } else if (doubleCountValue > 0) { + this.doubleValue = doubleValue; + } else if (longCountValue > 0) { + this.longValue = longValue; + } + break; + } + } + } + + public TimescaleTsKvEntity(Long tsBucket, Long interval, Long booleanValueCount, Long strValueCount, Long longValueCount, Long doubleValueCount) { + if (!isAllNull(tsBucket, interval, booleanValueCount, strValueCount, longValueCount, doubleValueCount)) { + this.ts = tsBucket + interval/2; + if (booleanValueCount != 0) { + this.longValue = booleanValueCount; + } else if (strValueCount != 0) { + this.longValue = strValueCount; + } else { + this.longValue = longValueCount + doubleValueCount; + } + } + } + + @Override + public boolean isNotEmpty() { + return ts != null && (strValue != null || longValue != null || doubleValue != null || booleanValue != null); + } +} \ No newline at end of file diff --git a/dao/src/main/java/org/thingsboard/server/dao/model/sql/TsKvCompositeKey.java b/dao/src/main/java/org/thingsboard/server/dao/model/sqlts/ts/TsKvCompositeKey.java similarity index 95% rename from dao/src/main/java/org/thingsboard/server/dao/model/sql/TsKvCompositeKey.java rename to dao/src/main/java/org/thingsboard/server/dao/model/sqlts/ts/TsKvCompositeKey.java index 67bfe40370..0b7ae78e07 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/model/sql/TsKvCompositeKey.java +++ b/dao/src/main/java/org/thingsboard/server/dao/model/sqlts/ts/TsKvCompositeKey.java @@ -13,7 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.server.dao.model.sql; +package org.thingsboard.server.dao.model.sqlts.ts; import lombok.AllArgsConstructor; import lombok.Data; diff --git a/dao/src/main/java/org/thingsboard/server/dao/model/sql/TsKvEntity.java b/dao/src/main/java/org/thingsboard/server/dao/model/sqlts/ts/TsKvEntity.java similarity index 58% rename from dao/src/main/java/org/thingsboard/server/dao/model/sql/TsKvEntity.java rename to dao/src/main/java/org/thingsboard/server/dao/model/sqlts/ts/TsKvEntity.java index 873f8e8c42..4440a3dd0c 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/model/sql/TsKvEntity.java +++ b/dao/src/main/java/org/thingsboard/server/dao/model/sqlts/ts/TsKvEntity.java @@ -13,18 +13,13 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.server.dao.model.sql; +package org.thingsboard.server.dao.model.sqlts.ts; import lombok.Data; import org.thingsboard.server.common.data.EntityType; -import org.thingsboard.server.common.data.kv.BasicTsKvEntry; -import org.thingsboard.server.common.data.kv.BooleanDataEntry; -import org.thingsboard.server.common.data.kv.DoubleDataEntry; -import org.thingsboard.server.common.data.kv.KvEntry; -import org.thingsboard.server.common.data.kv.LongDataEntry; -import org.thingsboard.server.common.data.kv.StringDataEntry; import org.thingsboard.server.common.data.kv.TsKvEntry; import org.thingsboard.server.dao.model.ToData; +import org.thingsboard.server.dao.model.sql.AbsractTsKvEntity; import javax.persistence.Column; import javax.persistence.Entity; @@ -34,25 +29,18 @@ import javax.persistence.Id; import javax.persistence.IdClass; import javax.persistence.Table; -import static org.thingsboard.server.dao.model.ModelConstants.BOOLEAN_VALUE_COLUMN; -import static org.thingsboard.server.dao.model.ModelConstants.DOUBLE_VALUE_COLUMN; -import static org.thingsboard.server.dao.model.ModelConstants.ENTITY_ID_COLUMN; import static org.thingsboard.server.dao.model.ModelConstants.ENTITY_TYPE_COLUMN; -import static org.thingsboard.server.dao.model.ModelConstants.KEY_COLUMN; -import static org.thingsboard.server.dao.model.ModelConstants.LONG_VALUE_COLUMN; -import static org.thingsboard.server.dao.model.ModelConstants.STRING_VALUE_COLUMN; -import static org.thingsboard.server.dao.model.ModelConstants.TS_COLUMN; @Data @Entity @Table(name = "ts_kv") @IdClass(TsKvCompositeKey.class) -public final class TsKvEntity implements ToData { +public final class TsKvEntity extends AbsractTsKvEntity implements ToData { - private static final String SUM = "SUM"; - private static final String AVG = "AVG"; - private static final String MIN = "MIN"; - private static final String MAX = "MAX"; + @Id + @Enumerated(EnumType.STRING) + @Column(name = ENTITY_TYPE_COLUMN) + private EntityType entityType; public TsKvEntity() { } @@ -62,7 +50,7 @@ public final class TsKvEntity implements ToData { } public TsKvEntity(Long longValue, Double doubleValue, Long longCountValue, Long doubleCountValue, String aggType) { - if(!isAllNull(longValue, doubleValue, longCountValue, doubleCountValue)) { + if (!isAllNull(longValue, doubleValue, longCountValue, doubleCountValue)) { switch (aggType) { case AVG: double sum = 0.0; @@ -101,7 +89,7 @@ public final class TsKvEntity implements ToData { } public TsKvEntity(Long booleanValueCount, Long strValueCount, Long longValueCount, Long doubleValueCount) { - if(!isAllNull(booleanValueCount, strValueCount, longValueCount, doubleValueCount)) { + if (!isAllNull(booleanValueCount, strValueCount, longValueCount, doubleValueCount)) { if (booleanValueCount != 0) { this.longValue = booleanValueCount; } else if (strValueCount != 0) { @@ -112,60 +100,9 @@ public final class TsKvEntity implements ToData { } } - @Id - @Enumerated(EnumType.STRING) - @Column(name = ENTITY_TYPE_COLUMN) - private EntityType entityType; - - @Id - @Column(name = ENTITY_ID_COLUMN) - private String entityId; - - @Id - @Column(name = KEY_COLUMN) - private String key; - - @Id - @Column(name = TS_COLUMN) - private long ts; - - @Column(name = BOOLEAN_VALUE_COLUMN) - private Boolean booleanValue; - - @Column(name = STRING_VALUE_COLUMN) - private String strValue; - - @Column(name = LONG_VALUE_COLUMN) - private Long longValue; - - @Column(name = DOUBLE_VALUE_COLUMN) - private Double doubleValue; @Override - public TsKvEntry toData() { - KvEntry kvEntry = null; - if (strValue != null) { - kvEntry = new StringDataEntry(key, strValue); - } else if (longValue != null) { - kvEntry = new LongDataEntry(key, longValue); - } else if (doubleValue != null) { - kvEntry = new DoubleDataEntry(key, doubleValue); - } else if (booleanValue != null) { - kvEntry = new BooleanDataEntry(key, booleanValue); - } - return new BasicTsKvEntry(ts, kvEntry); - } - public boolean isNotEmpty() { return strValue != null || longValue != null || doubleValue != null || booleanValue != null; } - - private static boolean isAllNull(Object... args) { - for (Object arg : args) { - if(arg != null) { - return false; - } - } - return true; - } } diff --git a/dao/src/main/java/org/thingsboard/server/dao/model/sql/TsKvLatestCompositeKey.java b/dao/src/main/java/org/thingsboard/server/dao/model/sqlts/ts/TsKvLatestCompositeKey.java similarity index 95% rename from dao/src/main/java/org/thingsboard/server/dao/model/sql/TsKvLatestCompositeKey.java rename to dao/src/main/java/org/thingsboard/server/dao/model/sqlts/ts/TsKvLatestCompositeKey.java index 671fc0c662..004efa8b8b 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/model/sql/TsKvLatestCompositeKey.java +++ b/dao/src/main/java/org/thingsboard/server/dao/model/sqlts/ts/TsKvLatestCompositeKey.java @@ -13,7 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.server.dao.model.sql; +package org.thingsboard.server.dao.model.sqlts.ts; import lombok.AllArgsConstructor; import lombok.Data; diff --git a/dao/src/main/java/org/thingsboard/server/dao/model/sql/TsKvLatestEntity.java b/dao/src/main/java/org/thingsboard/server/dao/model/sqlts/ts/TsKvLatestEntity.java similarity index 98% rename from dao/src/main/java/org/thingsboard/server/dao/model/sql/TsKvLatestEntity.java rename to dao/src/main/java/org/thingsboard/server/dao/model/sqlts/ts/TsKvLatestEntity.java index 462abf1ed5..77a7d64fbb 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/model/sql/TsKvLatestEntity.java +++ b/dao/src/main/java/org/thingsboard/server/dao/model/sqlts/ts/TsKvLatestEntity.java @@ -13,7 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.server.dao.model.sql; +package org.thingsboard.server.dao.model.sqlts.ts; import lombok.Data; import org.thingsboard.server.common.data.EntityType; diff --git a/dao/src/main/java/org/thingsboard/server/dao/sqlts/AbstractSqlTimeseriesDao.java b/dao/src/main/java/org/thingsboard/server/dao/sqlts/AbstractSqlTimeseriesDao.java new file mode 100644 index 0000000000..cf3596a864 --- /dev/null +++ b/dao/src/main/java/org/thingsboard/server/dao/sqlts/AbstractSqlTimeseriesDao.java @@ -0,0 +1,130 @@ +/** + * Copyright © 2016-2019 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; + +import com.google.common.base.Function; +import com.google.common.util.concurrent.Futures; +import com.google.common.util.concurrent.ListenableFuture; +import com.google.common.util.concurrent.ListeningExecutorService; +import com.google.common.util.concurrent.MoreExecutors; +import org.springframework.beans.factory.annotation.Value; +import org.thingsboard.server.common.data.id.EntityId; +import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.data.kv.Aggregation; +import org.thingsboard.server.common.data.kv.BaseReadTsKvQuery; +import org.thingsboard.server.common.data.kv.DeleteTsKvQuery; +import org.thingsboard.server.common.data.kv.ReadTsKvQuery; +import org.thingsboard.server.common.data.kv.TsKvEntry; +import org.thingsboard.server.dao.sql.JpaAbstractDaoListeningExecutorService; +import org.thingsboard.server.dao.timeseries.TsInsertExecutorType; + +import javax.annotation.Nullable; +import javax.annotation.PostConstruct; +import javax.annotation.PreDestroy; +import java.util.List; +import java.util.Objects; +import java.util.Optional; +import java.util.concurrent.Executors; +import java.util.stream.Collectors; + +public abstract class AbstractSqlTimeseriesDao extends JpaAbstractDaoListeningExecutorService { + + private static final String DESC_ORDER = "DESC"; + + @Value("${sql.ts_inserts_executor_type}") + private String insertExecutorType; + + @Value("${sql.ts_inserts_fixed_thread_pool_size}") + private int insertFixedThreadPoolSize; + + @Value("${spring.datasource.hikari.maximumPoolSize}") + private int maximumPoolSize; + + protected ListeningExecutorService insertService; + + @PostConstruct + void init() { + Optional executorTypeOptional = TsInsertExecutorType.parse(insertExecutorType); + TsInsertExecutorType executorType; + executorType = executorTypeOptional.orElse(TsInsertExecutorType.FIXED); + switch (executorType) { + case SINGLE: + insertService = MoreExecutors.listeningDecorator(Executors.newSingleThreadExecutor()); + break; + case FIXED: + case CACHED: + int poolSize = insertFixedThreadPoolSize; + if (poolSize <= 0) { + poolSize = maximumPoolSize * 4; + } + insertService = MoreExecutors.listeningDecorator(Executors.newFixedThreadPool(poolSize)); + break; + } + } + + @PreDestroy + void preDestroy() { + if (insertService != null) { + insertService.shutdown(); + } + } + + protected ListenableFuture> processFindAllAsync(TenantId tenantId, EntityId entityId, List queries) { + List>> futures = queries + .stream() + .map(query -> findAllAsync(tenantId, entityId, query)) + .collect(Collectors.toList()); + return Futures.transform(Futures.allAsList(futures), new Function>, List>() { + @Nullable + @Override + public List apply(@Nullable List> results) { + if (results == null || results.isEmpty()) { + return null; + } + return results.stream() + .filter(Objects::nonNull) + .flatMap(List::stream) + .collect(Collectors.toList()); + } + }, service); + } + + protected abstract ListenableFuture> findAllAsync(TenantId tenantId, EntityId entityId, ReadTsKvQuery query); + + protected ListenableFuture> getTskvEntriesFuture(ListenableFuture>> future) { + return Futures.transform(future, new Function>, List>() { + @Nullable + @Override + public List apply(@Nullable List> results) { + if (results == null || results.isEmpty()) { + return null; + } + return results.stream() + .filter(Optional::isPresent) + .map(Optional::get) + .collect(Collectors.toList()); + } + }, service); + } + + protected ListenableFuture> findNewLatestEntryFuture(TenantId tenantId, EntityId entityId, DeleteTsKvQuery query) { + long startTs = 0; + long endTs = query.getStartTs() - 1; + ReadTsKvQuery findNewLatestQuery = new BaseReadTsKvQuery(query.getKey(), startTs, endTs, endTs - startTs, 1, + Aggregation.NONE, DESC_ORDER); + return findAllAsync(tenantId, entityId, findNewLatestQuery); + } +} \ No newline at end of file diff --git a/dao/src/main/java/org/thingsboard/server/dao/sqlts/AbstractTimeseriesInsertRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sqlts/AbstractTimeseriesInsertRepository.java new file mode 100644 index 0000000000..608933c3f7 --- /dev/null +++ b/dao/src/main/java/org/thingsboard/server/dao/sqlts/AbstractTimeseriesInsertRepository.java @@ -0,0 +1,65 @@ +/** + * Copyright © 2016-2019 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; + +import org.springframework.data.jpa.repository.Modifying; +import org.springframework.stereotype.Repository; +import org.thingsboard.server.dao.model.sql.AbsractTsKvEntity; + +import javax.persistence.EntityManager; +import javax.persistence.PersistenceContext; + +@Repository +public abstract class AbstractTimeseriesInsertRepository { + + protected static final String BOOL_V = "bool_v"; + protected static final String STR_V = "str_v"; + protected static final String LONG_V = "long_v"; + protected static final String DBL_V = "dbl_v"; + + @PersistenceContext + protected EntityManager entityManager; + + public abstract void saveOrUpdate(T entity); + + protected void processSaveOrUpdate(T entity, String requestBoolValue, String requestStrValue, String requestLongValue, String requestDblValue) { + if (entity.getBooleanValue() != null) { + saveOrUpdateBoolean(entity, requestBoolValue); + } + if (entity.getStrValue() != null) { + saveOrUpdateString(entity, requestStrValue); + } + if (entity.getLongValue() != null) { + saveOrUpdateLong(entity, requestLongValue); + } + if (entity.getDoubleValue() != null) { + saveOrUpdateDouble(entity, requestDblValue); + } + } + + @Modifying + protected abstract void saveOrUpdateBoolean(T entity, String query); + + @Modifying + protected abstract void saveOrUpdateString(T entity, String query); + + @Modifying + protected abstract void saveOrUpdateLong(T entity, String query); + + @Modifying + protected abstract void saveOrUpdateDouble(T entity, String query); + +} \ No newline at end of file diff --git a/dao/src/main/java/org/thingsboard/server/dao/sqlts/timescale/AggregationRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sqlts/timescale/AggregationRepository.java new file mode 100644 index 0000000000..ab2cf1fff5 --- /dev/null +++ b/dao/src/main/java/org/thingsboard/server/dao/sqlts/timescale/AggregationRepository.java @@ -0,0 +1,101 @@ +/** + * Copyright © 2016-2019 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.timescale; + +import org.springframework.scheduling.annotation.Async; +import org.springframework.stereotype.Repository; +import org.thingsboard.server.dao.model.sqlts.timescale.TimescaleTsKvEntity; +import org.thingsboard.server.dao.util.TimescaleDBTsDao; + +import javax.persistence.EntityManager; +import javax.persistence.PersistenceContext; +import java.util.List; +import java.util.concurrent.CompletableFuture; + +@Repository +@TimescaleDBTsDao +public class AggregationRepository { + + public static final String FIND_AVG = "findAvg"; + public static final String FIND_MAX = "findMax"; + public static final String FIND_MIN = "findMin"; + public static final String FIND_SUM = "findSum"; + public static final String FIND_COUNT = "findCount"; + + + public static final String FROM_WHERE_CLAUSE = "FROM tenant_ts_kv tskv WHERE tskv.tenant_id = cast(:tenantId AS varchar) AND tskv.entity_id = cast(:entityId AS varchar) AND tskv.key= cast(:entityKey AS varchar) AND tskv.ts > :startTs AND tskv.ts <= :endTs GROUP BY tskv.tenant_id, tskv.entity_id, tskv.key, tsBucket ORDER BY tskv.tenant_id, tskv.entity_id, tskv.key, tsBucket"; + + public static final String FIND_AVG_QUERY = "SELECT time_bucket(:timeBucket, tskv.ts) AS tsBucket, :timeBucket AS interval, SUM(COALESCE(tskv.long_v, 0)) AS longValue, SUM(COALESCE(tskv.dbl_v, 0.0)) AS doubleValue, SUM(CASE WHEN tskv.long_v IS NULL THEN 0 ELSE 1 END) AS longCountValue, SUM(CASE WHEN tskv.dbl_v IS NULL THEN 0 ELSE 1 END) AS doubleCountValue, null AS strValue, 'AVG' AS aggType "; + + public static final String FIND_MAX_QUERY = "SELECT time_bucket(:timeBucket, tskv.ts) AS tsBucket, :timeBucket AS interval, MAX(COALESCE(tskv.long_v, -9223372036854775807)) AS longValue, MAX(COALESCE(tskv.dbl_v, -1.79769E+308)) as doubleValue, SUM(CASE WHEN tskv.long_v IS NULL THEN 0 ELSE 1 END) AS longCountValue, SUM(CASE WHEN tskv.dbl_v IS NULL THEN 0 ELSE 1 END) AS doubleCountValue, MAX(tskv.str_v) AS strValue, 'MAX' AS aggType "; + + public static final String FIND_MIN_QUERY = "SELECT time_bucket(:timeBucket, tskv.ts) AS tsBucket, :timeBucket AS interval, MIN(COALESCE(tskv.long_v, 9223372036854775807)) AS longValue, MIN(COALESCE(tskv.dbl_v, 1.79769E+308)) as doubleValue, SUM(CASE WHEN tskv.long_v IS NULL THEN 0 ELSE 1 END) AS longCountValue, SUM(CASE WHEN tskv.dbl_v IS NULL THEN 0 ELSE 1 END) AS doubleCountValue, MIN(tskv.str_v) AS strValue, 'MIN' AS aggType "; + + public static final String FIND_SUM_QUERY = "SELECT time_bucket(:timeBucket, tskv.ts) AS tsBucket, :timeBucket AS interval, SUM(COALESCE(tskv.long_v, 0)) AS longValue, SUM(COALESCE(tskv.dbl_v, 0.0)) AS doubleValue, SUM(CASE WHEN tskv.long_v IS NULL THEN 0 ELSE 1 END) AS longCountValue, SUM(CASE WHEN tskv.dbl_v IS NULL THEN 0 ELSE 1 END) AS doubleCountValue, null AS strValue, 'SUM' AS aggType "; + + public static final String FIND_COUNT_QUERY = "SELECT time_bucket(:timeBucket, tskv.ts) AS tsBucket, :timeBucket AS interval, SUM(CASE WHEN tskv.bool_v IS NULL THEN 0 ELSE 1 END) AS booleanValueCount, SUM(CASE WHEN tskv.str_v IS NULL THEN 0 ELSE 1 END) AS strValueCount, SUM(CASE WHEN tskv.long_v IS NULL THEN 0 ELSE 1 END) AS longValueCount, SUM(CASE WHEN tskv.dbl_v IS NULL THEN 0 ELSE 1 END) AS doubleValueCount "; + + @PersistenceContext + private EntityManager entityManager; + + @Async + public CompletableFuture> findAvg(String tenantId, String entityId, String entityKey, long timeBucket, long startTs, long endTs) { + @SuppressWarnings("unchecked") + List resultList = getResultList(tenantId, entityId, entityKey, timeBucket, startTs, endTs, FIND_AVG); + return CompletableFuture.supplyAsync(() -> resultList); + } + + @Async + public CompletableFuture> findMax(String tenantId, String entityId, String entityKey, long timeBucket, long startTs, long endTs) { + @SuppressWarnings("unchecked") + List resultList = getResultList(tenantId, entityId, entityKey, timeBucket, startTs, endTs, FIND_MAX); + return CompletableFuture.supplyAsync(() -> resultList); + } + + @Async + public CompletableFuture> findMin(String tenantId, String entityId, String entityKey, long timeBucket, long startTs, long endTs) { + @SuppressWarnings("unchecked") + List resultList = getResultList(tenantId, entityId, entityKey, timeBucket, startTs, endTs, FIND_MIN); + return CompletableFuture.supplyAsync(() -> resultList); + } + + @Async + public CompletableFuture> findSum(String tenantId, String entityId, String entityKey, long timeBucket, long startTs, long endTs) { + @SuppressWarnings("unchecked") + List resultList = getResultList(tenantId, entityId, entityKey, timeBucket, startTs, endTs, FIND_SUM); + return CompletableFuture.supplyAsync(() -> resultList); + } + + @Async + public CompletableFuture> findCount(String tenantId, String entityId, String entityKey, long timeBucket, long startTs, long endTs) { + @SuppressWarnings("unchecked") + List resultList = getResultList(tenantId, entityId, entityKey, timeBucket, startTs, endTs, FIND_COUNT); + return CompletableFuture.supplyAsync(() -> resultList); + } + + private List getResultList(String tenantId, String entityId, String entityKey, long timeBucket, long startTs, long endTs, String query) { + return entityManager.createNamedQuery(query) + .setParameter("tenantId", tenantId) + .setParameter("entityId", entityId) + .setParameter("entityKey", entityKey) + .setParameter("timeBucket", timeBucket) + .setParameter("startTs", startTs) + .setParameter("endTs", endTs) + .getResultList(); + } + + +} \ No newline at end of file diff --git a/dao/src/main/java/org/thingsboard/server/dao/sqlts/timescale/TimescaleInsertRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sqlts/timescale/TimescaleInsertRepository.java new file mode 100644 index 0000000000..3c87e9f909 --- /dev/null +++ b/dao/src/main/java/org/thingsboard/server/dao/sqlts/timescale/TimescaleInsertRepository.java @@ -0,0 +1,93 @@ +/** + * Copyright © 2016-2019 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.timescale; + +import org.springframework.stereotype.Repository; +import org.springframework.transaction.annotation.Transactional; +import org.thingsboard.server.dao.model.sqlts.timescale.TimescaleTsKvEntity; +import org.thingsboard.server.dao.sqlts.AbstractTimeseriesInsertRepository; +import org.thingsboard.server.dao.util.PsqlDao; +import org.thingsboard.server.dao.util.TimescaleDBTsDao; + +@TimescaleDBTsDao +@PsqlDao +@Repository +@Transactional +public class TimescaleInsertRepository extends AbstractTimeseriesInsertRepository { + + private static final String ON_BOOL_VALUE_UPDATE_SET_NULLS = "str_v = null, long_v = null, dbl_v = null"; + private static final String ON_STR_VALUE_UPDATE_SET_NULLS = "bool_v = null, long_v = null, dbl_v = null"; + private static final String ON_LONG_VALUE_UPDATE_SET_NULLS = "str_v = null, bool_v = null, dbl_v = null"; + private static final String ON_DBL_VALUE_UPDATE_SET_NULLS = "str_v = null, long_v = null, bool_v = null"; + + private static final String INSERT_OR_UPDATE_BOOL_STATEMENT = getInsertOrUpdateString(BOOL_V, ON_BOOL_VALUE_UPDATE_SET_NULLS); + private static final String INSERT_OR_UPDATE_STR_STATEMENT = getInsertOrUpdateString(STR_V, ON_STR_VALUE_UPDATE_SET_NULLS); + private static final String INSERT_OR_UPDATE_LONG_STATEMENT = getInsertOrUpdateString(LONG_V , ON_LONG_VALUE_UPDATE_SET_NULLS); + private static final String INSERT_OR_UPDATE_DBL_STATEMENT = getInsertOrUpdateString(DBL_V, ON_DBL_VALUE_UPDATE_SET_NULLS); + + @Override + public void saveOrUpdate(TimescaleTsKvEntity entity) { + processSaveOrUpdate(entity, INSERT_OR_UPDATE_BOOL_STATEMENT, INSERT_OR_UPDATE_STR_STATEMENT, INSERT_OR_UPDATE_LONG_STATEMENT, INSERT_OR_UPDATE_DBL_STATEMENT); + } + + @Override + protected void saveOrUpdateBoolean(TimescaleTsKvEntity entity, String query) { + entityManager.createNativeQuery(query) + .setParameter("tenant_id", entity.getTenantId()) + .setParameter("entity_id", entity.getEntityId()) + .setParameter("key", entity.getKey()) + .setParameter("ts", entity.getTs()) + .setParameter("bool_v", entity.getBooleanValue()) + .executeUpdate(); + } + + @Override + protected void saveOrUpdateString(TimescaleTsKvEntity entity, String query) { + entityManager.createNativeQuery(query) + .setParameter("tenant_id", entity.getTenantId()) + .setParameter("entity_id", entity.getEntityId()) + .setParameter("key", entity.getKey()) + .setParameter("ts", entity.getTs()) + .setParameter("str_v", entity.getStrValue()) + .executeUpdate(); + } + + @Override + protected void saveOrUpdateLong(TimescaleTsKvEntity entity, String query) { + entityManager.createNativeQuery(query) + .setParameter("tenant_id", entity.getTenantId()) + .setParameter("entity_id", entity.getEntityId()) + .setParameter("key", entity.getKey()) + .setParameter("ts", entity.getTs()) + .setParameter("long_v", entity.getLongValue()) + .executeUpdate(); + } + + @Override + protected void saveOrUpdateDouble(TimescaleTsKvEntity entity, String query) { + entityManager.createNativeQuery(query) + .setParameter("tenant_id", entity.getTenantId()) + .setParameter("entity_id", entity.getEntityId()) + .setParameter("key", entity.getKey()) + .setParameter("ts", entity.getTs()) + .setParameter("dbl_v", entity.getDoubleValue()) + .executeUpdate(); + } + + private static String getInsertOrUpdateString(String value, String nullValues) { + return "INSERT INTO tenant_ts_kv(tenant_id, entity_id, key, ts, " + value + ") VALUES (:tenant_id, :entity_id, :key, :ts, :" + value + ") ON CONFLICT (tenant_id, entity_id, key, ts) DO UPDATE SET " + value + " = :" + value + ", ts = :ts," + nullValues; + } +} \ No newline at end of file diff --git a/dao/src/main/java/org/thingsboard/server/dao/sqlts/timescale/TimescaleTimeseriesDao.java b/dao/src/main/java/org/thingsboard/server/dao/sqlts/timescale/TimescaleTimeseriesDao.java new file mode 100644 index 0000000000..844f22a31c --- /dev/null +++ b/dao/src/main/java/org/thingsboard/server/dao/sqlts/timescale/TimescaleTimeseriesDao.java @@ -0,0 +1,295 @@ +/** + * Copyright © 2016-2019 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.timescale; + +import com.google.common.collect.Lists; +import com.google.common.util.concurrent.Futures; +import com.google.common.util.concurrent.ListenableFuture; +import com.google.common.util.concurrent.SettableFuture; +import lombok.extern.slf4j.Slf4j; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.data.domain.PageRequest; +import org.springframework.data.domain.Sort; +import org.springframework.stereotype.Component; +import org.springframework.util.CollectionUtils; +import org.thingsboard.server.common.data.id.EntityId; +import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.data.kv.Aggregation; +import org.thingsboard.server.common.data.kv.BasicTsKvEntry; +import org.thingsboard.server.common.data.kv.DeleteTsKvQuery; +import org.thingsboard.server.common.data.kv.ReadTsKvQuery; +import org.thingsboard.server.common.data.kv.StringDataEntry; +import org.thingsboard.server.common.data.kv.TsKvEntry; +import org.thingsboard.server.common.data.kv.TsKvQuery; +import org.thingsboard.server.dao.DaoUtil; +import org.thingsboard.server.dao.model.sqlts.timescale.TimescaleTsKvEntity; +import org.thingsboard.server.dao.sqlts.AbstractSqlTimeseriesDao; +import org.thingsboard.server.dao.sqlts.AbstractTimeseriesInsertRepository; +import org.thingsboard.server.dao.timeseries.TimeseriesDao; +import org.thingsboard.server.dao.util.TimescaleDBTsDao; + +import java.util.ArrayList; +import java.util.Collections; +import java.util.List; +import java.util.Optional; +import java.util.concurrent.CompletableFuture; + +import static org.thingsboard.server.common.data.UUIDConverter.fromTimeUUID; + + +@Component +@Slf4j +@TimescaleDBTsDao +public class TimescaleTimeseriesDao extends AbstractSqlTimeseriesDao implements TimeseriesDao { + + private static final String TS = "ts"; + + @Autowired + private TsKvTimescaleRepository tsKvRepository; + + @Autowired + private AggregationRepository aggregationRepository; + + @Autowired + private AbstractTimeseriesInsertRepository insertRepository; + + @Override + public ListenableFuture> findAllAsync(TenantId tenantId, EntityId entityId, List queries) { + return processFindAllAsync(tenantId, entityId, queries); + } + + protected ListenableFuture> findAllAsync(TenantId tenantId, EntityId entityId, ReadTsKvQuery query) { + if (query.getAggregation() == Aggregation.NONE) { + return findAllAsyncWithLimit(tenantId, entityId, query); + } else { + long startTs = query.getStartTs(); + long endTs = query.getEndTs(); + long timeBucket = query.getInterval(); + ListenableFuture>> future = findAndAggregateAsync(tenantId, entityId, query.getKey(), startTs, endTs, timeBucket, query.getAggregation()); + return getTskvEntriesFuture(future); + } + } + + private ListenableFuture> findAllAsyncWithLimit(TenantId tenantId, EntityId entityId, ReadTsKvQuery query) { + return Futures.immediateFuture( + DaoUtil.convertDataList( + tsKvRepository.findAllWithLimit( + fromTimeUUID(tenantId.getId()), + fromTimeUUID(entityId.getId()), + query.getKey(), + query.getStartTs(), + query.getEndTs(), + new PageRequest(0, query.getLimit(), + new Sort(Sort.Direction.fromString( + query.getOrderBy()), "ts"))))); + } + + + @Override + public ListenableFuture findLatest(TenantId tenantId, EntityId entityId, String key) { + ListenableFuture> future = getLatest(tenantId, entityId, key, 0L, System.currentTimeMillis()); + return Futures.transform(future, latest -> { + if (!CollectionUtils.isEmpty(latest)) { + return DaoUtil.getData(latest.get(0)); + } else { + return new BasicTsKvEntry(System.currentTimeMillis(), new StringDataEntry(key, null)); + } + }, service); + } + + @Override + public ListenableFuture> findAllLatest(TenantId tenantId, EntityId entityId) { + return Futures.immediateFuture(DaoUtil.convertDataList(Lists.newArrayList(tsKvRepository.findAllLatestValues(fromTimeUUID(tenantId.getId()), fromTimeUUID(entityId.getId()))))); + } + + @Override + public ListenableFuture save(TenantId tenantId, EntityId entityId, TsKvEntry tsKvEntry, long ttl) { + TimescaleTsKvEntity entity = new TimescaleTsKvEntity(); + entity.setTenantId(fromTimeUUID(tenantId.getId())); + entity.setEntityId(fromTimeUUID(entityId.getId())); + entity.setTs(tsKvEntry.getTs()); + entity.setKey(tsKvEntry.getKey()); + entity.setStrValue(tsKvEntry.getStrValue().orElse(null)); + entity.setDoubleValue(tsKvEntry.getDoubleValue().orElse(null)); + entity.setLongValue(tsKvEntry.getLongValue().orElse(null)); + entity.setBooleanValue(tsKvEntry.getBooleanValue().orElse(null)); + log.trace("Saving entity to timescale db: {}", entity); + return insertService.submit(() -> { + insertRepository.saveOrUpdate(entity); + return null; + }); + } + + @Override + public ListenableFuture savePartition(TenantId tenantId, EntityId entityId, long tsKvEntryTs, String key, long ttl) { + return insertService.submit(() -> null); + } + + @Override + public ListenableFuture saveLatest(TenantId tenantId, EntityId entityId, TsKvEntry tsKvEntry) { + return insertService.submit(() -> null); + } + + @Override + public ListenableFuture remove(TenantId tenantId, EntityId entityId, DeleteTsKvQuery query) { + return service.submit(() -> { + tsKvRepository.delete( + fromTimeUUID(tenantId.getId()), + fromTimeUUID(entityId.getId()), + query.getKey(), + query.getStartTs(), + query.getEndTs()); + return null; + }); + } + + @Override + public ListenableFuture removeLatest(TenantId tenantId, EntityId entityId, DeleteTsKvQuery query) { + return service.submit(() -> null); + } + + @Override + public ListenableFuture removePartition(TenantId tenantId, EntityId entityId, DeleteTsKvQuery query) { + return service.submit(() -> null); + } + + private ListenableFuture getNewLatestEntryFuture(TenantId tenantId, EntityId entityId, DeleteTsKvQuery query) { + ListenableFuture> future = findNewLatestEntryFuture(tenantId, entityId, query); + return Futures.transformAsync(future, entryList -> { + if (entryList.size() == 1) { + return save(tenantId, entityId, entryList.get(0), 0L); + } else { + log.trace("Could not find new latest value for [{}], key - {}", entityId, query.getKey()); + } + return Futures.immediateFuture(null); + }, service); + } + + private ListenableFuture> findLatestByQuery(TenantId tenantId, EntityId entityId, TsKvQuery query) { + return getLatest(tenantId, entityId, query.getKey(), query.getStartTs(), query.getEndTs()); + } + + private ListenableFuture> getLatest(TenantId tenantId, EntityId entityId, String key, long start, long end) { + return Futures.immediateFuture(tsKvRepository.findAllWithLimit( + fromTimeUUID(tenantId.getId()), + fromTimeUUID(entityId.getId()), + key, + start, + end, + new PageRequest(0, 1, + new Sort(Sort.Direction.DESC, TS)))); + } + + private ListenableFuture>> findAndAggregateAsync(TenantId tenantId, EntityId entityId, String key, long startTs, long endTs, long timeBucket, Aggregation aggregation) { + String entityIdStr = fromTimeUUID(entityId.getId()); + String tenantIdStr = fromTimeUUID(tenantId.getId()); + CompletableFuture> listCompletableFuture = switchAgregation(key, startTs, endTs, timeBucket, aggregation, entityIdStr, tenantIdStr); + SettableFuture> listenableFuture = SettableFuture.create(); + listCompletableFuture.whenComplete((timescaleTsKvEntities, throwable) -> { + if (throwable != null) { + listenableFuture.setException(throwable); + } else { + listenableFuture.set(timescaleTsKvEntities); + } + }); + return Futures.transform(listenableFuture, timescaleTsKvEntities -> { + if (!CollectionUtils.isEmpty(timescaleTsKvEntities)) { + List> result = new ArrayList<>(); + timescaleTsKvEntities.forEach(entity -> { + if(entity != null && entity.isNotEmpty()) { + entity.setEntityId(entityIdStr); + entity.setTenantId(tenantIdStr); + entity.setKey(key); + result.add(Optional.of(DaoUtil.getData(entity))); + } else { + result.add(Optional.empty()); + } + }); + return result; + } else { + return Collections.emptyList(); + } + }); + } + + private CompletableFuture> switchAgregation(String key, long startTs, long endTs, long timeBucket, Aggregation aggregation, String entityIdStr, String tenantIdStr) { + switch (aggregation) { + case AVG: + return findAvg(key, startTs, endTs, timeBucket, entityIdStr, tenantIdStr); + case MAX: + return findMax(key, startTs, endTs, timeBucket, entityIdStr, tenantIdStr); + case MIN: + return findMin(key, startTs, endTs, timeBucket, entityIdStr, tenantIdStr); + case SUM: + return findSum(key, startTs, endTs, timeBucket, entityIdStr, tenantIdStr); + case COUNT: + return findCount(key, startTs, endTs, timeBucket, entityIdStr, tenantIdStr); + default: + throw new IllegalArgumentException("Not supported aggregation type: " + aggregation); + } + } + + private CompletableFuture> findAvg(String key, long startTs, long endTs, long timeBucket, String entityIdStr, String tenantIdStr) { + return aggregationRepository.findAvg( + tenantIdStr, + entityIdStr, + key, + timeBucket, + startTs, + endTs); + } + + private CompletableFuture> findMax(String key, long startTs, long endTs, long timeBucket, String entityIdStr, String tenantIdStr) { + return aggregationRepository.findMax( + tenantIdStr, + entityIdStr, + key, + timeBucket, + startTs, + endTs); + } + + private CompletableFuture> findMin(String key, long startTs, long endTs, long timeBucket, String entityIdStr, String tenantIdStr) { + return aggregationRepository.findMin( + tenantIdStr, + entityIdStr, + key, + timeBucket, + startTs, + endTs); + + } + + private CompletableFuture> findSum(String key, long startTs, long endTs, long timeBucket, String entityIdStr, String tenantIdStr) { + return aggregationRepository.findSum( + tenantIdStr, + entityIdStr, + key, + timeBucket, + startTs, + endTs); + } + + private CompletableFuture> findCount(String key, long startTs, long endTs, long timeBucket, String entityIdStr, String tenantIdStr) { + return aggregationRepository.findCount( + tenantIdStr, + entityIdStr, + key, + timeBucket, + startTs, + endTs); + } +} \ No newline at end of file diff --git a/dao/src/main/java/org/thingsboard/server/dao/sqlts/timescale/TsKvTimescaleRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sqlts/timescale/TsKvTimescaleRepository.java new file mode 100644 index 0000000000..af15e8546e --- /dev/null +++ b/dao/src/main/java/org/thingsboard/server/dao/sqlts/timescale/TsKvTimescaleRepository.java @@ -0,0 +1,66 @@ +/** + * Copyright © 2016-2019 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.timescale; + +import org.springframework.data.domain.Pageable; +import org.springframework.data.jpa.repository.Modifying; +import org.springframework.data.jpa.repository.Query; +import org.springframework.data.repository.CrudRepository; +import org.springframework.data.repository.query.Param; +import org.springframework.transaction.annotation.Transactional; +import org.thingsboard.server.dao.model.sqlts.timescale.TimescaleTsKvCompositeKey; +import org.thingsboard.server.dao.model.sqlts.timescale.TimescaleTsKvEntity; +import org.thingsboard.server.dao.util.TimescaleDBTsDao; + +import java.util.List; + +@TimescaleDBTsDao +public interface TsKvTimescaleRepository extends CrudRepository { + + @Query("SELECT tskv FROM TimescaleTsKvEntity tskv WHERE tskv.tenantId = :tenantId " + + "AND tskv.entityId = :entityId " + + "AND tskv.key = :entityKey " + + "AND tskv.ts > :startTs AND tskv.ts <= :endTs") + List findAllWithLimit( + @Param("tenantId") String tenantId, + @Param("entityId") String entityId, + @Param("entityKey") String key, + @Param("startTs") long startTs, + @Param("endTs") long endTs, Pageable pageable); + + @Query(value = "SELECT tskv.tenant_id as tenant_id, tskv.entity_id as entity_id, tskv.key as key, last(tskv.ts,tskv.ts) as ts," + + " last(tskv.bool_v, tskv.ts) as bool_v, last(tskv.str_v, tskv.ts) as str_v," + + " last(tskv.long_v, tskv.ts) as long_v, last(tskv.dbl_v, tskv.ts) as dbl_v" + + " FROM tenant_ts_kv tskv WHERE tskv.tenant_id = cast(:tenantId AS varchar) " + + "AND tskv.entity_id = cast(:entityId AS varchar) " + + "GROUP BY tskv.tenant_id, tskv.entity_id, tskv.key", nativeQuery = true) + List findAllLatestValues( + @Param("tenantId") String tenantId, + @Param("entityId") String entityId); + + @Transactional + @Modifying + @Query("DELETE FROM TimescaleTsKvEntity tskv WHERE tskv.tenantId = :tenantId " + + "AND tskv.entityId = :entityId " + + "AND tskv.key = :entityKey " + + "AND tskv.ts > :startTs AND tskv.ts <= :endTs") + void delete(@Param("tenantId") String tenantId, + @Param("entityId") String entityId, + @Param("entityKey") String key, + @Param("startTs") long startTs, + @Param("endTs") long endTs); + +} \ No newline at end of file diff --git a/dao/src/main/java/org/thingsboard/server/dao/sqlts/ts/HsqlTimeseriesInsertRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sqlts/ts/HsqlTimeseriesInsertRepository.java new file mode 100644 index 0000000000..3ac5cc67a9 --- /dev/null +++ b/dao/src/main/java/org/thingsboard/server/dao/sqlts/ts/HsqlTimeseriesInsertRepository.java @@ -0,0 +1,93 @@ +/** + * Copyright © 2016-2019 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.ts; + +import org.springframework.stereotype.Repository; +import org.springframework.transaction.annotation.Transactional; +import org.thingsboard.server.dao.model.sqlts.ts.TsKvEntity; +import org.thingsboard.server.dao.sqlts.AbstractTimeseriesInsertRepository; +import org.thingsboard.server.dao.util.HsqlDao; +import org.thingsboard.server.dao.util.SqlTsDao; + +@SqlTsDao +@HsqlDao +@Repository +@Transactional +public class HsqlTimeseriesInsertRepository extends AbstractTimeseriesInsertRepository { + + private static final String ON_BOOL_VALUE_UPDATE_SET_NULLS = " ts_kv.str_v = null, ts_kv.long_v = null, ts_kv.dbl_v = null "; + private static final String ON_STR_VALUE_UPDATE_SET_NULLS = " ts_kv.bool_v = null, ts_kv.long_v = null, ts_kv.dbl_v = null "; + private static final String ON_LONG_VALUE_UPDATE_SET_NULLS = " ts_kv.str_v = null, ts_kv.bool_v = null, ts_kv.dbl_v = null "; + private static final String ON_DBL_VALUE_UPDATE_SET_NULLS = " ts_kv.str_v = null, ts_kv.long_v = null, ts_kv.bool_v = null "; + + private static final String INSERT_OR_UPDATE_BOOL_STATEMENT = getInsertOrUpdateString(BOOL_V, ON_BOOL_VALUE_UPDATE_SET_NULLS); + private static final String INSERT_OR_UPDATE_STR_STATEMENT = getInsertOrUpdateString(STR_V, ON_STR_VALUE_UPDATE_SET_NULLS); + private static final String INSERT_OR_UPDATE_LONG_STATEMENT = getInsertOrUpdateString(LONG_V , ON_LONG_VALUE_UPDATE_SET_NULLS); + private static final String INSERT_OR_UPDATE_DBL_STATEMENT = getInsertOrUpdateString(DBL_V, ON_DBL_VALUE_UPDATE_SET_NULLS); + + private static String getInsertOrUpdateString(String value, String nullValues) { + return "MERGE INTO ts_kv USING(VALUES :entity_type, :entity_id, :key, :ts, :" + value + ") A (entity_type, entity_id, key, ts, " + value + ") ON (ts_kv.entity_type=A.entity_type AND ts_kv.entity_id=A.entity_id AND ts_kv.key=A.key AND ts_kv.ts=A.ts) WHEN MATCHED THEN UPDATE SET ts_kv." + value + " = A." + value + ", ts_kv.ts = A.ts," + nullValues + "WHEN NOT MATCHED THEN INSERT (entity_type, entity_id, key, ts, " + value + ") VALUES (A.entity_type, A.entity_id, A.key, A.ts, A." + value + ")"; + } + + @Override + public void saveOrUpdate(TsKvEntity entity) { + processSaveOrUpdate(entity, INSERT_OR_UPDATE_BOOL_STATEMENT, INSERT_OR_UPDATE_STR_STATEMENT, INSERT_OR_UPDATE_LONG_STATEMENT, INSERT_OR_UPDATE_DBL_STATEMENT); + } + + @Override + protected void saveOrUpdateBoolean(TsKvEntity entity, String query) { + entityManager.createNativeQuery(query) + .setParameter("entity_type", entity.getEntityType().name()) + .setParameter("entity_id", entity.getEntityId()) + .setParameter("key", entity.getKey()) + .setParameter("ts", entity.getTs()) + .setParameter("bool_v", entity.getBooleanValue()) + .executeUpdate(); + } + + @Override + protected void saveOrUpdateString(TsKvEntity entity, String query) { + entityManager.createNativeQuery(query) + .setParameter("entity_type", entity.getEntityType().name()) + .setParameter("entity_id", entity.getEntityId()) + .setParameter("key", entity.getKey()) + .setParameter("ts", entity.getTs()) + .setParameter("str_v", entity.getStrValue()) + .executeUpdate(); + } + + @Override + protected void saveOrUpdateLong(TsKvEntity entity, String query) { + entityManager.createNativeQuery(query) + .setParameter("entity_type", entity.getEntityType().name()) + .setParameter("entity_id", entity.getEntityId()) + .setParameter("key", entity.getKey()) + .setParameter("ts", entity.getTs()) + .setParameter("long_v", entity.getLongValue()) + .executeUpdate(); + } + + @Override + protected void saveOrUpdateDouble(TsKvEntity entity, String query) { + entityManager.createNativeQuery(query) + .setParameter("entity_type", entity.getEntityType().name()) + .setParameter("entity_id", entity.getEntityId()) + .setParameter("key", entity.getKey()) + .setParameter("ts", entity.getTs()) + .setParameter("dbl_v", entity.getDoubleValue()) + .executeUpdate(); + } +} \ No newline at end of file diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/timeseries/JpaTimeseriesDao.java b/dao/src/main/java/org/thingsboard/server/dao/sqlts/ts/JpaTimeseriesDao.java similarity index 62% rename from dao/src/main/java/org/thingsboard/server/dao/sql/timeseries/JpaTimeseriesDao.java rename to dao/src/main/java/org/thingsboard/server/dao/sqlts/ts/JpaTimeseriesDao.java index 2b36eb80b9..9e4e281ebd 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/timeseries/JpaTimeseriesDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sqlts/ts/JpaTimeseriesDao.java @@ -13,19 +13,15 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.server.dao.sql.timeseries; +package org.thingsboard.server.dao.sqlts.ts; -import com.google.common.base.Function; import com.google.common.collect.Lists; import com.google.common.util.concurrent.FutureCallback; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; -import com.google.common.util.concurrent.ListeningExecutorService; -import com.google.common.util.concurrent.MoreExecutors; import com.google.common.util.concurrent.SettableFuture; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; -import org.springframework.beans.factory.annotation.Value; import org.springframework.data.domain.PageRequest; import org.springframework.data.domain.Sort; import org.springframework.stereotype.Component; @@ -33,31 +29,27 @@ import org.thingsboard.server.common.data.UUIDConverter; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.kv.Aggregation; -import org.thingsboard.server.common.data.kv.BaseReadTsKvQuery; import org.thingsboard.server.common.data.kv.BasicTsKvEntry; import org.thingsboard.server.common.data.kv.DeleteTsKvQuery; import org.thingsboard.server.common.data.kv.ReadTsKvQuery; import org.thingsboard.server.common.data.kv.StringDataEntry; import org.thingsboard.server.common.data.kv.TsKvEntry; import org.thingsboard.server.dao.DaoUtil; -import org.thingsboard.server.dao.model.sql.TsKvEntity; -import org.thingsboard.server.dao.model.sql.TsKvLatestCompositeKey; -import org.thingsboard.server.dao.model.sql.TsKvLatestEntity; -import org.thingsboard.server.dao.sql.JpaAbstractDaoListeningExecutorService; +import org.thingsboard.server.dao.model.sqlts.ts.TsKvEntity; +import org.thingsboard.server.dao.model.sqlts.ts.TsKvLatestCompositeKey; +import org.thingsboard.server.dao.model.sqlts.ts.TsKvLatestEntity; +import org.thingsboard.server.dao.sqlts.AbstractSqlTimeseriesDao; +import org.thingsboard.server.dao.sqlts.AbstractTimeseriesInsertRepository; import org.thingsboard.server.dao.timeseries.SimpleListenableFuture; import org.thingsboard.server.dao.timeseries.TimeseriesDao; -import org.thingsboard.server.dao.timeseries.TsInsertExecutorType; import org.thingsboard.server.dao.util.SqlTsDao; import javax.annotation.Nullable; -import javax.annotation.PostConstruct; -import javax.annotation.PreDestroy; import java.util.ArrayList; import java.util.List; import java.util.Optional; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ExecutionException; -import java.util.concurrent.Executors; import java.util.stream.Collectors; import static org.thingsboard.server.common.data.UUIDConverter.fromTimeUUID; @@ -66,17 +58,7 @@ import static org.thingsboard.server.common.data.UUIDConverter.fromTimeUUID; @Component @Slf4j @SqlTsDao -public class JpaTimeseriesDao extends JpaAbstractDaoListeningExecutorService implements TimeseriesDao { - - private static final String DESC_ORDER = "DESC"; - - @Value("${sql.ts_inserts_executor_type}") - private String insertExecutorType; - - @Value("${sql.ts_inserts_fixed_thread_pool_size}") - private int insertFixedThreadPoolSize; - - private ListeningExecutorService insertService; +public class JpaTimeseriesDao extends AbstractSqlTimeseriesDao implements TimeseriesDao { @Autowired private TsKvRepository tsKvRepository; @@ -84,53 +66,15 @@ public class JpaTimeseriesDao extends JpaAbstractDaoListeningExecutorService imp @Autowired private TsKvLatestRepository tsKvLatestRepository; - @PostConstruct - public void init() { - Optional executorTypeOptional = TsInsertExecutorType.parse(insertExecutorType); - TsInsertExecutorType executorType; - if (executorTypeOptional.isPresent()) { - executorType = executorTypeOptional.get(); - } else { - executorType = TsInsertExecutorType.FIXED; - } - switch (executorType) { - case SINGLE: - insertService = MoreExecutors.listeningDecorator(Executors.newSingleThreadExecutor()); - break; - case FIXED: - int poolSize = insertFixedThreadPoolSize; - if (poolSize <= 0) { - poolSize = 10; - } - insertService = MoreExecutors.listeningDecorator(Executors.newFixedThreadPool(poolSize)); - break; - case CACHED: - insertService = MoreExecutors.listeningDecorator(Executors.newCachedThreadPool()); - break; - } - } + @Autowired + private AbstractTimeseriesInsertRepository insertRepository; @Override public ListenableFuture> findAllAsync(TenantId tenantId, EntityId entityId, List queries) { - List>> futures = queries - .stream() - .map(query -> findAllAsync(tenantId, entityId, query)) - .collect(Collectors.toList()); - return Futures.transform(Futures.allAsList(futures), new Function>, List>() { - @Nullable - @Override - public List apply(@Nullable List> results) { - if (results == null || results.isEmpty()) { - return null; - } - return results.stream() - .flatMap(List::stream) - .collect(Collectors.toList()); - } - }, service); + return processFindAllAsync(tenantId, entityId, queries); } - private ListenableFuture> findAllAsync(TenantId tenantId, EntityId entityId, ReadTsKvQuery query) { + protected ListenableFuture> findAllAsync(TenantId tenantId, EntityId entityId, ReadTsKvQuery query) { if (query.getAggregation() == Aggregation.NONE) { return findAllAsyncWithLimit(entityId, query); } else { @@ -140,98 +84,25 @@ public class JpaTimeseriesDao extends JpaAbstractDaoListeningExecutorService imp long startTs = stepTs; long endTs = stepTs + query.getInterval(); long ts = startTs + (endTs - startTs) / 2; - futures.add(findAndAggregateAsync(entityId, query.getKey(), startTs, endTs, ts, query.getAggregation())); + futures.add(findAndAggregateAsync(tenantId, entityId, query.getKey(), startTs, endTs, ts, query.getAggregation())); stepTs = endTs; } - ListenableFuture>> future = Futures.allAsList(futures); - return Futures.transform(future, new Function>, List>() { - @Nullable - @Override - public List apply(@Nullable List> results) { - if (results == null || results.isEmpty()) { - return null; - } - return results.stream() - .filter(Optional::isPresent) - .map(Optional::get) - .collect(Collectors.toList()); - } - }, service); + return getTskvEntriesFuture(Futures.allAsList(futures)); } } - private ListenableFuture> findAndAggregateAsync(EntityId entityId, String key, long startTs, long endTs, long ts, Aggregation aggregation) { + private ListenableFuture> findAndAggregateAsync(TenantId tenantId, EntityId entityId, String key, long startTs, long endTs, long ts, Aggregation aggregation) { List> entitiesFutures = new ArrayList<>(); String entityIdStr = fromTimeUUID(entityId.getId()); - switch (aggregation) { - case AVG: - entitiesFutures.add(tsKvRepository.findAvg( - entityIdStr, - entityId.getEntityType(), - key, - startTs, - endTs)); - - break; - case MAX: - entitiesFutures.add(tsKvRepository.findStringMax( - entityIdStr, - entityId.getEntityType(), - key, - startTs, - endTs)); - entitiesFutures.add(tsKvRepository.findNumericMax( - entityIdStr, - entityId.getEntityType(), - key, - startTs, - endTs)); - - break; - case MIN: - entitiesFutures.add(tsKvRepository.findStringMin( - entityIdStr, - entityId.getEntityType(), - key, - startTs, - endTs)); - entitiesFutures.add(tsKvRepository.findNumericMin( - entityIdStr, - entityId.getEntityType(), - key, - startTs, - endTs)); - break; - case SUM: - entitiesFutures.add(tsKvRepository.findSum( - entityIdStr, - entityId.getEntityType(), - key, - startTs, - endTs)); - break; - case COUNT: - entitiesFutures.add(tsKvRepository.findCount( - entityIdStr, - entityId.getEntityType(), - key, - startTs, - endTs)); - - break; - default: - throw new IllegalArgumentException("Not supported aggregation type: " + aggregation); - } + switchAgregation(entityId, key, startTs, endTs, aggregation, entitiesFutures, entityIdStr); SettableFuture listenableFuture = SettableFuture.create(); - CompletableFuture> entities = CompletableFuture.allOf(entitiesFutures.toArray(new CompletableFuture[entitiesFutures.size()])) - .thenApply(v -> entitiesFutures.stream() - .map(CompletableFuture::join) - .collect(Collectors.toList())); - + .thenApply(v -> entitiesFutures.stream() + .map(CompletableFuture::join) + .collect(Collectors.toList())); entities.whenComplete((tsKvEntities, throwable) -> { if (throwable != null) { @@ -247,22 +118,98 @@ public class JpaTimeseriesDao extends JpaAbstractDaoListeningExecutorService imp listenableFuture.set(result); } }); - return Futures.transform(listenableFuture, new Function>() { - @Override - public Optional apply(@Nullable TsKvEntity entity) { - if (entity != null && entity.isNotEmpty()) { - entity.setEntityId(entityIdStr); - entity.setEntityType(entityId.getEntityType()); - entity.setKey(key); - entity.setTs(ts); - return Optional.of(DaoUtil.getData(entity)); - } else { - return Optional.empty(); - } + return Futures.transform(listenableFuture, entity -> { + if (entity != null && entity.isNotEmpty()) { + entity.setEntityId(entityIdStr); + entity.setEntityType(entityId.getEntityType()); + entity.setKey(key); + entity.setTs(ts); + return Optional.of(DaoUtil.getData(entity)); + } else { + return Optional.empty(); } }); } + private void switchAgregation(EntityId entityId, String key, long startTs, long endTs, Aggregation aggregation, List> entitiesFutures, String entityIdStr) { + switch (aggregation) { + case AVG: + findAvg(entityId, key, startTs, endTs, entitiesFutures, entityIdStr); + break; + case MAX: + findMax(entityId, key, startTs, endTs, entitiesFutures, entityIdStr); + break; + case MIN: + findMin(entityId, key, startTs, endTs, entitiesFutures, entityIdStr); + break; + case SUM: + findSum(entityId, key, startTs, endTs, entitiesFutures, entityIdStr); + break; + case COUNT: + findCount(entityId, key, startTs, endTs, entitiesFutures, entityIdStr); + break; + default: + throw new IllegalArgumentException("Not supported aggregation type: " + aggregation); + } + } + + private void findCount(EntityId entityId, String key, long startTs, long endTs, List> entitiesFutures, String entityIdStr) { + entitiesFutures.add(tsKvRepository.findCount( + entityIdStr, + entityId.getEntityType(), + key, + startTs, + endTs)); + } + + private void findSum(EntityId entityId, String key, long startTs, long endTs, List> entitiesFutures, String entityIdStr) { + entitiesFutures.add(tsKvRepository.findSum( + entityIdStr, + entityId.getEntityType(), + key, + startTs, + endTs)); + } + + private void findMin(EntityId entityId, String key, long startTs, long endTs, List> entitiesFutures, String entityIdStr) { + entitiesFutures.add(tsKvRepository.findStringMin( + entityIdStr, + entityId.getEntityType(), + key, + startTs, + endTs)); + entitiesFutures.add(tsKvRepository.findNumericMin( + entityIdStr, + entityId.getEntityType(), + key, + startTs, + endTs)); + } + + private void findMax(EntityId entityId, String key, long startTs, long endTs, List> entitiesFutures, String entityIdStr) { + entitiesFutures.add(tsKvRepository.findStringMax( + entityIdStr, + entityId.getEntityType(), + key, + startTs, + endTs)); + entitiesFutures.add(tsKvRepository.findNumericMax( + entityIdStr, + entityId.getEntityType(), + key, + startTs, + endTs)); + } + + private void findAvg(EntityId entityId, String key, long startTs, long endTs, List> entitiesFutures, String entityIdStr) { + entitiesFutures.add(tsKvRepository.findAvg( + entityIdStr, + entityId.getEntityType(), + key, + startTs, + endTs)); + } + private ListenableFuture> findAllAsyncWithLimit(EntityId entityId, ReadTsKvQuery query) { return Futures.immediateFuture( DaoUtil.convertDataList( @@ -316,7 +263,7 @@ public class JpaTimeseriesDao extends JpaAbstractDaoListeningExecutorService imp entity.setBooleanValue(tsKvEntry.getBooleanValue().orElse(null)); log.trace("Saving entity: {}", entity); return insertService.submit(() -> { - tsKvRepository.save(entity); + insertRepository.saveOrUpdate(entity); return null; }); } @@ -410,12 +357,7 @@ public class JpaTimeseriesDao extends JpaAbstractDaoListeningExecutorService imp } private ListenableFuture getNewLatestEntryFuture(TenantId tenantId, EntityId entityId, DeleteTsKvQuery query) { - long startTs = 0; - long endTs = query.getStartTs() - 1; - ReadTsKvQuery findNewLatestQuery = new BaseReadTsKvQuery(query.getKey(), startTs, endTs, endTs - startTs, 1, - Aggregation.NONE, DESC_ORDER); - ListenableFuture> future = findAllAsync(tenantId, entityId, findNewLatestQuery); - + ListenableFuture> future = findNewLatestEntryFuture(tenantId, entityId, query); return Futures.transformAsync(future, entryList -> { if (entryList.size() == 1) { return saveLatest(tenantId, entityId, entryList.get(0)); @@ -430,12 +372,4 @@ public class JpaTimeseriesDao extends JpaAbstractDaoListeningExecutorService imp public ListenableFuture removePartition(TenantId tenantId, EntityId entityId, DeleteTsKvQuery query) { return service.submit(() -> null); } - - @PreDestroy - void onDestroy() { - if (insertService != null) { - insertService.shutdown(); - } - } - -} +} \ No newline at end of file diff --git a/dao/src/main/java/org/thingsboard/server/dao/sqlts/ts/PsqlTimeseriesInsertRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sqlts/ts/PsqlTimeseriesInsertRepository.java new file mode 100644 index 0000000000..4ed91c28f2 --- /dev/null +++ b/dao/src/main/java/org/thingsboard/server/dao/sqlts/ts/PsqlTimeseriesInsertRepository.java @@ -0,0 +1,93 @@ +/** + * Copyright © 2016-2019 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.ts; + +import org.springframework.stereotype.Repository; +import org.springframework.transaction.annotation.Transactional; +import org.thingsboard.server.dao.model.sqlts.ts.TsKvEntity; +import org.thingsboard.server.dao.sqlts.AbstractTimeseriesInsertRepository; +import org.thingsboard.server.dao.util.PsqlDao; +import org.thingsboard.server.dao.util.SqlTsDao; + +@SqlTsDao +@PsqlDao +@Repository +@Transactional +public class PsqlTimeseriesInsertRepository extends AbstractTimeseriesInsertRepository { + + private static final String ON_BOOL_VALUE_UPDATE_SET_NULLS = "str_v = null, long_v = null, dbl_v = null"; + private static final String ON_STR_VALUE_UPDATE_SET_NULLS = "bool_v = null, long_v = null, dbl_v = null"; + private static final String ON_LONG_VALUE_UPDATE_SET_NULLS = "str_v = null, bool_v = null, dbl_v = null"; + private static final String ON_DBL_VALUE_UPDATE_SET_NULLS = "str_v = null, long_v = null, bool_v = null"; + + private static final String INSERT_OR_UPDATE_BOOL_STATEMENT = getInsertOrUpdateString(BOOL_V, ON_BOOL_VALUE_UPDATE_SET_NULLS); + private static final String INSERT_OR_UPDATE_STR_STATEMENT = getInsertOrUpdateString(STR_V, ON_STR_VALUE_UPDATE_SET_NULLS); + private static final String INSERT_OR_UPDATE_LONG_STATEMENT = getInsertOrUpdateString(LONG_V , ON_LONG_VALUE_UPDATE_SET_NULLS); + private static final String INSERT_OR_UPDATE_DBL_STATEMENT = getInsertOrUpdateString(DBL_V, ON_DBL_VALUE_UPDATE_SET_NULLS); + + private static String getInsertOrUpdateString(String value, String nullValues) { + return "INSERT INTO ts_kv (entity_type, entity_id, key, ts, " + value + ") VALUES (:entity_type, :entity_id, :key, :ts, :" + value + ") ON CONFLICT (entity_type, entity_id, key, ts) DO UPDATE SET " + value + " = :" + value + ", ts = :ts," + nullValues; + } + + @Override + public void saveOrUpdate(TsKvEntity entity) { + processSaveOrUpdate(entity, INSERT_OR_UPDATE_BOOL_STATEMENT, INSERT_OR_UPDATE_STR_STATEMENT, INSERT_OR_UPDATE_LONG_STATEMENT, INSERT_OR_UPDATE_DBL_STATEMENT); + } + + @Override + protected void saveOrUpdateBoolean(TsKvEntity entity, String query) { + entityManager.createNativeQuery(query) + .setParameter("entity_type", entity.getEntityType().name()) + .setParameter("entity_id", entity.getEntityId()) + .setParameter("key", entity.getKey()) + .setParameter("ts", entity.getTs()) + .setParameter("bool_v", entity.getBooleanValue()) + .executeUpdate(); + } + + @Override + protected void saveOrUpdateString(TsKvEntity entity, String query) { + entityManager.createNativeQuery(query) + .setParameter("entity_type", entity.getEntityType().name()) + .setParameter("entity_id", entity.getEntityId()) + .setParameter("key", entity.getKey()) + .setParameter("ts", entity.getTs()) + .setParameter("str_v", entity.getStrValue()) + .executeUpdate(); + } + + @Override + protected void saveOrUpdateLong(TsKvEntity entity, String query) { + entityManager.createNativeQuery(query) + .setParameter("entity_type", entity.getEntityType().name()) + .setParameter("entity_id", entity.getEntityId()) + .setParameter("key", entity.getKey()) + .setParameter("ts", entity.getTs()) + .setParameter("long_v", entity.getLongValue()) + .executeUpdate(); + } + + @Override + protected void saveOrUpdateDouble(TsKvEntity entity, String query) { + entityManager.createNativeQuery(query) + .setParameter("entity_type", entity.getEntityType().name()) + .setParameter("entity_id", entity.getEntityId()) + .setParameter("key", entity.getKey()) + .setParameter("ts", entity.getTs()) + .setParameter("dbl_v", entity.getDoubleValue()) + .executeUpdate(); + } +} \ No newline at end of file diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/timeseries/TsKvLatestRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sqlts/ts/TsKvLatestRepository.java similarity index 84% rename from dao/src/main/java/org/thingsboard/server/dao/sql/timeseries/TsKvLatestRepository.java rename to dao/src/main/java/org/thingsboard/server/dao/sqlts/ts/TsKvLatestRepository.java index f378e1c3cf..28ef8bfdf5 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/timeseries/TsKvLatestRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sqlts/ts/TsKvLatestRepository.java @@ -13,12 +13,12 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.server.dao.sql.timeseries; +package org.thingsboard.server.dao.sqlts.ts; import org.springframework.data.repository.CrudRepository; import org.thingsboard.server.common.data.EntityType; -import org.thingsboard.server.dao.model.sql.TsKvLatestCompositeKey; -import org.thingsboard.server.dao.model.sql.TsKvLatestEntity; +import org.thingsboard.server.dao.model.sqlts.ts.TsKvLatestCompositeKey; +import org.thingsboard.server.dao.model.sqlts.ts.TsKvLatestEntity; import org.thingsboard.server.dao.util.SqlDao; import java.util.List; diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/timeseries/TsKvRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sqlts/ts/TsKvRepository.java similarity index 97% rename from dao/src/main/java/org/thingsboard/server/dao/sql/timeseries/TsKvRepository.java rename to dao/src/main/java/org/thingsboard/server/dao/sqlts/ts/TsKvRepository.java index 536f367359..5787087a0c 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/timeseries/TsKvRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sqlts/ts/TsKvRepository.java @@ -13,7 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.server.dao.sql.timeseries; +package org.thingsboard.server.dao.sqlts.ts; import org.springframework.data.domain.Pageable; import org.springframework.data.jpa.repository.Modifying; @@ -23,8 +23,8 @@ import org.springframework.data.repository.query.Param; import org.springframework.scheduling.annotation.Async; import org.springframework.transaction.annotation.Transactional; import org.thingsboard.server.common.data.EntityType; -import org.thingsboard.server.dao.model.sql.TsKvCompositeKey; -import org.thingsboard.server.dao.model.sql.TsKvEntity; +import org.thingsboard.server.dao.model.sqlts.ts.TsKvCompositeKey; +import org.thingsboard.server.dao.model.sqlts.ts.TsKvEntity; import org.thingsboard.server.dao.util.SqlDao; import java.util.List; diff --git a/dao/src/main/java/org/thingsboard/server/dao/util/TimescaleDBTsDao.java b/dao/src/main/java/org/thingsboard/server/dao/util/TimescaleDBTsDao.java new file mode 100644 index 0000000000..4302fd09fb --- /dev/null +++ b/dao/src/main/java/org/thingsboard/server/dao/util/TimescaleDBTsDao.java @@ -0,0 +1,22 @@ +/** + * Copyright © 2016-2019 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.util; + +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; + +@ConditionalOnProperty(prefix = "database.ts", value = "type", havingValue = "timescale") +public @interface TimescaleDBTsDao { +} diff --git a/dao/src/main/resources/sql/schema-timescale-idx.sql b/dao/src/main/resources/sql/schema-timescale-idx.sql new file mode 100644 index 0000000000..6dc045f0fa --- /dev/null +++ b/dao/src/main/resources/sql/schema-timescale-idx.sql @@ -0,0 +1,17 @@ +-- +-- Copyright © 2016-2019 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. +-- + +CREATE INDEX IF NOT EXISTS idx_tenant_ts_kv ON tenant_ts_kv(tenant_id, entity_id, key, ts); \ No newline at end of file diff --git a/dao/src/main/resources/sql/schema-timescale.sql b/dao/src/main/resources/sql/schema-timescale.sql new file mode 100644 index 0000000000..d5ac5775f8 --- /dev/null +++ b/dao/src/main/resources/sql/schema-timescale.sql @@ -0,0 +1,31 @@ +-- +-- Copyright © 2016-2019 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. +-- + +CREATE EXTENSION IF NOT EXISTS timescaledb CASCADE; + +CREATE TABLE IF NOT EXISTS tenant_ts_kv ( + tenant_id varchar(31) NOT NULL, + entity_id varchar(31) NOT NULL, + key varchar(255) NOT NULL, + ts bigint NOT NULL, + bool_v boolean, + str_v varchar(10000000), + long_v bigint, + dbl_v double precision, + CONSTRAINT ts_kv_pkey PRIMARY KEY (tenant_id, entity_id, key, ts) +); + +SELECT create_hypertable('tenant_ts_kv', 'ts', chunk_time_interval => 86400000, if_not_exists => true); \ No newline at end of file diff --git a/dao/src/test/java/org/thingsboard/server/dao/AbstractJpaDaoTest.java b/dao/src/test/java/org/thingsboard/server/dao/AbstractJpaDaoTest.java index b4e528cbe9..85ca613223 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/AbstractJpaDaoTest.java +++ b/dao/src/test/java/org/thingsboard/server/dao/AbstractJpaDaoTest.java @@ -30,7 +30,7 @@ import org.springframework.test.context.support.DirtiesContextTestExecutionListe * Created by Valerii Sosliuk on 4/22/2017. */ @RunWith(SpringRunner.class) -@ContextConfiguration(classes = {JpaDaoConfig.class, JpaDbunitTestConfig.class}) +@ContextConfiguration(classes = {JpaDaoConfig.class, SqlTsDaoConfig.class, JpaDbunitTestConfig.class}) @TestPropertySource("classpath:sql-test.properties") @TestExecutionListeners({ DependencyInjectionTestExecutionListener.class, diff --git a/dao/src/test/java/org/thingsboard/server/dao/service/timeseries/BaseTimeseriesServiceTest.java b/dao/src/test/java/org/thingsboard/server/dao/service/timeseries/BaseTimeseriesServiceTest.java index fedd95ed94..0135b2a308 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/service/timeseries/BaseTimeseriesServiceTest.java +++ b/dao/src/test/java/org/thingsboard/server/dao/service/timeseries/BaseTimeseriesServiceTest.java @@ -201,26 +201,6 @@ public abstract class BaseTimeseriesServiceTest extends AbstractServiceTest { Assert.assertEquals(toTsEntry(TS - 2, stringKvEntry), entries.get(2)); } - @Test - public void testDeleteDeviceTsDataWithoutOverwritingLatest() throws Exception { - DeviceId deviceId = new DeviceId(UUIDs.timeBased()); - - saveEntries(deviceId, 10000); - saveEntries(deviceId, 20000); - saveEntries(deviceId, 30000); - saveEntries(deviceId, 40000); - - tsService.remove(tenantId, deviceId, Collections.singletonList( - new BaseDeleteTsKvQuery(STRING_KEY, 15000, 45000))).get(); - - List list = tsService.findAll(tenantId, deviceId, Collections.singletonList( - new BaseReadTsKvQuery(STRING_KEY, 5000, 45000, 10000, 10, Aggregation.NONE))).get(); - Assert.assertEquals(1, list.size()); - - List latest = tsService.findLatest(tenantId, deviceId, Collections.singletonList(STRING_KEY)).get(); - Assert.assertEquals(null, latest.get(0).getValueAsString()); - } - @Test public void testDeleteDeviceTsDataWithOverwritingLatest() throws Exception { DeviceId deviceId = new DeviceId(UUIDs.timeBased()); diff --git a/dao/src/test/resources/sql-test.properties b/dao/src/test/resources/sql-test.properties index 745aa9e1e0..a2fd6bb17b 100644 --- a/dao/src/test/resources/sql-test.properties +++ b/dao/src/test/resources/sql-test.properties @@ -12,4 +12,5 @@ spring.jpa.database-platform=org.hibernate.dialect.HSQLDialect spring.datasource.username=sa spring.datasource.password= spring.datasource.url=jdbc:hsqldb:file:/tmp/testDb;sql.enforce_size=false -spring.datasource.driverClassName=org.hsqldb.jdbc.JDBCDriver \ No newline at end of file +spring.datasource.driverClassName=org.hsqldb.jdbc.JDBCDriver +spring.datasource.hikari.maximumPoolSize = 50 \ No newline at end of file diff --git a/dao/src/test/resources/sql/timescale/drop-all-tables.sql b/dao/src/test/resources/sql/timescale/drop-all-tables.sql new file mode 100644 index 0000000000..e50d307c8e --- /dev/null +++ b/dao/src/test/resources/sql/timescale/drop-all-tables.sql @@ -0,0 +1,21 @@ +DROP TABLE IF EXISTS admin_settings; +DROP TABLE IF EXISTS alarm; +DROP TABLE IF EXISTS asset; +DROP TABLE IF EXISTS audit_log; +DROP TABLE IF EXISTS attribute_kv; +DROP TABLE IF EXISTS component_descriptor; +DROP TABLE IF EXISTS customer; +DROP TABLE IF EXISTS dashboard; +DROP TABLE IF EXISTS device; +DROP TABLE IF EXISTS device_credentials; +DROP TABLE IF EXISTS event; +DROP TABLE IF EXISTS relation; +DROP TABLE IF EXISTS tb_user; +DROP TABLE IF EXISTS tenant; +DROP TABLE IF EXISTS tenant_ts_kv; +DROP TABLE IF EXISTS user_credentials; +DROP TABLE IF EXISTS widget_type; +DROP TABLE IF EXISTS widgets_bundle; +DROP TABLE IF EXISTS rule_node; +DROP TABLE IF EXISTS rule_chain; +DROP TABLE IF EXISTS entity_view; \ No newline at end of file