Browse Source
* created annotations for ts latest and created LatestDatabaseSchemaService for cassandra, hsql, psql, timescale * created ts latest dao * fixed tests and refactored * refactored * refactored * created latest dao configs * fix SqlTimeseriesDaoConfig * refactoring, and deleted annotation @SqlDao * refactoring * created migration script for ts latest * refactoring Co-authored-by: Andrew Shvayka <ashvayka@thingsboard.io>pull/3217/head
committed by
GitHub
116 changed files with 1623 additions and 1039 deletions
@ -0,0 +1,30 @@ |
|||
/** |
|||
* Copyright © 2016-2020 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.NoSqlTsLatestDao; |
|||
|
|||
@Service |
|||
@NoSqlTsLatestDao |
|||
@Profile("install") |
|||
public class CassandraTsLatestDatabaseSchemaService extends CassandraAbstractDatabaseSchemaService |
|||
implements TsLatestDatabaseSchemaService { |
|||
public CassandraTsLatestDatabaseSchemaService() { |
|||
super("schema-ts-latest.cql"); |
|||
} |
|||
} |
|||
@ -0,0 +1,19 @@ |
|||
/** |
|||
* Copyright © 2016-2020 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; |
|||
|
|||
public interface TsLatestDatabaseSchemaService extends DatabaseSchemaService { |
|||
} |
|||
@ -0,0 +1,205 @@ |
|||
/** |
|||
* Copyright © 2016-2020 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.migrate; |
|||
|
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.apache.commons.lang3.StringUtils; |
|||
import org.hibernate.exception.ConstraintViolationException; |
|||
import org.springframework.beans.factory.annotation.Autowired; |
|||
import org.springframework.beans.factory.annotation.Value; |
|||
import org.springframework.context.annotation.Profile; |
|||
import org.springframework.stereotype.Service; |
|||
import org.thingsboard.server.common.data.UUIDConverter; |
|||
import org.thingsboard.server.dao.cassandra.CassandraCluster; |
|||
import org.thingsboard.server.dao.model.sqlts.dictionary.TsKvDictionary; |
|||
import org.thingsboard.server.dao.model.sqlts.dictionary.TsKvDictionaryCompositeKey; |
|||
import org.thingsboard.server.dao.model.sqlts.latest.TsKvLatestEntity; |
|||
import org.thingsboard.server.dao.sqlts.dictionary.TsKvDictionaryRepository; |
|||
import org.thingsboard.server.dao.sqlts.insert.latest.InsertLatestTsRepository; |
|||
import org.thingsboard.server.dao.util.NoSqlAnyDao; |
|||
import org.thingsboard.server.service.install.EntityDatabaseSchemaService; |
|||
|
|||
import java.sql.Connection; |
|||
import java.sql.DriverManager; |
|||
import java.util.Arrays; |
|||
import java.util.List; |
|||
import java.util.Optional; |
|||
import java.util.concurrent.ConcurrentHashMap; |
|||
import java.util.concurrent.ConcurrentMap; |
|||
import java.util.concurrent.locks.ReentrantLock; |
|||
import java.util.stream.Collectors; |
|||
|
|||
import static org.thingsboard.server.service.install.migrate.CassandraToSqlColumn.bigintColumn; |
|||
import static org.thingsboard.server.service.install.migrate.CassandraToSqlColumn.booleanColumn; |
|||
import static org.thingsboard.server.service.install.migrate.CassandraToSqlColumn.doubleColumn; |
|||
import static org.thingsboard.server.service.install.migrate.CassandraToSqlColumn.idColumn; |
|||
import static org.thingsboard.server.service.install.migrate.CassandraToSqlColumn.jsonColumn; |
|||
import static org.thingsboard.server.service.install.migrate.CassandraToSqlColumn.stringColumn; |
|||
|
|||
@Service |
|||
@Profile("install") |
|||
@NoSqlAnyDao |
|||
@Slf4j |
|||
public class CassandraTsLatestToSqlMigrateService implements TsLatestMigrateService { |
|||
|
|||
@Autowired |
|||
private EntityDatabaseSchemaService entityDatabaseSchemaService; |
|||
|
|||
@Autowired |
|||
private InsertLatestTsRepository insertLatestTsRepository; |
|||
|
|||
@Autowired |
|||
protected CassandraCluster cluster; |
|||
|
|||
@Autowired |
|||
protected TsKvDictionaryRepository dictionaryRepository; |
|||
|
|||
@Value("${spring.datasource.url}") |
|||
protected String dbUrl; |
|||
|
|||
@Value("${spring.datasource.username}") |
|||
protected String dbUserName; |
|||
|
|||
@Value("${spring.datasource.password}") |
|||
protected String dbPassword; |
|||
|
|||
private final ConcurrentMap<String, Integer> tsKvDictionaryMap = new ConcurrentHashMap<>(); |
|||
|
|||
protected static final ReentrantLock tsCreationLock = new ReentrantLock(); |
|||
|
|||
@Override |
|||
public void migrate() throws Exception { |
|||
log.info("Performing migration of latest timeseries data from cassandra to SQL database ..."); |
|||
entityDatabaseSchemaService.createDatabaseSchema(false); |
|||
try (Connection conn = DriverManager.getConnection(dbUrl, dbUserName, dbPassword)) { |
|||
conn.setAutoCommit(false); |
|||
for (CassandraToSqlTable table : tables) { |
|||
table.migrateToSql(cluster.getSession(), conn); |
|||
} |
|||
} catch (Exception e) { |
|||
log.error("Unexpected error during ThingsBoard entities data migration!", e); |
|||
throw e; |
|||
} |
|||
entityDatabaseSchemaService.createDatabaseIndexes(); |
|||
} |
|||
|
|||
private List<CassandraToSqlTable> tables = Arrays.asList( |
|||
new CassandraToSqlTable("ts_kv_latest_cf", |
|||
idColumn("entity_id"), |
|||
stringColumn("key"), |
|||
bigintColumn("ts"), |
|||
booleanColumn("bool_v"), |
|||
stringColumn("str_v"), |
|||
bigintColumn("long_v"), |
|||
doubleColumn("dbl_v"), |
|||
jsonColumn("json_v")) { |
|||
|
|||
@Override |
|||
protected void batchInsert(List<CassandraToSqlColumnData[]> batchData, Connection conn) { |
|||
insertLatestTsRepository |
|||
.saveOrUpdate(batchData.stream().map(data -> getTsKvLatestEntity(data)).collect(Collectors.toList())); |
|||
} |
|||
|
|||
@Override |
|||
protected CassandraToSqlColumnData[] validateColumnData(CassandraToSqlColumnData[] data) { |
|||
return data; |
|||
} |
|||
}); |
|||
|
|||
private TsKvLatestEntity getTsKvLatestEntity(CassandraToSqlColumnData[] data) { |
|||
TsKvLatestEntity latestEntity = new TsKvLatestEntity(); |
|||
latestEntity.setEntityId(UUIDConverter.fromString(data[0].getValue())); |
|||
latestEntity.setKey(getOrSaveKeyId(data[1].getValue())); |
|||
latestEntity.setTs(Long.parseLong(data[2].getValue())); |
|||
|
|||
String strV = data[4].getValue(); |
|||
if (strV != null) { |
|||
latestEntity.setStrValue(strV); |
|||
} else { |
|||
Long longV = null; |
|||
try { |
|||
longV = Long.parseLong(data[5].getValue()); |
|||
} catch (Exception e) { |
|||
} |
|||
if (longV != null) { |
|||
latestEntity.setLongValue(longV); |
|||
} else { |
|||
Double doubleV = null; |
|||
try { |
|||
doubleV = Double.parseDouble(data[6].getValue()); |
|||
} catch (Exception e) { |
|||
} |
|||
if (doubleV != null) { |
|||
latestEntity.setDoubleValue(doubleV); |
|||
} else { |
|||
|
|||
String jsonV = data[7].getValue(); |
|||
if (StringUtils.isNoneEmpty(jsonV)) { |
|||
latestEntity.setJsonValue(jsonV); |
|||
} else { |
|||
Boolean boolV = null; |
|||
try { |
|||
boolV = Boolean.parseBoolean(data[3].getValue()); |
|||
} catch (Exception e) { |
|||
} |
|||
if (boolV != null) { |
|||
latestEntity.setBooleanValue(boolV); |
|||
} else { |
|||
log.warn("All values in key-value row are nullable "); |
|||
} |
|||
} |
|||
} |
|||
} |
|||
} |
|||
return latestEntity; |
|||
} |
|||
|
|||
protected Integer getOrSaveKeyId(String strKey) { |
|||
Integer keyId = tsKvDictionaryMap.get(strKey); |
|||
if (keyId == null) { |
|||
Optional<TsKvDictionary> tsKvDictionaryOptional; |
|||
tsKvDictionaryOptional = dictionaryRepository.findById(new TsKvDictionaryCompositeKey(strKey)); |
|||
if (!tsKvDictionaryOptional.isPresent()) { |
|||
tsCreationLock.lock(); |
|||
try { |
|||
tsKvDictionaryOptional = dictionaryRepository.findById(new TsKvDictionaryCompositeKey(strKey)); |
|||
if (!tsKvDictionaryOptional.isPresent()) { |
|||
TsKvDictionary tsKvDictionary = new TsKvDictionary(); |
|||
tsKvDictionary.setKey(strKey); |
|||
try { |
|||
TsKvDictionary saved = dictionaryRepository.save(tsKvDictionary); |
|||
tsKvDictionaryMap.put(saved.getKey(), saved.getKeyId()); |
|||
keyId = saved.getKeyId(); |
|||
} catch (ConstraintViolationException e) { |
|||
tsKvDictionaryOptional = dictionaryRepository.findById(new TsKvDictionaryCompositeKey(strKey)); |
|||
TsKvDictionary dictionary = tsKvDictionaryOptional.orElseThrow(() -> new RuntimeException("Failed to get TsKvDictionary entity from DB!")); |
|||
tsKvDictionaryMap.put(dictionary.getKey(), dictionary.getKeyId()); |
|||
keyId = dictionary.getKeyId(); |
|||
} |
|||
} else { |
|||
keyId = tsKvDictionaryOptional.get().getKeyId(); |
|||
} |
|||
} finally { |
|||
tsCreationLock.unlock(); |
|||
} |
|||
} else { |
|||
keyId = tsKvDictionaryOptional.get().getKeyId(); |
|||
tsKvDictionaryMap.put(strKey, keyId); |
|||
} |
|||
} |
|||
return keyId; |
|||
} |
|||
} |
|||
@ -0,0 +1,21 @@ |
|||
/** |
|||
* Copyright © 2016-2020 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.migrate; |
|||
|
|||
public interface TsLatestMigrateService { |
|||
|
|||
void migrate() throws Exception; |
|||
} |
|||
@ -0,0 +1,26 @@ |
|||
/** |
|||
* Copyright © 2016-2020 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; |
|||
|
|||
import java.lang.annotation.Retention; |
|||
import java.lang.annotation.RetentionPolicy; |
|||
|
|||
@Retention(RetentionPolicy.RUNTIME) |
|||
@ConditionalOnProperty(prefix = "database.ts_latest", value = "type", havingValue = "cassandra") |
|||
public @interface NoSqlTsLatestDao { |
|||
} |
|||
@ -0,0 +1,26 @@ |
|||
/** |
|||
* Copyright © 2016-2020 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; |
|||
|
|||
import java.lang.annotation.Retention; |
|||
import java.lang.annotation.RetentionPolicy; |
|||
|
|||
@Retention(RetentionPolicy.RUNTIME) |
|||
@ConditionalOnProperty(prefix = "database.ts_latest", value = "type", havingValue = "timescale") |
|||
public @interface TimescaleDBTsLatestDao { |
|||
} |
|||
@ -0,0 +1,26 @@ |
|||
/** |
|||
* Copyright © 2016-2020 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.ConditionalOnExpression; |
|||
|
|||
import java.lang.annotation.Retention; |
|||
import java.lang.annotation.RetentionPolicy; |
|||
|
|||
@Retention(RetentionPolicy.RUNTIME) |
|||
@ConditionalOnExpression("'${database.ts.type}'=='timescale' || '${database.ts_latest.type}'=='timescale'") |
|||
public @interface TimescaleDBTsOrTsLatestDao { |
|||
} |
|||
@ -0,0 +1,37 @@ |
|||
/** |
|||
* Copyright © 2016-2020 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.HsqlDao; |
|||
import org.thingsboard.server.dao.util.SqlTsLatestDao; |
|||
|
|||
@Configuration |
|||
@EnableAutoConfiguration |
|||
@ComponentScan({"org.thingsboard.server.dao.sqlts.hsql"}) |
|||
@EnableJpaRepositories({"org.thingsboard.server.dao.sqlts.insert.latest.hsql", "org.thingsboard.server.dao.sqlts.latest"}) |
|||
@EntityScan({"org.thingsboard.server.dao.model.sqlts.latest"}) |
|||
@EnableTransactionManagement |
|||
@SqlTsLatestDao |
|||
@HsqlDao |
|||
public class HsqlTsLatestDaoConfig { |
|||
|
|||
} |
|||
@ -0,0 +1,37 @@ |
|||
/** |
|||
* Copyright © 2016-2020 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.PsqlDao; |
|||
import org.thingsboard.server.dao.util.SqlTsLatestDao; |
|||
|
|||
@Configuration |
|||
@EnableAutoConfiguration |
|||
@ComponentScan({"org.thingsboard.server.dao.sqlts.psql"}) |
|||
@EnableJpaRepositories({"org.thingsboard.server.dao.sqlts.insert.latest.psql", "org.thingsboard.server.dao.sqlts.latest"}) |
|||
@EntityScan({"org.thingsboard.server.dao.model.sqlts.latest"}) |
|||
@EnableTransactionManagement |
|||
@SqlTsLatestDao |
|||
@PsqlDao |
|||
public class PsqlTsLatestDaoConfig { |
|||
|
|||
} |
|||
@ -0,0 +1,34 @@ |
|||
/** |
|||
* Copyright © 2016-2020 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.Configuration; |
|||
import org.springframework.data.jpa.repository.config.EnableJpaRepositories; |
|||
import org.springframework.transaction.annotation.EnableTransactionManagement; |
|||
import org.thingsboard.server.dao.util.PsqlDao; |
|||
import org.thingsboard.server.dao.util.SqlTsOrTsLatestAnyDao; |
|||
|
|||
@Configuration |
|||
@EnableAutoConfiguration |
|||
@EnableJpaRepositories({"org.thingsboard.server.dao.sqlts.dictionary"}) |
|||
@EntityScan({"org.thingsboard.server.dao.model.sqlts.dictionary"}) |
|||
@EnableTransactionManagement |
|||
@SqlTsOrTsLatestAnyDao |
|||
public class SqlTimeseriesDaoConfig { |
|||
|
|||
} |
|||
@ -0,0 +1,37 @@ |
|||
/** |
|||
* Copyright © 2016-2020 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.PsqlDao; |
|||
import org.thingsboard.server.dao.util.TimescaleDBTsLatestDao; |
|||
|
|||
@Configuration |
|||
@EnableAutoConfiguration |
|||
@ComponentScan({"org.thingsboard.server.dao.sqlts.timescale"}) |
|||
@EnableJpaRepositories({"org.thingsboard.server.dao.sqlts.insert.latest.psql", "org.thingsboard.server.dao.sqlts.latest"}) |
|||
@EntityScan({"org.thingsboard.server.dao.model.sqlts.latest"}) |
|||
@EnableTransactionManagement |
|||
@TimescaleDBTsLatestDao |
|||
@PsqlDao |
|||
public class TimescaleTsLatestDaoConfig { |
|||
|
|||
} |
|||
@ -0,0 +1,29 @@ |
|||
/** |
|||
* Copyright © 2016-2020 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.util.concurrent.ListenableFuture; |
|||
import org.thingsboard.server.common.data.id.EntityId; |
|||
import org.thingsboard.server.common.data.id.TenantId; |
|||
import org.thingsboard.server.common.data.kv.ReadTsKvQuery; |
|||
import org.thingsboard.server.common.data.kv.TsKvEntry; |
|||
|
|||
import java.util.List; |
|||
|
|||
public interface AggregationTimeseriesDao { |
|||
|
|||
ListenableFuture<List<TsKvEntry>> findAllAsync(TenantId tenantId, EntityId entityId, ReadTsKvQuery query); |
|||
} |
|||
@ -0,0 +1,99 @@ |
|||
/** |
|||
* Copyright © 2016-2020 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 lombok.extern.slf4j.Slf4j; |
|||
import org.hibernate.exception.ConstraintViolationException; |
|||
import org.springframework.beans.factory.annotation.Autowired; |
|||
import org.thingsboard.server.common.data.kv.TsKvEntry; |
|||
import org.thingsboard.server.dao.model.sqlts.dictionary.TsKvDictionary; |
|||
import org.thingsboard.server.dao.model.sqlts.dictionary.TsKvDictionaryCompositeKey; |
|||
import org.thingsboard.server.dao.sql.JpaAbstractDaoListeningExecutorService; |
|||
import org.thingsboard.server.dao.sqlts.dictionary.TsKvDictionaryRepository; |
|||
|
|||
import javax.annotation.Nullable; |
|||
import java.util.List; |
|||
import java.util.Optional; |
|||
import java.util.concurrent.ConcurrentHashMap; |
|||
import java.util.concurrent.ConcurrentMap; |
|||
import java.util.concurrent.locks.ReentrantLock; |
|||
import java.util.stream.Collectors; |
|||
|
|||
@Slf4j |
|||
public abstract class BaseAbstractSqlTimeseriesDao extends JpaAbstractDaoListeningExecutorService { |
|||
|
|||
private final ConcurrentMap<String, Integer> tsKvDictionaryMap = new ConcurrentHashMap<>(); |
|||
|
|||
protected static final ReentrantLock tsCreationLock = new ReentrantLock(); |
|||
|
|||
@Autowired |
|||
protected TsKvDictionaryRepository dictionaryRepository; |
|||
|
|||
protected Integer getOrSaveKeyId(String strKey) { |
|||
Integer keyId = tsKvDictionaryMap.get(strKey); |
|||
if (keyId == null) { |
|||
Optional<TsKvDictionary> tsKvDictionaryOptional; |
|||
tsKvDictionaryOptional = dictionaryRepository.findById(new TsKvDictionaryCompositeKey(strKey)); |
|||
if (!tsKvDictionaryOptional.isPresent()) { |
|||
tsCreationLock.lock(); |
|||
try { |
|||
tsKvDictionaryOptional = dictionaryRepository.findById(new TsKvDictionaryCompositeKey(strKey)); |
|||
if (!tsKvDictionaryOptional.isPresent()) { |
|||
TsKvDictionary tsKvDictionary = new TsKvDictionary(); |
|||
tsKvDictionary.setKey(strKey); |
|||
try { |
|||
TsKvDictionary saved = dictionaryRepository.save(tsKvDictionary); |
|||
tsKvDictionaryMap.put(saved.getKey(), saved.getKeyId()); |
|||
keyId = saved.getKeyId(); |
|||
} catch (ConstraintViolationException e) { |
|||
tsKvDictionaryOptional = dictionaryRepository.findById(new TsKvDictionaryCompositeKey(strKey)); |
|||
TsKvDictionary dictionary = tsKvDictionaryOptional.orElseThrow(() -> new RuntimeException("Failed to get TsKvDictionary entity from DB!")); |
|||
tsKvDictionaryMap.put(dictionary.getKey(), dictionary.getKeyId()); |
|||
keyId = dictionary.getKeyId(); |
|||
} |
|||
} else { |
|||
keyId = tsKvDictionaryOptional.get().getKeyId(); |
|||
} |
|||
} finally { |
|||
tsCreationLock.unlock(); |
|||
} |
|||
} else { |
|||
keyId = tsKvDictionaryOptional.get().getKeyId(); |
|||
tsKvDictionaryMap.put(strKey, keyId); |
|||
} |
|||
} |
|||
return keyId; |
|||
} |
|||
|
|||
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); |
|||
} |
|||
} |
|||
@ -0,0 +1,261 @@ |
|||
/** |
|||
* Copyright © 2016-2020 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.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.MoreExecutors; |
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.springframework.beans.factory.annotation.Autowired; |
|||
import org.springframework.beans.factory.annotation.Value; |
|||
import org.springframework.stereotype.Component; |
|||
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.common.stats.StatsFactory; |
|||
import org.thingsboard.server.dao.DaoUtil; |
|||
import org.thingsboard.server.dao.model.sqlts.latest.TsKvLatestCompositeKey; |
|||
import org.thingsboard.server.dao.model.sqlts.latest.TsKvLatestEntity; |
|||
import org.thingsboard.server.dao.sql.ScheduledLogExecutorComponent; |
|||
import org.thingsboard.server.dao.sql.TbSqlBlockingQueueParams; |
|||
import org.thingsboard.server.dao.sql.TbSqlBlockingQueueWrapper; |
|||
import org.thingsboard.server.dao.sqlts.insert.latest.InsertLatestTsRepository; |
|||
import org.thingsboard.server.dao.sqlts.latest.SearchTsKvLatestRepository; |
|||
import org.thingsboard.server.dao.sqlts.latest.TsKvLatestRepository; |
|||
import org.thingsboard.server.dao.timeseries.SimpleListenableFuture; |
|||
import org.thingsboard.server.dao.timeseries.TimeseriesLatestDao; |
|||
import org.thingsboard.server.dao.util.SqlTsLatestAnyDao; |
|||
|
|||
import javax.annotation.Nullable; |
|||
import javax.annotation.PostConstruct; |
|||
import javax.annotation.PreDestroy; |
|||
import java.util.ArrayList; |
|||
import java.util.HashMap; |
|||
import java.util.List; |
|||
import java.util.Map; |
|||
import java.util.Optional; |
|||
import java.util.concurrent.ExecutionException; |
|||
|
|||
@Slf4j |
|||
@Component |
|||
@SqlTsLatestAnyDao |
|||
public class SqlTimeseriesLatestDao extends BaseAbstractSqlTimeseriesDao implements TimeseriesLatestDao { |
|||
|
|||
private static final String DESC_ORDER = "DESC"; |
|||
|
|||
@Autowired |
|||
private TsKvLatestRepository tsKvLatestRepository; |
|||
|
|||
@Autowired |
|||
protected AggregationTimeseriesDao aggregationTimeseriesDao; |
|||
|
|||
@Autowired |
|||
private SearchTsKvLatestRepository searchTsKvLatestRepository; |
|||
|
|||
@Autowired |
|||
private InsertLatestTsRepository insertLatestTsRepository; |
|||
|
|||
private TbSqlBlockingQueueWrapper<TsKvLatestEntity> tsLatestQueue; |
|||
|
|||
@Value("${sql.ts_latest.batch_size:1000}") |
|||
private int tsLatestBatchSize; |
|||
|
|||
@Value("${sql.ts_latest.batch_max_delay:100}") |
|||
private long tsLatestMaxDelay; |
|||
|
|||
@Value("${sql.ts_latest.stats_print_interval_ms:1000}") |
|||
private long tsLatestStatsPrintIntervalMs; |
|||
|
|||
@Value("${sql.ts_latest.batch_threads:4}") |
|||
private int tsLatestBatchThreads; |
|||
|
|||
@Autowired |
|||
protected ScheduledLogExecutorComponent logExecutor; |
|||
|
|||
@Autowired |
|||
private StatsFactory statsFactory; |
|||
|
|||
@PostConstruct |
|||
protected void init() { |
|||
TbSqlBlockingQueueParams tsLatestParams = TbSqlBlockingQueueParams.builder() |
|||
.logName("TS Latest") |
|||
.batchSize(tsLatestBatchSize) |
|||
.maxDelay(tsLatestMaxDelay) |
|||
.statsPrintIntervalMs(tsLatestStatsPrintIntervalMs) |
|||
.statsNamePrefix("ts.latest") |
|||
.build(); |
|||
|
|||
java.util.function.Function<TsKvLatestEntity, Integer> hashcodeFunction = entity -> entity.getEntityId().hashCode(); |
|||
tsLatestQueue = new TbSqlBlockingQueueWrapper<>(tsLatestParams, hashcodeFunction, tsLatestBatchThreads, statsFactory); |
|||
|
|||
tsLatestQueue.init(logExecutor, v -> { |
|||
Map<TsKey, TsKvLatestEntity> trueLatest = new HashMap<>(); |
|||
v.forEach(ts -> { |
|||
TsKey key = new TsKey(ts.getEntityId(), ts.getKey()); |
|||
TsKvLatestEntity old = trueLatest.get(key); |
|||
if (old == null || old.getTs() < ts.getTs()) { |
|||
trueLatest.put(key, ts); |
|||
} |
|||
}); |
|||
List<TsKvLatestEntity> latestEntities = new ArrayList<>(trueLatest.values()); |
|||
insertLatestTsRepository.saveOrUpdate(latestEntities); |
|||
}); |
|||
} |
|||
|
|||
@PreDestroy |
|||
protected void destroy() { |
|||
if (tsLatestQueue != null) { |
|||
tsLatestQueue.destroy(); |
|||
} |
|||
} |
|||
|
|||
@Override |
|||
public ListenableFuture<Void> saveLatest(TenantId tenantId, EntityId entityId, TsKvEntry tsKvEntry) { |
|||
return getSaveLatestFuture(entityId, tsKvEntry); |
|||
} |
|||
|
|||
@Override |
|||
public ListenableFuture<Void> removeLatest(TenantId tenantId, EntityId entityId, DeleteTsKvQuery query) { |
|||
return getRemoveLatestFuture(tenantId, entityId, query); |
|||
} |
|||
|
|||
@Override |
|||
public ListenableFuture<TsKvEntry> findLatest(TenantId tenantId, EntityId entityId, String key) { |
|||
return getFindLatestFuture(entityId, key); |
|||
} |
|||
|
|||
@Override |
|||
public ListenableFuture<List<TsKvEntry>> findAllLatest(TenantId tenantId, EntityId entityId) { |
|||
return getFindAllLatestFuture(entityId); |
|||
} |
|||
|
|||
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 getSaveLatestFuture(entityId, entryList.get(0)); |
|||
} else { |
|||
log.trace("Could not find new latest value for [{}], key - {}", entityId, query.getKey()); |
|||
} |
|||
return Futures.immediateFuture(null); |
|||
}, service); |
|||
} |
|||
|
|||
private 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 aggregationTimeseriesDao.findAllAsync(tenantId, entityId, findNewLatestQuery); |
|||
} |
|||
|
|||
protected ListenableFuture<TsKvEntry> getFindLatestFuture(EntityId entityId, String key) { |
|||
TsKvLatestCompositeKey compositeKey = |
|||
new TsKvLatestCompositeKey( |
|||
entityId.getId(), |
|||
getOrSaveKeyId(key)); |
|||
Optional<TsKvLatestEntity> entry = tsKvLatestRepository.findById(compositeKey); |
|||
TsKvEntry result; |
|||
if (entry.isPresent()) { |
|||
TsKvLatestEntity tsKvLatestEntity = entry.get(); |
|||
tsKvLatestEntity.setStrKey(key); |
|||
result = DaoUtil.getData(tsKvLatestEntity); |
|||
} else { |
|||
result = new BasicTsKvEntry(System.currentTimeMillis(), new StringDataEntry(key, null)); |
|||
} |
|||
return Futures.immediateFuture(result); |
|||
} |
|||
|
|||
protected ListenableFuture<Void> getRemoveLatestFuture(TenantId tenantId, EntityId entityId, DeleteTsKvQuery query) { |
|||
ListenableFuture<TsKvEntry> latestFuture = getFindLatestFuture(entityId, query.getKey()); |
|||
|
|||
ListenableFuture<Boolean> booleanFuture = Futures.transform(latestFuture, tsKvEntry -> { |
|||
long ts = tsKvEntry.getTs(); |
|||
return ts > query.getStartTs() && ts <= query.getEndTs(); |
|||
}, service); |
|||
|
|||
ListenableFuture<Void> removedLatestFuture = Futures.transformAsync(booleanFuture, isRemove -> { |
|||
if (isRemove) { |
|||
TsKvLatestEntity latestEntity = new TsKvLatestEntity(); |
|||
latestEntity.setEntityId(entityId.getId()); |
|||
latestEntity.setKey(getOrSaveKeyId(query.getKey())); |
|||
return service.submit(() -> { |
|||
tsKvLatestRepository.delete(latestEntity); |
|||
return null; |
|||
}); |
|||
} |
|||
return Futures.immediateFuture(null); |
|||
}, service); |
|||
|
|||
final SimpleListenableFuture<Void> resultFuture = new SimpleListenableFuture<>(); |
|||
Futures.addCallback(removedLatestFuture, new FutureCallback<Void>() { |
|||
@Override |
|||
public void onSuccess(@Nullable Void result) { |
|||
if (query.getRewriteLatestIfDeleted()) { |
|||
ListenableFuture<Void> savedLatestFuture = Futures.transformAsync(booleanFuture, isRemove -> { |
|||
if (isRemove) { |
|||
return getNewLatestEntryFuture(tenantId, entityId, query); |
|||
} |
|||
return Futures.immediateFuture(null); |
|||
}, service); |
|||
|
|||
try { |
|||
resultFuture.set(savedLatestFuture.get()); |
|||
} catch (InterruptedException | ExecutionException e) { |
|||
log.warn("Could not get latest saved value for [{}], {}", entityId, query.getKey(), e); |
|||
} |
|||
} else { |
|||
resultFuture.set(null); |
|||
} |
|||
} |
|||
|
|||
@Override |
|||
public void onFailure(Throwable t) { |
|||
log.warn("[{}] Failed to process remove of the latest value", entityId, t); |
|||
} |
|||
}, MoreExecutors.directExecutor()); |
|||
return resultFuture; |
|||
} |
|||
|
|||
protected ListenableFuture<List<TsKvEntry>> getFindAllLatestFuture(EntityId entityId) { |
|||
return Futures.immediateFuture( |
|||
DaoUtil.convertDataList(Lists.newArrayList( |
|||
searchTsKvLatestRepository.findAllByEntityId(entityId.getId())))); |
|||
} |
|||
|
|||
protected ListenableFuture<Void> getSaveLatestFuture(EntityId entityId, TsKvEntry tsKvEntry) { |
|||
TsKvLatestEntity latestEntity = new TsKvLatestEntity(); |
|||
latestEntity.setEntityId(entityId.getId()); |
|||
latestEntity.setTs(tsKvEntry.getTs()); |
|||
latestEntity.setKey(getOrSaveKeyId(tsKvEntry.getKey())); |
|||
latestEntity.setStrValue(tsKvEntry.getStrValue().orElse(null)); |
|||
latestEntity.setDoubleValue(tsKvEntry.getDoubleValue().orElse(null)); |
|||
latestEntity.setLongValue(tsKvEntry.getLongValue().orElse(null)); |
|||
latestEntity.setBooleanValue(tsKvEntry.getBooleanValue().orElse(null)); |
|||
latestEntity.setJsonValue(tsKvEntry.getJsonValue().orElse(null)); |
|||
|
|||
return tsLatestQueue.add(latestEntity); |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,88 @@ |
|||
/** |
|||
* Copyright © 2016-2020 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.timeseries; |
|||
|
|||
import com.datastax.oss.driver.api.core.cql.Row; |
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.apache.commons.lang3.StringUtils; |
|||
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.JsonDataEntry; |
|||
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.ModelConstants; |
|||
import org.thingsboard.server.dao.nosql.CassandraAbstractAsyncDao; |
|||
|
|||
import java.util.ArrayList; |
|||
import java.util.List; |
|||
|
|||
@Slf4j |
|||
public abstract class AbstractCassandraBaseTimeseriesDao extends CassandraAbstractAsyncDao { |
|||
public static final String DESC_ORDER = "DESC"; |
|||
public static final String GENERATED_QUERY_FOR_ENTITY_TYPE_AND_ENTITY_ID = "Generated query [{}] for entityType {} and entityId {}"; |
|||
public static final String INSERT_INTO = "INSERT INTO "; |
|||
public static final String SELECT_PREFIX = "SELECT "; |
|||
public static final String EQUALS_PARAM = " = ? "; |
|||
|
|||
public static KvEntry toKvEntry(Row row, String key) { |
|||
KvEntry kvEntry = null; |
|||
String strV = row.get(ModelConstants.STRING_VALUE_COLUMN, String.class); |
|||
if (strV != null) { |
|||
kvEntry = new StringDataEntry(key, strV); |
|||
} else { |
|||
Long longV = row.get(ModelConstants.LONG_VALUE_COLUMN, Long.class); |
|||
if (longV != null) { |
|||
kvEntry = new LongDataEntry(key, longV); |
|||
} else { |
|||
Double doubleV = row.get(ModelConstants.DOUBLE_VALUE_COLUMN, Double.class); |
|||
if (doubleV != null) { |
|||
kvEntry = new DoubleDataEntry(key, doubleV); |
|||
} else { |
|||
Boolean boolV = row.get(ModelConstants.BOOLEAN_VALUE_COLUMN, Boolean.class); |
|||
if (boolV != null) { |
|||
kvEntry = new BooleanDataEntry(key, boolV); |
|||
} else { |
|||
String jsonV = row.get(ModelConstants.JSON_VALUE_COLUMN, String.class); |
|||
if (StringUtils.isNoneEmpty(jsonV)) { |
|||
kvEntry = new JsonDataEntry(key, jsonV); |
|||
} else { |
|||
log.warn("All values in key-value row are nullable "); |
|||
} |
|||
} |
|||
} |
|||
} |
|||
} |
|||
return kvEntry; |
|||
} |
|||
|
|||
protected List<TsKvEntry> convertResultToTsKvEntryList(List<Row> rows) { |
|||
List<TsKvEntry> entries = new ArrayList<>(rows.size()); |
|||
if (!rows.isEmpty()) { |
|||
rows.forEach(row -> entries.add(convertResultToTsKvEntry(row))); |
|||
} |
|||
return entries; |
|||
} |
|||
|
|||
private TsKvEntry convertResultToTsKvEntry(Row row) { |
|||
String key = row.getString(ModelConstants.KEY_COLUMN); |
|||
long ts = row.getLong(ModelConstants.TS_COLUMN); |
|||
return new BasicTsKvEntry(ts, toKvEntry(row, key)); |
|||
} |
|||
|
|||
} |
|||
Some files were not shown because too many files changed in this diff
Loading…
Reference in new issue