Browse Source
* 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 typopull/2053/head
committed by
Igor Kulikov
32 changed files with 1586 additions and 283 deletions
@ -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"); |
|||
} |
|||
} |
|||
@ -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 { |
|||
|
|||
} |
|||
@ -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 { |
|||
|
|||
} |
|||
@ -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<TsKvEntry> { |
|||
|
|||
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; |
|||
} |
|||
} |
|||
@ -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; |
|||
} |
|||
@ -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<TsKvEntry> { |
|||
|
|||
@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); |
|||
} |
|||
} |
|||
@ -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<TsInsertExecutorType> 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<List<TsKvEntry>> processFindAllAsync(TenantId tenantId, EntityId entityId, List<ReadTsKvQuery> queries) { |
|||
List<ListenableFuture<List<TsKvEntry>>> futures = queries |
|||
.stream() |
|||
.map(query -> findAllAsync(tenantId, entityId, query)) |
|||
.collect(Collectors.toList()); |
|||
return Futures.transform(Futures.allAsList(futures), new Function<List<List<TsKvEntry>>, List<TsKvEntry>>() { |
|||
@Nullable |
|||
@Override |
|||
public List<TsKvEntry> apply(@Nullable List<List<TsKvEntry>> results) { |
|||
if (results == null || results.isEmpty()) { |
|||
return null; |
|||
} |
|||
return results.stream() |
|||
.filter(Objects::nonNull) |
|||
.flatMap(List::stream) |
|||
.collect(Collectors.toList()); |
|||
} |
|||
}, service); |
|||
} |
|||
|
|||
protected abstract ListenableFuture<List<TsKvEntry>> findAllAsync(TenantId tenantId, EntityId entityId, ReadTsKvQuery query); |
|||
|
|||
protected ListenableFuture<List<TsKvEntry>> getTskvEntriesFuture(ListenableFuture<List<Optional<TsKvEntry>>> future) { |
|||
return Futures.transform(future, new Function<List<Optional<TsKvEntry>>, List<TsKvEntry>>() { |
|||
@Nullable |
|||
@Override |
|||
public List<TsKvEntry> apply(@Nullable List<Optional<TsKvEntry>> results) { |
|||
if (results == null || results.isEmpty()) { |
|||
return null; |
|||
} |
|||
return results.stream() |
|||
.filter(Optional::isPresent) |
|||
.map(Optional::get) |
|||
.collect(Collectors.toList()); |
|||
} |
|||
}, service); |
|||
} |
|||
|
|||
protected ListenableFuture<List<TsKvEntry>> 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); |
|||
} |
|||
} |
|||
@ -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<T extends AbsractTsKvEntity> { |
|||
|
|||
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); |
|||
|
|||
} |
|||
@ -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<List<TimescaleTsKvEntity>> findAvg(String tenantId, String entityId, String entityKey, long timeBucket, long startTs, long endTs) { |
|||
@SuppressWarnings("unchecked") |
|||
List<TimescaleTsKvEntity> resultList = getResultList(tenantId, entityId, entityKey, timeBucket, startTs, endTs, FIND_AVG); |
|||
return CompletableFuture.supplyAsync(() -> resultList); |
|||
} |
|||
|
|||
@Async |
|||
public CompletableFuture<List<TimescaleTsKvEntity>> findMax(String tenantId, String entityId, String entityKey, long timeBucket, long startTs, long endTs) { |
|||
@SuppressWarnings("unchecked") |
|||
List<TimescaleTsKvEntity> resultList = getResultList(tenantId, entityId, entityKey, timeBucket, startTs, endTs, FIND_MAX); |
|||
return CompletableFuture.supplyAsync(() -> resultList); |
|||
} |
|||
|
|||
@Async |
|||
public CompletableFuture<List<TimescaleTsKvEntity>> findMin(String tenantId, String entityId, String entityKey, long timeBucket, long startTs, long endTs) { |
|||
@SuppressWarnings("unchecked") |
|||
List<TimescaleTsKvEntity> resultList = getResultList(tenantId, entityId, entityKey, timeBucket, startTs, endTs, FIND_MIN); |
|||
return CompletableFuture.supplyAsync(() -> resultList); |
|||
} |
|||
|
|||
@Async |
|||
public CompletableFuture<List<TimescaleTsKvEntity>> findSum(String tenantId, String entityId, String entityKey, long timeBucket, long startTs, long endTs) { |
|||
@SuppressWarnings("unchecked") |
|||
List<TimescaleTsKvEntity> resultList = getResultList(tenantId, entityId, entityKey, timeBucket, startTs, endTs, FIND_SUM); |
|||
return CompletableFuture.supplyAsync(() -> resultList); |
|||
} |
|||
|
|||
@Async |
|||
public CompletableFuture<List<TimescaleTsKvEntity>> findCount(String tenantId, String entityId, String entityKey, long timeBucket, long startTs, long endTs) { |
|||
@SuppressWarnings("unchecked") |
|||
List<TimescaleTsKvEntity> 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(); |
|||
} |
|||
|
|||
|
|||
} |
|||
@ -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<TimescaleTsKvEntity> { |
|||
|
|||
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; |
|||
} |
|||
} |
|||
@ -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<List<TsKvEntry>> findAllAsync(TenantId tenantId, EntityId entityId, List<ReadTsKvQuery> queries) { |
|||
return processFindAllAsync(tenantId, entityId, queries); |
|||
} |
|||
|
|||
protected ListenableFuture<List<TsKvEntry>> 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<List<Optional<TsKvEntry>>> future = findAndAggregateAsync(tenantId, entityId, query.getKey(), startTs, endTs, timeBucket, query.getAggregation()); |
|||
return getTskvEntriesFuture(future); |
|||
} |
|||
} |
|||
|
|||
private ListenableFuture<List<TsKvEntry>> 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<TsKvEntry> findLatest(TenantId tenantId, EntityId entityId, String key) { |
|||
ListenableFuture<List<TimescaleTsKvEntity>> 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<List<TsKvEntry>> findAllLatest(TenantId tenantId, EntityId entityId) { |
|||
return Futures.immediateFuture(DaoUtil.convertDataList(Lists.newArrayList(tsKvRepository.findAllLatestValues(fromTimeUUID(tenantId.getId()), fromTimeUUID(entityId.getId()))))); |
|||
} |
|||
|
|||
@Override |
|||
public ListenableFuture<Void> 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<Void> savePartition(TenantId tenantId, EntityId entityId, long tsKvEntryTs, String key, long ttl) { |
|||
return insertService.submit(() -> null); |
|||
} |
|||
|
|||
@Override |
|||
public ListenableFuture<Void> saveLatest(TenantId tenantId, EntityId entityId, TsKvEntry tsKvEntry) { |
|||
return insertService.submit(() -> null); |
|||
} |
|||
|
|||
@Override |
|||
public ListenableFuture<Void> 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<Void> removeLatest(TenantId tenantId, EntityId entityId, DeleteTsKvQuery query) { |
|||
return service.submit(() -> null); |
|||
} |
|||
|
|||
@Override |
|||
public ListenableFuture<Void> removePartition(TenantId tenantId, EntityId entityId, DeleteTsKvQuery query) { |
|||
return service.submit(() -> null); |
|||
} |
|||
|
|||
private ListenableFuture<Void> getNewLatestEntryFuture(TenantId tenantId, EntityId entityId, DeleteTsKvQuery query) { |
|||
ListenableFuture<List<TsKvEntry>> 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<List<TimescaleTsKvEntity>> findLatestByQuery(TenantId tenantId, EntityId entityId, TsKvQuery query) { |
|||
return getLatest(tenantId, entityId, query.getKey(), query.getStartTs(), query.getEndTs()); |
|||
} |
|||
|
|||
private ListenableFuture<List<TimescaleTsKvEntity>> 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<List<Optional<TsKvEntry>>> 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<List<TimescaleTsKvEntity>> listCompletableFuture = switchAgregation(key, startTs, endTs, timeBucket, aggregation, entityIdStr, tenantIdStr); |
|||
SettableFuture<List<TimescaleTsKvEntity>> 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<Optional<TsKvEntry>> 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<List<TimescaleTsKvEntity>> 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<List<TimescaleTsKvEntity>> findAvg(String key, long startTs, long endTs, long timeBucket, String entityIdStr, String tenantIdStr) { |
|||
return aggregationRepository.findAvg( |
|||
tenantIdStr, |
|||
entityIdStr, |
|||
key, |
|||
timeBucket, |
|||
startTs, |
|||
endTs); |
|||
} |
|||
|
|||
private CompletableFuture<List<TimescaleTsKvEntity>> findMax(String key, long startTs, long endTs, long timeBucket, String entityIdStr, String tenantIdStr) { |
|||
return aggregationRepository.findMax( |
|||
tenantIdStr, |
|||
entityIdStr, |
|||
key, |
|||
timeBucket, |
|||
startTs, |
|||
endTs); |
|||
} |
|||
|
|||
private CompletableFuture<List<TimescaleTsKvEntity>> findMin(String key, long startTs, long endTs, long timeBucket, String entityIdStr, String tenantIdStr) { |
|||
return aggregationRepository.findMin( |
|||
tenantIdStr, |
|||
entityIdStr, |
|||
key, |
|||
timeBucket, |
|||
startTs, |
|||
endTs); |
|||
|
|||
} |
|||
|
|||
private CompletableFuture<List<TimescaleTsKvEntity>> findSum(String key, long startTs, long endTs, long timeBucket, String entityIdStr, String tenantIdStr) { |
|||
return aggregationRepository.findSum( |
|||
tenantIdStr, |
|||
entityIdStr, |
|||
key, |
|||
timeBucket, |
|||
startTs, |
|||
endTs); |
|||
} |
|||
|
|||
private CompletableFuture<List<TimescaleTsKvEntity>> findCount(String key, long startTs, long endTs, long timeBucket, String entityIdStr, String tenantIdStr) { |
|||
return aggregationRepository.findCount( |
|||
tenantIdStr, |
|||
entityIdStr, |
|||
key, |
|||
timeBucket, |
|||
startTs, |
|||
endTs); |
|||
} |
|||
} |
|||
@ -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<TimescaleTsKvEntity, TimescaleTsKvCompositeKey> { |
|||
|
|||
@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<TimescaleTsKvEntity> 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<TimescaleTsKvEntity> 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); |
|||
|
|||
} |
|||
@ -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<TsKvEntity> { |
|||
|
|||
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(); |
|||
} |
|||
} |
|||
@ -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<TsKvEntity> { |
|||
|
|||
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(); |
|||
} |
|||
} |
|||
@ -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 { |
|||
} |
|||
@ -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); |
|||
@ -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); |
|||
@ -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; |
|||
Loading…
Reference in new issue