Browse Source

Merge remote-tracking branch 'origin/lts-4.2' into lts-4.3

pull/15069/head
Viacheslav Klimov 7 months ago
parent
commit
9a2a7a4468
Failed to extract signature
  1. 2
      application/src/main/resources/thingsboard.yml
  2. 32
      common/dao-api/src/main/java/org/thingsboard/server/dao/nosql/ResultSetSizeLimitExceededException.java
  3. 24
      common/dao-api/src/main/java/org/thingsboard/server/dao/nosql/TbResultSet.java
  4. 140
      common/dao-api/src/test/java/org/thingsboard/server/dao/nosql/TbResultSetTest.java
  5. 16
      common/version-control/src/main/java/org/thingsboard/server/service/sync/DefaultGitSyncService.java
  6. 35
      common/version-control/src/main/java/org/thingsboard/server/service/sync/vc/GitRepository.java
  7. 3
      dao/src/main/java/org/thingsboard/server/dao/nosql/CassandraAbstractAsyncDao.java
  8. 92
      dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java
  9. 2
      dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesLatestDao.java
  10. 4
      dao/src/main/java/org/thingsboard/server/dao/timeseries/SimpleListenableFuture.java
  11. 42
      dao/src/test/java/org/thingsboard/server/dao/service/timeseries/nosql/TimeseriesServiceNoSqlTest.java
  12. 34
      rest-client/src/main/resources/logback.xml

2
application/src/main/resources/thingsboard.yml

@ -339,6 +339,8 @@ cassandra:
set_null_values_enabled: "${CASSANDRA_QUERY_SET_NULL_VALUES_ENABLED:true}" set_null_values_enabled: "${CASSANDRA_QUERY_SET_NULL_VALUES_ENABLED:true}"
# log one of cassandra queries with specified frequency (0 - logging is disabled) # log one of cassandra queries with specified frequency (0 - logging is disabled)
print_queries_freq: "${CASSANDRA_QUERY_PRINT_FREQ:0}" print_queries_freq: "${CASSANDRA_QUERY_PRINT_FREQ:0}"
# Maximum total size in bytes of a Cassandra query result set across all pages. Default is 50MB. 0 means unlimited
max_result_set_size_in_bytes: "${CASSANDRA_QUERY_MAX_RESULT_SET_SIZE_IN_BYTES:52428800}"
tenant_rate_limits: tenant_rate_limits:
# Whether to print rate-limited tenant names when printing Cassandra query queue statistic # Whether to print rate-limited tenant names when printing Cassandra query queue statistic
print_tenant_names: "${CASSANDRA_QUERY_TENANT_RATE_LIMITS_PRINT_TENANT_NAMES:false}" print_tenant_names: "${CASSANDRA_QUERY_TENANT_RATE_LIMITS_PRINT_TENANT_NAMES:false}"

32
common/dao-api/src/main/java/org/thingsboard/server/dao/nosql/ResultSetSizeLimitExceededException.java

@ -0,0 +1,32 @@
/**
* Copyright © 2016-2026 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.nosql;
import lombok.Getter;
@Getter
public class ResultSetSizeLimitExceededException extends IllegalArgumentException {
private final long limitBytes;
private final long actualBytes;
public ResultSetSizeLimitExceededException(long limitBytes, long actualBytes) {
super("Result set size exceeds the maximum allowed limit. Please narrow your query");
this.limitBytes = limitBytes;
this.actualBytes = actualBytes;
}
}

24
common/dao-api/src/main/java/org/thingsboard/server/dao/nosql/TbResultSet.java

@ -34,6 +34,7 @@ import java.util.ArrayList;
import java.util.List; import java.util.List;
import java.util.concurrent.CompletionStage; import java.util.concurrent.CompletionStage;
import java.util.concurrent.Executor; import java.util.concurrent.Executor;
import java.util.concurrent.atomic.AtomicLong;
import java.util.function.Function; import java.util.function.Function;
public class TbResultSet implements AsyncResultSet { public class TbResultSet implements AsyncResultSet {
@ -89,9 +90,14 @@ public class TbResultSet implements AsyncResultSet {
} }
public ListenableFuture<List<Row>> allRows(Executor executor) { public ListenableFuture<List<Row>> allRows(Executor executor) {
return allRows(executor, 0);
}
public ListenableFuture<List<Row>> allRows(Executor executor, long maxResultSetSizeBytes) {
List<Row> allRows = new ArrayList<>(); List<Row> allRows = new ArrayList<>();
SettableFuture<List<Row>> resultFuture = SettableFuture.create(); SettableFuture<List<Row>> resultFuture = SettableFuture.create();
this.processRows(originalStatement, delegate, allRows, resultFuture, executor); AtomicLong accumulatedBytes = new AtomicLong(0);
this.processRows(originalStatement, delegate, allRows, resultFuture, executor, maxResultSetSizeBytes, accumulatedBytes);
return resultFuture; return resultFuture;
} }
@ -99,7 +105,19 @@ public class TbResultSet implements AsyncResultSet {
AsyncResultSet resultSet, AsyncResultSet resultSet,
List<Row> allRows, List<Row> allRows,
SettableFuture<List<Row>> resultFuture, SettableFuture<List<Row>> resultFuture,
Executor executor) { Executor executor,
long maxResultSetSizeBytes,
AtomicLong accumulatedBytes) {
if (maxResultSetSizeBytes > 0) {
int pageSizeInBytes = resultSet.getExecutionInfo().getResponseSizeInBytes();
if (pageSizeInBytes > 0) {
accumulatedBytes.addAndGet(pageSizeInBytes);
}
if (accumulatedBytes.get() > maxResultSetSizeBytes) {
resultFuture.setException(new ResultSetSizeLimitExceededException(maxResultSetSizeBytes, accumulatedBytes.get()));
return;
}
}
allRows.addAll(loadRows(resultSet)); allRows.addAll(loadRows(resultSet));
if (resultSet.hasMorePages()) { if (resultSet.hasMorePages()) {
ByteBuffer nextPagingState = resultSet.getExecutionInfo().getPagingState(); ByteBuffer nextPagingState = resultSet.getExecutionInfo().getPagingState();
@ -110,7 +128,7 @@ public class TbResultSet implements AsyncResultSet {
@Override @Override
public void onSuccess(@Nullable TbResultSet result) { public void onSuccess(@Nullable TbResultSet result) {
processRows(nextStatement, result, processRows(nextStatement, result,
allRows, resultFuture, executor); allRows, resultFuture, executor, maxResultSetSizeBytes, accumulatedBytes);
} }
@Override @Override

140
common/dao-api/src/test/java/org/thingsboard/server/dao/nosql/TbResultSetTest.java

@ -0,0 +1,140 @@
/**
* Copyright © 2016-2026 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.nosql;
import com.datastax.oss.driver.api.core.cql.AsyncResultSet;
import com.datastax.oss.driver.api.core.cql.ColumnDefinitions;
import com.datastax.oss.driver.api.core.cql.ExecutionInfo;
import com.datastax.oss.driver.api.core.cql.Row;
import com.datastax.oss.driver.api.core.cql.Statement;
import com.google.common.util.concurrent.ListenableFuture;
import com.google.common.util.concurrent.MoreExecutors;
import com.google.common.util.concurrent.SettableFuture;
import org.junit.jupiter.api.Test;
import java.nio.ByteBuffer;
import java.util.List;
import java.util.concurrent.ExecutionException;
import java.util.function.Function;
import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.assertThatThrownBy;
import static org.mockito.Mockito.doReturn;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.when;
class TbResultSetTest {
@Test
void allRows_withinLimit_returnsAllRows() throws Exception {
Row row = mock(Row.class);
AsyncResultSet asyncResultSet = createMockResultSet(List.of(row), false, 1000);
Statement<?> statement = mock(Statement.class);
TbResultSet tbResultSet = new TbResultSet(statement, asyncResultSet, s -> null);
ListenableFuture<List<Row>> future = tbResultSet.allRows(MoreExecutors.directExecutor(), 5000);
List<Row> result = future.get();
assertThat(result).hasSize(1);
assertThat(result.get(0)).isSameAs(row);
}
@Test
void allRows_exceedsLimitOnFirstPage_failsWithException() {
Row row = mock(Row.class);
AsyncResultSet asyncResultSet = createMockResultSet(List.of(row), false, 6000);
Statement<?> statement = mock(Statement.class);
TbResultSet tbResultSet = new TbResultSet(statement, asyncResultSet, s -> null);
ListenableFuture<List<Row>> future = tbResultSet.allRows(MoreExecutors.directExecutor(), 5000);
assertThatThrownBy(future::get)
.isInstanceOf(ExecutionException.class)
.hasCauseInstanceOf(ResultSetSizeLimitExceededException.class);
}
@Test
void allRows_exceedsLimitOnSecondPage_failsAfterSecondPage() {
Row row1 = mock(Row.class);
Row row2 = mock(Row.class);
Statement<?> statement = mock(Statement.class);
doReturn(statement).when(statement).setPagingState((ByteBuffer) null);
AsyncResultSet page2 = createMockResultSet(List.of(row2), false, 3000);
TbResultSet tbResultSetPage2 = new TbResultSet(statement, page2, s -> null);
SettableFuture<TbResultSet> page2Future = SettableFuture.create();
page2Future.set(tbResultSetPage2);
TbResultSetFuture tbPage2Future = new TbResultSetFuture(page2Future);
ExecutionInfo page1ExecInfo = mock(ExecutionInfo.class);
when(page1ExecInfo.getResponseSizeInBytes()).thenReturn(3000);
when(page1ExecInfo.getPagingState()).thenReturn(null);
AsyncResultSet page1 = createMockResultSet(List.of(row1), true, 3000);
when(page1.getExecutionInfo()).thenReturn(page1ExecInfo);
Function<Statement, TbResultSetFuture> executeAsync = s -> tbPage2Future;
TbResultSet tbResultSet = new TbResultSet(statement, page1, executeAsync);
ListenableFuture<List<Row>> future = tbResultSet.allRows(MoreExecutors.directExecutor(), 5000);
assertThatThrownBy(future::get)
.isInstanceOf(ExecutionException.class)
.hasCauseInstanceOf(ResultSetSizeLimitExceededException.class);
}
@Test
void allRows_unlimitedWithZero_returnsAllRowsRegardlessOfSize() throws Exception {
Row row = mock(Row.class);
AsyncResultSet asyncResultSet = createMockResultSet(List.of(row), false, 999999);
Statement<?> statement = mock(Statement.class);
TbResultSet tbResultSet = new TbResultSet(statement, asyncResultSet, s -> null);
ListenableFuture<List<Row>> future = tbResultSet.allRows(MoreExecutors.directExecutor(), 0);
List<Row> result = future.get();
assertThat(result).hasSize(1);
}
@Test
void allRows_noLimitOverload_returnsAllRows() throws Exception {
Row row = mock(Row.class);
AsyncResultSet asyncResultSet = createMockResultSet(List.of(row), false, 999999);
Statement<?> statement = mock(Statement.class);
TbResultSet tbResultSet = new TbResultSet(statement, asyncResultSet, s -> null);
ListenableFuture<List<Row>> future = tbResultSet.allRows(MoreExecutors.directExecutor());
List<Row> result = future.get();
assertThat(result).hasSize(1);
}
private AsyncResultSet createMockResultSet(List<Row> rows, boolean hasMorePages, int responseSizeInBytes) {
AsyncResultSet resultSet = mock(AsyncResultSet.class);
ExecutionInfo executionInfo = mock(ExecutionInfo.class);
ColumnDefinitions columnDefs = mock(ColumnDefinitions.class);
when(executionInfo.getResponseSizeInBytes()).thenReturn(responseSizeInBytes);
when(executionInfo.getPagingState()).thenReturn(null);
when(resultSet.getExecutionInfo()).thenReturn(executionInfo);
when(resultSet.getColumnDefinitions()).thenReturn(columnDefs);
when(resultSet.currentPage()).thenReturn(rows);
when(resultSet.hasMorePages()).thenReturn(hasMorePages);
when(resultSet.remaining()).thenReturn(rows.size());
return resultSet;
}
}

16
common/version-control/src/main/java/org/thingsboard/server/service/sync/DefaultGitSyncService.java

@ -18,6 +18,7 @@ package org.thingsboard.server.service.sync;
import jakarta.annotation.PreDestroy; import jakarta.annotation.PreDestroy;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.apache.commons.lang3.StringUtils; import org.apache.commons.lang3.StringUtils;
import org.eclipse.jgit.revwalk.RevCommit;
import org.springframework.beans.factory.annotation.Value; import org.springframework.beans.factory.annotation.Value;
import org.springframework.stereotype.Service; import org.springframework.stereotype.Service;
import org.thingsboard.common.util.ThingsBoardExecutors; import org.thingsboard.common.util.ThingsBoardExecutors;
@ -47,6 +48,8 @@ public class DefaultGitSyncService implements GitSyncService {
private final Map<String, GitRepository> repositories = new ConcurrentHashMap<>(); private final Map<String, GitRepository> repositories = new ConcurrentHashMap<>();
private final Map<String, Runnable> updateListeners = new ConcurrentHashMap<>(); private final Map<String, Runnable> updateListeners = new ConcurrentHashMap<>();
private RevCommit lastCommit;
@Override @Override
public void registerSync(String key, String repoUri, String branch, long fetchFrequencyMs, Runnable onUpdate) { public void registerSync(String key, String repoUri, String branch, long fetchFrequencyMs, Runnable onUpdate) {
RepositorySettings settings = new RepositorySettings(); RepositorySettings settings = new RepositorySettings();
@ -84,7 +87,7 @@ public class DefaultGitSyncService implements GitSyncService {
@Override @Override
public List<RepoFile> listFiles(String key, String path, int depth, FileType type) { public List<RepoFile> listFiles(String key, String path, int depth, FileType type) {
GitRepository repository = getRepository(key); GitRepository repository = getRepository(key);
return repository.listFilesAtCommit(getBranchRef(repository), path, depth).stream() return repository.listFilesAtCommit(lastCommit, path, depth).stream()
.filter(file -> type == null || file.type() == type) .filter(file -> type == null || file.type() == type)
.toList(); .toList();
} }
@ -93,7 +96,7 @@ public class DefaultGitSyncService implements GitSyncService {
@Override @Override
public byte[] getFileContent(String key, String path) { public byte[] getFileContent(String key, String path) {
GitRepository repository = getRepository(key); GitRepository repository = getRepository(key);
return repository.getFileContentAtCommit(path, getBranchRef(repository)); return repository.getFileContentAtCommit(path, lastCommit);
} }
@Override @Override
@ -137,6 +140,15 @@ public class DefaultGitSyncService implements GitSyncService {
} }
private void onUpdate(String key) { private void onUpdate(String key) {
GitRepository repository = getRepository(key);
String branchRef = getBranchRef(repository);
try {
lastCommit = repository.resolveCommit(branchRef);
} catch (Throwable e) {
log.error("[{}] Failed to resolve commit for ref {}", key, branchRef, e);
return;
}
Runnable listener = updateListeners.get(key); Runnable listener = updateListeners.get(key);
if (listener != null) { if (listener != null) {
log.debug("[{}] Handling repository update", key); log.debug("[{}] Handling repository update", key);

35
common/version-control/src/main/java/org/thingsboard/server/service/sync/vc/GitRepository.java

@ -277,13 +277,17 @@ public class GitRepository {
return listFilesAtCommit(commitId, path, -1).stream().map(RepoFile::path).toList(); return listFilesAtCommit(commitId, path, -1).stream().map(RepoFile::path).toList();
} }
@SneakyThrows
public List<RepoFile> listFilesAtCommit(String commitId, String path, int depth) { public List<RepoFile> listFilesAtCommit(String commitId, String path, int depth) {
log.debug("Executing listFilesAtCommit [{}][{}][{}]", settings.getRepositoryUri(), commitId, path); RevCommit commit = resolveCommit(commitId);
return listFilesAtCommit(commit, path, depth);
}
@SneakyThrows
public List<RepoFile> listFilesAtCommit(RevCommit commit, String path, int depth) {
log.debug("Executing listFilesAtCommit [{}][{}][{}]", settings.getRepositoryUri(), commit, path);
List<RepoFile> files = new ArrayList<>(); List<RepoFile> files = new ArrayList<>();
RevCommit revCommit = resolveCommit(commitId);
try (TreeWalk treeWalk = new TreeWalk(git.getRepository())) { try (TreeWalk treeWalk = new TreeWalk(git.getRepository())) {
treeWalk.reset(revCommit.getTree().getId()); treeWalk.reset(commit.getTree().getId());
if (StringUtils.isNotEmpty(path)) { if (StringUtils.isNotEmpty(path)) {
treeWalk.setFilter(PathFilter.create(path)); treeWalk.setFilter(PathFilter.create(path));
} }
@ -301,11 +305,15 @@ public class GitRepository {
return files; return files;
} }
@SneakyThrows
public byte[] getFileContentAtCommit(String file, String commitId) { public byte[] getFileContentAtCommit(String file, String commitId) {
log.debug("Executing getFileContentAtCommit [{}][{}][{}]", settings.getRepositoryUri(), commitId, file); RevCommit commit = resolveCommit(commitId);
RevCommit revCommit = resolveCommit(commitId); return getFileContentAtCommit(file, commit);
try (TreeWalk treeWalk = TreeWalk.forPath(git.getRepository(), file, revCommit.getTree())) { }
@SneakyThrows
public byte[] getFileContentAtCommit(String file, RevCommit commit) {
log.debug("Executing getFileContentAtCommit [{}][{}][{}]", settings.getRepositoryUri(), commit, file);
try (TreeWalk treeWalk = TreeWalk.forPath(git.getRepository(), file, commit.getTree())) {
if (treeWalk == null) { if (treeWalk == null) {
throw new IllegalArgumentException("File not found"); throw new IllegalArgumentException("File not found");
} }
@ -373,9 +381,9 @@ public class GitRepository {
for (RemoteRefUpdate update : pushResult.getRemoteUpdates()) { for (RemoteRefUpdate update : pushResult.getRemoteUpdates()) {
RemoteRefUpdate.Status status = update.getStatus(); RemoteRefUpdate.Status status = update.getStatus();
if (status == REJECTED_NONFASTFORWARD || status == REJECTED_NODELETE || if (status == REJECTED_NONFASTFORWARD || status == REJECTED_NODELETE ||
status == REJECTED_REMOTE_CHANGED || status == REJECTED_OTHER_REASON) { status == REJECTED_REMOTE_CHANGED || status == REJECTED_OTHER_REASON) {
throw new RuntimeException("Remote repository answered with error: " + throw new RuntimeException("Remote repository answered with error: " +
Optional.ofNullable(update.getMessage()).orElseGet(status::name)); Optional.ofNullable(update.getMessage()).orElseGet(status::name));
} }
} }
}); });
@ -450,7 +458,8 @@ public class GitRepository {
revCommit.getFullMessage(), revCommit.getAuthorIdent().getName(), revCommit.getAuthorIdent().getEmailAddress()); revCommit.getFullMessage(), revCommit.getAuthorIdent().getName(), revCommit.getAuthorIdent().getEmailAddress());
} }
private RevCommit resolveCommit(String id) throws IOException { @SneakyThrows
public RevCommit resolveCommit(String id) {
return git.getRepository().parseCommit(resolve(id)); return git.getRepository().parseCommit(resolve(id));
} }
@ -481,8 +490,8 @@ public class GitRepository {
private static final Function<PageLink, Comparator<RevCommit>> revCommitComparatorFunction = pageLink -> { private static final Function<PageLink, Comparator<RevCommit>> revCommitComparatorFunction = pageLink -> {
SortOrder sortOrder = pageLink.getSortOrder(); SortOrder sortOrder = pageLink.getSortOrder();
if (sortOrder != null if (sortOrder != null
&& sortOrder.getProperty().equals("timestamp") && sortOrder.getProperty().equals("timestamp")
&& SortOrder.Direction.ASC.equals(sortOrder.getDirection())) { && SortOrder.Direction.ASC.equals(sortOrder.getDirection())) {
return Comparator.comparingInt(RevCommit::getCommitTime); return Comparator.comparingInt(RevCommit::getCommitTime);
} }
return null; return null;

3
dao/src/main/java/org/thingsboard/server/dao/nosql/CassandraAbstractAsyncDao.java

@ -37,6 +37,9 @@ public abstract class CassandraAbstractAsyncDao extends CassandraAbstractDao {
@Value("${cassandra.query.result_processing_threads:50}") @Value("${cassandra.query.result_processing_threads:50}")
private int threadPoolSize; private int threadPoolSize;
@Value("${cassandra.query.max_result_set_size_in_bytes:52428800}")
protected long maxResultSetSizeBytes;
@PostConstruct @PostConstruct
public void startExecutor() { public void startExecutor() {
readResultsProcessingExecutor = ThingsBoardExecutors.newWorkStealingPool(threadPoolSize, "cassandra-callback"); readResultsProcessingExecutor = ThingsBoardExecutors.newWorkStealingPool(threadPoolSize, "cassandra-callback");

92
dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java

@ -52,6 +52,7 @@ import org.thingsboard.server.common.data.kv.TsKvEntry;
import org.thingsboard.server.common.data.kv.TsKvEntryAggWrapper; import org.thingsboard.server.common.data.kv.TsKvEntryAggWrapper;
import org.thingsboard.server.common.data.kv.TsKvQuery; import org.thingsboard.server.common.data.kv.TsKvQuery;
import org.thingsboard.server.dao.model.ModelConstants; import org.thingsboard.server.dao.model.ModelConstants;
import org.thingsboard.server.dao.nosql.ResultSetSizeLimitExceededException;
import org.thingsboard.server.dao.nosql.TbResultSet; import org.thingsboard.server.dao.nosql.TbResultSet;
import org.thingsboard.server.dao.nosql.TbResultSetFuture; import org.thingsboard.server.dao.nosql.TbResultSetFuture;
import org.thingsboard.server.dao.sqlts.AggregationTimeseriesDao; import org.thingsboard.server.dao.sqlts.AggregationTimeseriesDao;
@ -256,7 +257,8 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD
@Override @Override
public void onFailure(Throwable t) { public void onFailure(Throwable t) {
log.error("[{}][{}] Failed to fetch partitions for interval {}-{}", entityId.getEntityType().name(), entityId.getId(), minPartition, maxPartition, t); log.error("[{}][{}][{}] Failed to fetch partitions for interval {}-{}", tenantId, entityId.getEntityType(), entityId.getId(), minPartition, maxPartition, t);
resultFuture.setException(t);
} }
}, readResultsProcessingExecutor); }, readResultsProcessingExecutor);
return resultFuture; return resultFuture;
@ -330,7 +332,8 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD
@Override @Override
public void onFailure(Throwable t) { public void onFailure(Throwable t) {
log.error("[{}][{}] Failed to fetch partitions for interval {}-{}", entityId.getEntityType().name(), entityId.getId(), toPartitionTs(query.getStartTs()), toPartitionTs(query.getEndTs()), t); log.error("[{}][{}][{}] Failed to fetch partitions for interval {}-{}", tenantId, entityId.getEntityType(), entityId.getId(), toPartitionTs(query.getStartTs()), toPartitionTs(query.getEndTs()), t);
resultFuture.setException(t);
} }
}, readResultsProcessingExecutor); }, readResultsProcessingExecutor);
@ -372,8 +375,7 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD
cursor.addData(convertResultToTsKvEntryList(Collections.emptyList())); cursor.addData(convertResultToTsKvEntryList(Collections.emptyList()));
findAllAsyncSequentiallyWithLimit(tenantId, cursor, resultFuture); findAllAsyncSequentiallyWithLimit(tenantId, cursor, resultFuture);
} else { } else {
Futures.addCallback(result.allRows(readResultsProcessingExecutor), new FutureCallback<List<Row>>() { Futures.addCallback(result.allRows(readResultsProcessingExecutor, maxResultSetSizeBytes), new FutureCallback<List<Row>>() {
@Override @Override
public void onSuccess(@Nullable List<Row> result) { public void onSuccess(@Nullable List<Row> result) {
cursor.addData(convertResultToTsKvEntryList(result == null ? Collections.emptyList() : result)); cursor.addData(convertResultToTsKvEntryList(result == null ? Collections.emptyList() : result));
@ -382,7 +384,13 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD
@Override @Override
public void onFailure(Throwable t) { public void onFailure(Throwable t) {
log.error("[{}][{}] Failed to fetch data for query {}-{}", stmt, t); if (t instanceof ResultSetSizeLimitExceededException e) {
log.warn("[{}][{}][{}] Result set size limit exceeded for key [{}], query [{}]: {} bytes, limit {} bytes",
tenantId, cursor.getEntityType(), cursor.getEntityId(), cursor.getKey(), stmt.getPreparedStatement().getQuery(), e.getActualBytes(), e.getLimitBytes());
} else {
log.error("[{}][{}][{}] Failed to fetch data for key [{}], query [{}]", tenantId, cursor.getEntityType(), cursor.getEntityId(), cursor.getKey(), stmt.getPreparedStatement().getQuery(), t);
}
resultFuture.setException(t);
} }
}, readResultsProcessingExecutor); }, readResultsProcessingExecutor);
@ -392,7 +400,8 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD
@Override @Override
public void onFailure(Throwable t) { public void onFailure(Throwable t) {
log.error("[{}][{}] Failed to fetch data for query {}-{}", stmt, t); log.error("[{}][{}][{}] Failed to fetch data for key [{}], query [{}]", tenantId, cursor.getEntityType(), cursor.getEntityId(), cursor.getKey(), stmt.getPreparedStatement().getQuery(), t);
resultFuture.setException(t);
} }
}, readResultsProcessingExecutor); }, readResultsProcessingExecutor);
} }
@ -425,7 +434,7 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD
} }
if (!isUseTsKeyValuePartitioningOnRead()) { if (!isUseTsKeyValuePartitioningOnRead()) {
final long estimatedPartitionCount = estimatePartitionCount(minPartition, maxPartition); final long estimatedPartitionCount = estimatePartitionCount(minPartition, maxPartition);
if (estimatedPartitionCount <= useTsKeyValuePartitioningOnReadMaxEstimatedPartitionCount) { if (estimatedPartitionCount <= useTsKeyValuePartitioningOnReadMaxEstimatedPartitionCount) {
return Futures.immediateFuture(calculatePartitions(minPartition, maxPartition, (int) estimatedPartitionCount)); return Futures.immediateFuture(calculatePartitions(minPartition, maxPartition, (int) estimatedPartitionCount));
} }
} }
@ -446,7 +455,7 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD
} }
List<Long> calculatePartitions(long minPartition, long maxPartition) { List<Long> calculatePartitions(long minPartition, long maxPartition) {
return calculatePartitions(minPartition, maxPartition, 0); return calculatePartitions(minPartition, maxPartition, 0);
} }
List<Long> calculatePartitions(long minPartition, long maxPartition, int estimatedPartitionCount) { List<Long> calculatePartitions(long minPartition, long maxPartition, int estimatedPartitionCount) {
@ -533,6 +542,7 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD
public void onFailure(Throwable t) { public void onFailure(Throwable t) {
} }
} }
private long computeTtl(long ttl) { private long computeTtl(long ttl) {
@ -569,7 +579,8 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD
@Override @Override
public void onFailure(Throwable t) { public void onFailure(Throwable t) {
log.error("[{}][{}] Failed to delete data for query {}-{}", stmt, t); log.error("[{}][{}][{}] Failed to delete data for key [{}], query [{}]", tenantId, cursor.getEntityType(), cursor.getEntityId(), cursor.getKey(), stmt.getPreparedStatement().getQuery(), t);
resultFuture.setException(t);
} }
}, readResultsProcessingExecutor); }, readResultsProcessingExecutor);
} }
@ -581,12 +592,12 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD
try { try {
if (deleteStmt == null) { if (deleteStmt == null) {
deleteStmt = prepare("DELETE FROM " + ModelConstants.TS_KV_CF + deleteStmt = prepare("DELETE FROM " + ModelConstants.TS_KV_CF +
" WHERE " + ModelConstants.ENTITY_TYPE_COLUMN + EQUALS_PARAM " WHERE " + ModelConstants.ENTITY_TYPE_COLUMN + EQUALS_PARAM
+ "AND " + ModelConstants.ENTITY_ID_COLUMN + EQUALS_PARAM + "AND " + ModelConstants.ENTITY_ID_COLUMN + EQUALS_PARAM
+ "AND " + ModelConstants.KEY_COLUMN + EQUALS_PARAM + "AND " + ModelConstants.KEY_COLUMN + EQUALS_PARAM
+ "AND " + ModelConstants.PARTITION_COLUMN + EQUALS_PARAM + "AND " + ModelConstants.PARTITION_COLUMN + EQUALS_PARAM
+ "AND " + ModelConstants.TS_COLUMN + " >= ? " + "AND " + ModelConstants.TS_COLUMN + " >= ? "
+ "AND " + ModelConstants.TS_COLUMN + " < ?"); + "AND " + ModelConstants.TS_COLUMN + " < ?");
} }
} finally { } finally {
stmtCreationLock.unlock(); stmtCreationLock.unlock();
@ -661,13 +672,13 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD
private String getPreparedStatementQuery(DataType type) { private String getPreparedStatementQuery(DataType type) {
return INSERT_INTO + ModelConstants.TS_KV_CF + return INSERT_INTO + ModelConstants.TS_KV_CF +
"(" + ModelConstants.ENTITY_TYPE_COLUMN + "(" + ModelConstants.ENTITY_TYPE_COLUMN +
"," + ModelConstants.ENTITY_ID_COLUMN + "," + ModelConstants.ENTITY_ID_COLUMN +
"," + ModelConstants.KEY_COLUMN + "," + ModelConstants.KEY_COLUMN +
"," + ModelConstants.PARTITION_COLUMN + "," + ModelConstants.PARTITION_COLUMN +
"," + ModelConstants.TS_COLUMN + "," + ModelConstants.TS_COLUMN +
"," + getColumnName(type) + ")" + "," + getColumnName(type) + ")" +
" VALUES(?, ?, ?, ?, ?, ?)"; " VALUES(?, ?, ?, ?, ?, ?)";
} }
private String getPreparedStatementQueryWithTtl(DataType type) { private String getPreparedStatementQueryWithTtl(DataType type) {
@ -680,11 +691,11 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD
try { try {
if (partitionInsertStmt == null) { if (partitionInsertStmt == null) {
partitionInsertStmt = prepare(INSERT_INTO + ModelConstants.TS_KV_PARTITIONS_CF + partitionInsertStmt = prepare(INSERT_INTO + ModelConstants.TS_KV_PARTITIONS_CF +
"(" + ModelConstants.ENTITY_TYPE_COLUMN + "(" + ModelConstants.ENTITY_TYPE_COLUMN +
"," + ModelConstants.ENTITY_ID_COLUMN + "," + ModelConstants.ENTITY_ID_COLUMN +
"," + ModelConstants.PARTITION_COLUMN + "," + ModelConstants.PARTITION_COLUMN +
"," + ModelConstants.KEY_COLUMN + ")" + "," + ModelConstants.KEY_COLUMN + ")" +
" VALUES(?, ?, ?, ?)"); " VALUES(?, ?, ?, ?)");
} }
} finally { } finally {
stmtCreationLock.unlock(); stmtCreationLock.unlock();
@ -699,11 +710,11 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD
try { try {
if (partitionInsertTtlStmt == null) { if (partitionInsertTtlStmt == null) {
partitionInsertTtlStmt = prepare(INSERT_INTO + ModelConstants.TS_KV_PARTITIONS_CF + partitionInsertTtlStmt = prepare(INSERT_INTO + ModelConstants.TS_KV_PARTITIONS_CF +
"(" + ModelConstants.ENTITY_TYPE_COLUMN + "(" + ModelConstants.ENTITY_TYPE_COLUMN +
"," + ModelConstants.ENTITY_ID_COLUMN + "," + ModelConstants.ENTITY_ID_COLUMN +
"," + ModelConstants.PARTITION_COLUMN + "," + ModelConstants.PARTITION_COLUMN +
"," + ModelConstants.KEY_COLUMN + ")" + "," + ModelConstants.KEY_COLUMN + ")" +
" VALUES(?, ?, ?, ?) USING TTL ?"); " VALUES(?, ?, ?, ?) USING TTL ?");
} }
} finally { } finally {
stmtCreationLock.unlock(); stmtCreationLock.unlock();
@ -809,16 +820,17 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD
fetchStmts[type.ordinal()] = fetchStmts[Aggregation.SUM.ordinal()]; fetchStmts[type.ordinal()] = fetchStmts[Aggregation.SUM.ordinal()];
} else { } else {
fetchStmts[type.ordinal()] = prepare(SELECT_PREFIX + fetchStmts[type.ordinal()] = prepare(SELECT_PREFIX +
String.join(", ", ModelConstants.getFetchColumnNames(type)) + " FROM " + ModelConstants.TS_KV_CF String.join(", ", ModelConstants.getFetchColumnNames(type)) + " FROM " + ModelConstants.TS_KV_CF
+ " WHERE " + ModelConstants.ENTITY_TYPE_COLUMN + EQUALS_PARAM + " WHERE " + ModelConstants.ENTITY_TYPE_COLUMN + EQUALS_PARAM
+ "AND " + ModelConstants.ENTITY_ID_COLUMN + EQUALS_PARAM + "AND " + ModelConstants.ENTITY_ID_COLUMN + EQUALS_PARAM
+ "AND " + ModelConstants.KEY_COLUMN + EQUALS_PARAM + "AND " + ModelConstants.KEY_COLUMN + EQUALS_PARAM
+ "AND " + ModelConstants.PARTITION_COLUMN + EQUALS_PARAM + "AND " + ModelConstants.PARTITION_COLUMN + EQUALS_PARAM
+ "AND " + ModelConstants.TS_COLUMN + " >= ? " + "AND " + ModelConstants.TS_COLUMN + " >= ? "
+ "AND " + ModelConstants.TS_COLUMN + " < ?" + "AND " + ModelConstants.TS_COLUMN + " < ?"
+ (type == Aggregation.NONE ? " ORDER BY " + ModelConstants.TS_COLUMN + " " + orderBy + " LIMIT ?" : "")); + (type == Aggregation.NONE ? " ORDER BY " + ModelConstants.TS_COLUMN + " " + orderBy + " LIMIT ?" : ""));
} }
} }
return fetchStmts; return fetchStmts;
} }
} }

2
dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesLatestDao.java

@ -184,7 +184,7 @@ public class CassandraBaseTimeseriesLatestDao extends AbstractCassandraBaseTimes
} }
private ListenableFuture<List<TsKvEntry>> convertAsyncResultSetToTsKvEntryList(TbResultSet rs) { private ListenableFuture<List<TsKvEntry>> convertAsyncResultSetToTsKvEntryList(TbResultSet rs) {
return Futures.transform(rs.allRows(readResultsProcessingExecutor), return Futures.transform(rs.allRows(readResultsProcessingExecutor, maxResultSetSizeBytes),
rows -> this.convertResultToTsKvEntryList(rows), readResultsProcessingExecutor); rows -> this.convertResultToTsKvEntryList(rows), readResultsProcessingExecutor);
} }

4
dao/src/main/java/org/thingsboard/server/dao/timeseries/SimpleListenableFuture.java

@ -26,4 +26,8 @@ public class SimpleListenableFuture<V> extends AbstractFuture<V> {
return super.set(value); return super.set(value);
} }
public boolean setException(Throwable throwable) {
return super.setException(throwable);
}
} }

42
dao/src/test/java/org/thingsboard/server/dao/service/timeseries/nosql/TimeseriesServiceNoSqlTest.java

@ -16,7 +16,10 @@
package org.thingsboard.server.dao.service.timeseries.nosql; package org.thingsboard.server.dao.service.timeseries.nosql;
import com.datastax.oss.driver.api.core.uuid.Uuids; import com.datastax.oss.driver.api.core.uuid.Uuids;
import org.apache.commons.lang3.exception.ExceptionUtils;
import org.junit.Test; import org.junit.Test;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.test.util.ReflectionTestUtils;
import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.kv.Aggregation; import org.thingsboard.server.common.data.kv.Aggregation;
import org.thingsboard.server.common.data.kv.BaseReadTsKvQuery; import org.thingsboard.server.common.data.kv.BaseReadTsKvQuery;
@ -27,9 +30,12 @@ import org.thingsboard.server.common.data.kv.JsonDataEntry;
import org.thingsboard.server.common.data.kv.LongDataEntry; import org.thingsboard.server.common.data.kv.LongDataEntry;
import org.thingsboard.server.common.data.kv.StringDataEntry; import org.thingsboard.server.common.data.kv.StringDataEntry;
import org.thingsboard.server.common.data.kv.TsKvEntry; import org.thingsboard.server.common.data.kv.TsKvEntry;
import org.thingsboard.server.dao.nosql.ResultSetSizeLimitExceededException;
import org.thingsboard.server.dao.service.DaoNoSqlTest; import org.thingsboard.server.dao.service.DaoNoSqlTest;
import org.thingsboard.server.dao.service.timeseries.BaseTimeseriesServiceTest; import org.thingsboard.server.dao.service.timeseries.BaseTimeseriesServiceTest;
import org.thingsboard.server.dao.timeseries.CassandraBaseTimeseriesDao;
import java.util.ArrayList;
import java.util.Collections; import java.util.Collections;
import java.util.List; import java.util.List;
import java.util.concurrent.ExecutionException; import java.util.concurrent.ExecutionException;
@ -37,6 +43,7 @@ import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException; import java.util.concurrent.TimeoutException;
import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.assertThatThrownBy;
import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertFalse; import static org.junit.Assert.assertFalse;
import static org.junit.Assert.assertTrue; import static org.junit.Assert.assertTrue;
@ -44,6 +51,9 @@ import static org.junit.Assert.assertTrue;
@DaoNoSqlTest @DaoNoSqlTest
public class TimeseriesServiceNoSqlTest extends BaseTimeseriesServiceTest { public class TimeseriesServiceNoSqlTest extends BaseTimeseriesServiceTest {
@Autowired
private CassandraBaseTimeseriesDao cassandraBaseTimeseriesDao;
@Test @Test
public void shouldSaveEntryOfEachTypeWithTtl() throws ExecutionException, InterruptedException, TimeoutException { public void shouldSaveEntryOfEachTypeWithTtl() throws ExecutionException, InterruptedException, TimeoutException {
long ttlInSec = TimeUnit.SECONDS.toSeconds(3); long ttlInSec = TimeUnit.SECONDS.toSeconds(3);
@ -94,4 +104,36 @@ public class TimeseriesServiceNoSqlTest extends BaseTimeseriesServiceTest {
double expectedValue = (doubleValue + longValue)/ 2; double expectedValue = (doubleValue + longValue)/ 2;
assertThat(listWithAgg.get(0).getDoubleValue().get()).isEqualTo(expectedValue); assertThat(listWithAgg.get(0).getDoubleValue().get()).isEqualTo(expectedValue);
} }
@Test
public void testResultSetSizeLimitExceeded() throws Exception {
DeviceId deviceId = new DeviceId(Uuids.timeBased());
String value = "x".repeat(500);
List<TsKvEntry> entries = new ArrayList<>();
for (int i = 0; i < 5; i++) {
entries.add(new BasicTsKvEntry(TimeUnit.MINUTES.toMillis(i + 1), new StringDataEntry("bigKey", value)));
}
tsService.save(tenantId, deviceId, entries, 0).get(MAX_TIMEOUT, TimeUnit.SECONDS);
long originalLimit = (long) ReflectionTestUtils.getField(cassandraBaseTimeseriesDao, "maxResultSetSizeBytes");
try {
// Set a very low limit to trigger the exception
ReflectionTestUtils.setField(cassandraBaseTimeseriesDao, "maxResultSetSizeBytes", 1024L);
assertThatThrownBy(() -> tsService.findAll(tenantId, deviceId, Collections.singletonList(
new BaseReadTsKvQuery("bigKey", 0L, TimeUnit.MINUTES.toMillis(6), 1000, 10, Aggregation.NONE)
)).get(MAX_TIMEOUT, TimeUnit.SECONDS))
.isInstanceOf(ExecutionException.class)
.satisfies(e -> assertThat(ExceptionUtils.getRootCause(e)).isInstanceOf(ResultSetSizeLimitExceededException.class));
} finally {
ReflectionTestUtils.setField(cassandraBaseTimeseriesDao, "maxResultSetSizeBytes", originalLimit);
}
// Verify query succeeds with original limit restored
List<TsKvEntry> result = tsService.findAll(tenantId, deviceId, Collections.singletonList(
new BaseReadTsKvQuery("bigKey", 0L, TimeUnit.MINUTES.toMillis(6), 1000, 10, Aggregation.NONE)
)).get(MAX_TIMEOUT, TimeUnit.SECONDS);
assertThat(result).hasSize(5);
}
} }

34
rest-client/src/main/resources/logback.xml

@ -1,34 +0,0 @@
<?xml version="1.0" encoding="UTF-8" ?>
<!--
Copyright © 2016-2026 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.
-->
<!DOCTYPE configuration>
<configuration>
<appender name="STDOUT" class="ch.qos.logback.core.ConsoleAppender">
<encoder>
<pattern>%d{ISO8601} [%thread] %-5level %logger{36} - %msg%n</pattern>
</encoder>
</appender>
<logger name="org.thingsboard.server" level="INFO" />
<root level="INFO">
<appender-ref ref="STDOUT"/>
</root>
</configuration>
Loading…
Cancel
Save