diff --git a/application/pom.xml b/application/pom.xml
index 6979c35b53..eec0d1e90a 100644
--- a/application/pom.xml
+++ b/application/pom.xml
@@ -20,7 +20,7 @@
4.0.0
org.thingsboard
- 1.1.1-SNAPSHOT
+ 1.2.0-SNAPSHOT
thingsboard
org.thingsboard
diff --git a/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java b/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java
index 70bb4f2421..d4c42d82b7 100644
--- a/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java
+++ b/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java
@@ -58,6 +58,7 @@ import java.util.Map;
import java.util.Optional;
import java.util.Set;
import java.util.UUID;
+import java.util.concurrent.ExecutionException;
import java.util.concurrent.TimeoutException;
import java.util.function.Consumer;
import java.util.function.Predicate;
@@ -85,14 +86,22 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso
this.attributeSubscriptions = new HashMap<>();
this.rpcSubscriptions = new HashMap<>();
this.rpcPendingMap = new HashMap<>();
- refreshAttributes();
+ initAttributes();
}
- private void refreshAttributes() {
+ private void initAttributes() {
this.deviceAttributes = new DeviceAttributes(fetchAttributes(DataConstants.CLIENT_SCOPE),
fetchAttributes(DataConstants.SERVER_SCOPE), fetchAttributes(DataConstants.SHARED_SCOPE));
}
+ private void refreshAttributes(DeviceAttributesEventNotificationMsg msg) {
+ if (msg.isDeleted()) {
+ msg.getDeletedKeys().forEach(key -> deviceAttributes.remove(key));
+ } else {
+ deviceAttributes.update(msg.getScope(), msg.getValues());
+ }
+ }
+
void processRpcRequest(ActorContext context, ToDeviceRpcRequestPluginMsg msg) {
ToDeviceRpcRequest request = msg.getMsg();
ToDeviceRpcRequestBody body = request.getBody();
@@ -195,10 +204,8 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso
}
void processAttributesUpdate(ActorContext context, DeviceAttributesEventNotificationMsg msg) {
- //TODO: improve this procedure to fetch only changed attributes.
- refreshAttributes();
- //TODO: support attributes deletion
- Set keys = msg.getKeys();
+ refreshAttributes(msg);
+ Set keys = msg.getDeletedKeys();
if (attributeSubscriptions.size() > 0) {
ToDeviceMsg notification = null;
if (msg.isDeleted()) {
@@ -360,8 +367,14 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso
}
}
- private List fetchAttributes(String attributeType) {
- return systemContext.getAttributesService().findAll(this.deviceId, attributeType);
+ private List fetchAttributes(String scope) {
+ try {
+ //TODO: replace this with async operation. Happens only during actor creation, but is still criticla for performance,
+ return systemContext.getAttributesService().findAll(this.deviceId, scope).get();
+ } catch (InterruptedException | ExecutionException e) {
+ logger.warning("[{}] Failed to fetch attributes for scope: {}", deviceId, scope);
+ throw new RuntimeException(e);
+ }
}
public void processCredentialsUpdate(ActorContext context, DeviceCredentialsUpdateNotificationMsg msg) {
diff --git a/application/src/main/java/org/thingsboard/server/actors/plugin/PluginActorMessageProcessor.java b/application/src/main/java/org/thingsboard/server/actors/plugin/PluginActorMessageProcessor.java
index 2f2f32e233..c3887484c0 100644
--- a/application/src/main/java/org/thingsboard/server/actors/plugin/PluginActorMessageProcessor.java
+++ b/application/src/main/java/org/thingsboard/server/actors/plugin/PluginActorMessageProcessor.java
@@ -173,7 +173,7 @@ public class PluginActorMessageProcessor extends ComponentMsgProcessor
@Override
public void onUpdate(ActorContext context) throws Exception {
- PluginMetaData oldPluginMd = systemContext.getPluginService().findPluginById(entityId);
+ PluginMetaData oldPluginMd = pluginMd;
pluginMd = systemContext.getPluginService().findPluginById(entityId);
boolean requiresRestart = false;
logger.info("[{}] Plugin configuration was updated from {} to {}.", entityId, oldPluginMd, pluginMd);
diff --git a/application/src/main/java/org/thingsboard/server/actors/plugin/PluginProcessingContext.java b/application/src/main/java/org/thingsboard/server/actors/plugin/PluginProcessingContext.java
index 92e0ec4101..bea51dbedb 100644
--- a/application/src/main/java/org/thingsboard/server/actors/plugin/PluginProcessingContext.java
+++ b/application/src/main/java/org/thingsboard/server/actors/plugin/PluginProcessingContext.java
@@ -17,6 +17,7 @@ package org.thingsboard.server.actors.plugin;
import java.io.IOException;
import java.util.*;
+import java.util.concurrent.ExecutionException;
import java.util.concurrent.Executor;
import java.util.concurrent.Executors;
import java.util.stream.Collectors;
@@ -55,6 +56,7 @@ import org.thingsboard.server.extensions.api.plugins.ws.PluginWebsocketSessionRe
import org.thingsboard.server.extensions.api.plugins.ws.msg.PluginWebsocketMsg;
import akka.actor.ActorRef;
+import org.w3c.dom.Attr;
import javax.annotation.Nullable;
@@ -88,109 +90,120 @@ public final class PluginProcessingContext implements PluginContext {
}
@Override
- public void saveAttributes(DeviceId deviceId, String scope, List attributes, PluginCallback callback) {
- validate(deviceId);
- Set keys = new HashSet<>();
- for (AttributeKvEntry attribute : attributes) {
- keys.add(new AttributeKey(scope, attribute.getKey()));
- }
+ public void saveAttributes(final TenantId tenantId, final DeviceId deviceId, final String scope, final List attributes, final PluginCallback callback) {
+ validate(deviceId, new ValidationCallback(callback, ctx -> {
+ ListenableFuture> rsListFuture = pluginCtx.attributesService.save(deviceId, scope, attributes);
+ Futures.addCallback(rsListFuture, getListCallback(callback, v -> {
+ onDeviceAttributesChanged(tenantId, deviceId, scope, attributes);
+ return null;
+ }), executor);
+ }));
+ }
- ListenableFuture> rsListFuture = pluginCtx.attributesService.save(deviceId, scope, attributes);
- Futures.addCallback(rsListFuture, getListCallback(callback, v -> {
- onDeviceAttributesChanged(deviceId, keys);
- return null;
- }), executor);
+ @Override
+ public void removeAttributes(final TenantId tenantId, final DeviceId deviceId, final String scope, final List keys, final PluginCallback callback) {
+ validate(deviceId, new ValidationCallback(callback, ctx -> {
+ ListenableFuture> future = pluginCtx.attributesService.removeAll(deviceId, scope, keys);
+ Futures.addCallback(future, getCallback(callback, v -> null), executor);
+ onDeviceAttributesDeleted(tenantId, deviceId, keys.stream().map(key -> new AttributeKey(scope, key)).collect(Collectors.toSet()));
+ }));
}
@Override
- public Optional loadAttribute(DeviceId deviceId, String attributeType, String attributeKey) {
- validate(deviceId);
- AttributeKvEntry attribute = pluginCtx.attributesService.find(deviceId, attributeType, attributeKey);
- return Optional.ofNullable(attribute);
+ public void loadAttribute(DeviceId deviceId, String attributeType, String attributeKey, final PluginCallback> callback) {
+ validate(deviceId, new ValidationCallback(callback, ctx -> {
+ ListenableFuture> future = pluginCtx.attributesService.find(deviceId, attributeType, attributeKey);
+ Futures.addCallback(future, getCallback(callback, v -> v), executor);
+ }));
}
@Override
- public List loadAttributes(DeviceId deviceId, String attributeType, List attributeKeys) {
- validate(deviceId);
- List result = new ArrayList<>(attributeKeys.size());
- for (String attributeKey : attributeKeys) {
- AttributeKvEntry attribute = pluginCtx.attributesService.find(deviceId, attributeType, attributeKey);
- if (attribute != null) {
- result.add(attribute);
- }
- }
- return result;
+ public void loadAttributes(DeviceId deviceId, String attributeType, Collection attributeKeys, final PluginCallback> callback) {
+ validate(deviceId, new ValidationCallback(callback, ctx -> {
+ ListenableFuture> future = pluginCtx.attributesService.find(deviceId, attributeType, attributeKeys);
+ Futures.addCallback(future, getCallback(callback, v -> v), executor);
+ }));
}
@Override
- public List loadAttributes(DeviceId deviceId, String attributeType) {
- validate(deviceId);
- return pluginCtx.attributesService.findAll(deviceId, attributeType);
+ public void loadAttributes(DeviceId deviceId, String attributeType, PluginCallback> callback) {
+ validate(deviceId, new ValidationCallback(callback, ctx -> {
+ ListenableFuture> future = pluginCtx.attributesService.findAll(deviceId, attributeType);
+ Futures.addCallback(future, getCallback(callback, v -> v), executor);
+ }));
}
@Override
- public void removeAttributes(DeviceId deviceId, String scope, List keys) {
- validate(deviceId);
- pluginCtx.attributesService.removeAll(deviceId, scope, keys);
- onDeviceAttributesDeleted(deviceId, keys.stream().map(key -> new AttributeKey(scope, key)).collect(Collectors.toSet()));
+ public void loadAttributes(final DeviceId deviceId, final Collection attributeTypes, final PluginCallback> callback) {
+ validate(deviceId, new ValidationCallback(callback, ctx -> {
+ List>> futures = new ArrayList<>();
+ attributeTypes.forEach(attributeType -> futures.add(pluginCtx.attributesService.findAll(deviceId, attributeType)));
+ convertFuturesAndAddCallback(callback, futures);
+ }));
}
@Override
- public void saveTsData(DeviceId deviceId, TsKvEntry entry, PluginCallback callback) {
- validate(deviceId);
- ListenableFuture> rsListFuture = pluginCtx.tsService.save(DataConstants.DEVICE, deviceId, entry);
- Futures.addCallback(rsListFuture, getListCallback(callback, v -> null), executor);
+ public void loadAttributes(final DeviceId deviceId, final Collection attributeTypes, final Collection attributeKeys, final PluginCallback> callback) {
+ validate(deviceId, new ValidationCallback(callback, ctx -> {
+ List>> futures = new ArrayList<>();
+ attributeTypes.forEach(attributeType -> futures.add(pluginCtx.attributesService.find(deviceId, attributeType, attributeKeys)));
+ convertFuturesAndAddCallback(callback, futures);
+ }));
}
@Override
- public void saveTsData(DeviceId deviceId, List entries, PluginCallback callback) {
- validate(deviceId);
- ListenableFuture> rsListFuture = pluginCtx.tsService.save(DataConstants.DEVICE, deviceId, entries);
- Futures.addCallback(rsListFuture, getListCallback(callback, v -> null), executor);
+ public void saveTsData(final DeviceId deviceId, final TsKvEntry entry, final PluginCallback callback) {
+ validate(deviceId, new ValidationCallback(callback, ctx -> {
+ ListenableFuture> rsListFuture = pluginCtx.tsService.save(DataConstants.DEVICE, deviceId, entry);
+ Futures.addCallback(rsListFuture, getListCallback(callback, v -> null), executor);
+ }));
}
@Override
- public List loadTimeseries(DeviceId deviceId, TsKvQuery query) {
- validate(deviceId);
- return pluginCtx.tsService.find(DataConstants.DEVICE, deviceId, query);
+ public void saveTsData(final DeviceId deviceId, final List entries, final PluginCallback callback) {
+ validate(deviceId, new ValidationCallback(callback, ctx -> {
+ ListenableFuture> rsListFuture = pluginCtx.tsService.save(DataConstants.DEVICE, deviceId, entries);
+ Futures.addCallback(rsListFuture, getListCallback(callback, v -> null), executor);
+ }));
}
@Override
- public void loadLatestTimeseries(DeviceId deviceId, PluginCallback> callback) {
- validate(deviceId);
- ResultSetFuture future = pluginCtx.tsService.findAllLatest(DataConstants.DEVICE, deviceId);
- Futures.addCallback(future, getCallback(callback, pluginCtx.tsService::convertResultSetToTsKvEntryList), executor);
+ public void loadTimeseries(final DeviceId deviceId, final List queries, final PluginCallback> callback) {
+ validate(deviceId, new ValidationCallback(callback, ctx -> {
+ ListenableFuture> future = pluginCtx.tsService.findAll(DataConstants.DEVICE, deviceId, queries);
+ Futures.addCallback(future, getCallback(callback, v -> v), executor);
+ }));
}
@Override
- public void loadLatestTimeseries(DeviceId deviceId, Collection keys, PluginCallback> callback) {
- validate(deviceId);
- ListenableFuture> rsListFuture = pluginCtx.tsService.findLatest(DataConstants.DEVICE, deviceId, keys);
- Futures.addCallback(rsListFuture, getListCallback(callback, rsList ->
- {
- List result = new ArrayList<>();
- for (ResultSet rs : rsList) {
- Row row = rs.one();
- if (row != null) {
- result.add(pluginCtx.tsService.convertResultToTsKvEntry(row));
- }
- }
- return result;
- }), executor);
+ public void loadLatestTimeseries(final DeviceId deviceId, final PluginCallback> callback) {
+ validate(deviceId, new ValidationCallback(callback, ctx -> {
+ ResultSetFuture future = pluginCtx.tsService.findAllLatest(DataConstants.DEVICE, deviceId);
+ Futures.addCallback(future, getCallback(callback, pluginCtx.tsService::convertResultSetToTsKvEntryList), executor);
+ }));
}
@Override
- public void reply(PluginToRuleMsg> msg) {
- pluginCtx.parentActor.tell(msg, ActorRef.noSender());
+ public void loadLatestTimeseries(final DeviceId deviceId, final Collection keys, final PluginCallback> callback) {
+ validate(deviceId, new ValidationCallback(callback, ctx -> {
+ ListenableFuture> rsListFuture = pluginCtx.tsService.findLatest(DataConstants.DEVICE, deviceId, keys);
+ Futures.addCallback(rsListFuture, getListCallback(callback, rsList ->
+ {
+ List result = new ArrayList<>();
+ for (ResultSet rs : rsList) {
+ Row row = rs.one();
+ if (row != null) {
+ result.add(pluginCtx.tsService.convertResultToTsKvEntry(row));
+ }
+ }
+ return result;
+ }), executor);
+ }));
}
@Override
- public boolean checkAccess(DeviceId deviceId) {
- try {
- return validate(deviceId);
- } catch (IllegalStateException | IllegalArgumentException e) {
- return false;
- }
+ public void reply(PluginToRuleMsg> msg) {
+ pluginCtx.parentActor.tell(msg, ActorRef.noSender());
}
@Override
@@ -203,18 +216,12 @@ public final class PluginProcessingContext implements PluginContext {
return securityCtx;
}
- private void onDeviceAttributesChanged(DeviceId deviceId, AttributeKey key) {
- onDeviceAttributesChanged(deviceId, Collections.singleton(key));
+ private void onDeviceAttributesDeleted(TenantId tenantId, DeviceId deviceId, Set keys) {
+ pluginCtx.toDeviceActor(DeviceAttributesEventNotificationMsg.onDelete(tenantId, deviceId, keys));
}
- private void onDeviceAttributesDeleted(DeviceId deviceId, Set keys) {
- Device device = pluginCtx.deviceService.findDeviceById(deviceId);
- pluginCtx.toDeviceActor(DeviceAttributesEventNotificationMsg.onDelete(device.getTenantId(), deviceId, keys));
- }
-
- private void onDeviceAttributesChanged(DeviceId deviceId, Set keys) {
- Device device = pluginCtx.deviceService.findDeviceById(deviceId);
- pluginCtx.toDeviceActor(DeviceAttributesEventNotificationMsg.onUpdate(device.getTenantId(), deviceId, keys));
+ private void onDeviceAttributesChanged(TenantId tenantId, DeviceId deviceId, String scope, List values) {
+ pluginCtx.toDeviceActor(DeviceAttributesEventNotificationMsg.onUpdate(tenantId, deviceId, scope, values));
}
private FutureCallback> getListCallback(final PluginCallback callback, Function, T> transformer) {
@@ -235,11 +242,15 @@ public final class PluginProcessingContext implements PluginContext {
};
}
- private FutureCallback getCallback(final PluginCallback callback, Function transformer) {
- return new FutureCallback() {
+ private FutureCallback getCallback(final PluginCallback callback, Function transformer) {
+ return new FutureCallback() {
@Override
- public void onSuccess(@Nullable ResultSet result) {
- pluginCtx.self().tell(PluginCallbackMessage.onSuccess(callback, transformer.apply(result)), ActorRef.noSender());
+ public void onSuccess(@Nullable R result) {
+ try {
+ pluginCtx.self().tell(PluginCallbackMessage.onSuccess(callback, transformer.apply(result)), ActorRef.noSender());
+ } catch (Exception e) {
+ pluginCtx.self().tell(PluginCallbackMessage.onError(callback, e), ActorRef.noSender());
+ }
}
@Override
@@ -253,26 +264,35 @@ public final class PluginProcessingContext implements PluginContext {
};
}
- // TODO: replace with our own exceptions
- private boolean validate(DeviceId deviceId) {
+ @Override
+ public void checkAccess(DeviceId deviceId, PluginCallback callback) {
+ validate(deviceId, new ValidationCallback(callback, ctx -> callback.onSuccess(ctx, null)));
+ }
+
+ private void validate(DeviceId deviceId, ValidationCallback callback) {
if (securityCtx.isPresent()) {
- PluginApiCallSecurityContext ctx = securityCtx.get();
+ final PluginApiCallSecurityContext ctx = securityCtx.get();
if (ctx.isTenantAdmin() || ctx.isCustomerUser()) {
- Device device = pluginCtx.deviceService.findDeviceById(deviceId);
- if (device == null) {
- throw new IllegalStateException("Device not found!");
- } else {
- if (!device.getTenantId().equals(ctx.getTenantId())) {
- throw new IllegalArgumentException("Device belongs to different tenant!");
- } else if (ctx.isCustomerUser() && !device.getCustomerId().equals(ctx.getCustomerId())) {
- throw new IllegalArgumentException("Device belongs to different customer!");
+ ListenableFuture deviceFuture = pluginCtx.deviceService.findDeviceByIdAsync(deviceId);
+ Futures.addCallback(deviceFuture, getCallback(callback, device -> {
+ if (device == null) {
+ return Boolean.FALSE;
+ } else {
+ if (!device.getTenantId().equals(ctx.getTenantId())) {
+ return Boolean.FALSE;
+ } else if (ctx.isCustomerUser() && !device.getCustomerId().equals(ctx.getCustomerId())) {
+ return Boolean.FALSE;
+ } else {
+ return Boolean.TRUE;
+ }
}
- }
+ }));
} else {
- return false;
+ callback.onSuccess(this, Boolean.FALSE);
}
+ } else {
+ callback.onSuccess(this, Boolean.TRUE);
}
- return true;
}
@Override
@@ -282,9 +302,8 @@ public final class PluginProcessingContext implements PluginContext {
@Override
public void getDevice(DeviceId deviceId, PluginCallback callback) {
- //TODO: add caching here with async api.
- Device device = pluginCtx.deviceService.findDeviceById(deviceId);
- pluginCtx.self().tell(PluginCallbackMessage.onSuccess(callback, device), ActorRef.noSender());
+ ListenableFuture deviceFuture = pluginCtx.deviceService.findDeviceByIdAsync(deviceId);
+ Futures.addCallback(deviceFuture, getCallback(callback, v -> v));
}
@Override
@@ -303,4 +322,15 @@ public final class PluginProcessingContext implements PluginContext {
public void scheduleTimeoutMsg(TimeoutMsg msg) {
pluginCtx.scheduleTimeoutMsg(msg);
}
+
+
+ private void convertFuturesAndAddCallback(PluginCallback> callback, List>> futures) {
+ ListenableFuture> future = Futures.transform(Futures.successfulAsList(futures),
+ (Function super List>, ? extends List>) input -> {
+ List result = new ArrayList<>();
+ input.forEach(r -> result.addAll(r));
+ return result;
+ }, executor);
+ Futures.addCallback(future, getCallback(callback, v -> v), executor);
+ }
}
diff --git a/application/src/main/java/org/thingsboard/server/actors/plugin/ValidationCallback.java b/application/src/main/java/org/thingsboard/server/actors/plugin/ValidationCallback.java
new file mode 100644
index 0000000000..707afa50dd
--- /dev/null
+++ b/application/src/main/java/org/thingsboard/server/actors/plugin/ValidationCallback.java
@@ -0,0 +1,49 @@
+/**
+ * Copyright © 2016-2017 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.actors.plugin;
+
+import com.hazelcast.util.function.Consumer;
+import org.thingsboard.server.extensions.api.exception.UnauthorizedException;
+import org.thingsboard.server.extensions.api.plugins.PluginCallback;
+import org.thingsboard.server.extensions.api.plugins.PluginContext;
+
+/**
+ * Created by ashvayka on 21.02.17.
+ */
+public class ValidationCallback implements PluginCallback {
+
+ private final PluginCallback> callback;
+ private final Consumer action;
+
+ public ValidationCallback(PluginCallback> callback, Consumer action) {
+ this.callback = callback;
+ this.action = action;
+ }
+
+ @Override
+ public void onSuccess(PluginContext ctx, Boolean value) {
+ if (value) {
+ action.accept(ctx);
+ } else {
+ onFailure(ctx, new UnauthorizedException());
+ }
+ }
+
+ @Override
+ public void onFailure(PluginContext ctx, Exception e) {
+ callback.onFailure(ctx, e);
+ }
+}
diff --git a/application/src/main/java/org/thingsboard/server/controller/BaseController.java b/application/src/main/java/org/thingsboard/server/controller/BaseController.java
index 5647b52e96..b66cd89084 100644
--- a/application/src/main/java/org/thingsboard/server/controller/BaseController.java
+++ b/application/src/main/java/org/thingsboard/server/controller/BaseController.java
@@ -157,6 +157,16 @@ public abstract class BaseController {
}
}
+ void checkArrayParameter(String name, String[] params) throws ThingsboardException {
+ if (params == null || params.length == 0) {
+ throw new ThingsboardException("Parameter '" + name + "' can't be empty!", ThingsboardErrorCode.BAD_REQUEST_PARAMS);
+ } else {
+ for (String param : params) {
+ checkParameter(name, param);
+ }
+ }
+ }
+
UUID toUUID(String id) {
return UUID.fromString(id);
}
diff --git a/application/src/main/java/org/thingsboard/server/controller/DashboardController.java b/application/src/main/java/org/thingsboard/server/controller/DashboardController.java
index 77898ceb7a..fc419b7020 100644
--- a/application/src/main/java/org/thingsboard/server/controller/DashboardController.java
+++ b/application/src/main/java/org/thingsboard/server/controller/DashboardController.java
@@ -32,6 +32,13 @@ import org.thingsboard.server.exception.ThingsboardException;
@RequestMapping("/api")
public class DashboardController extends BaseController {
+ @PreAuthorize("hasAnyAuthority('SYS_ADMIN', 'TENANT_ADMIN', 'CUSTOMER_USER')")
+ @RequestMapping(value = "/dashboard/serverTime", method = RequestMethod.GET)
+ @ResponseBody
+ public long getServerTime() throws ThingsboardException {
+ return System.currentTimeMillis();
+ }
+
@PreAuthorize("hasAnyAuthority('TENANT_ADMIN', 'CUSTOMER_USER')")
@RequestMapping(value = "/dashboard/{dashboardId}", method = RequestMethod.GET)
@ResponseBody
diff --git a/application/src/main/java/org/thingsboard/server/controller/DeviceController.java b/application/src/main/java/org/thingsboard/server/controller/DeviceController.java
index 57433214ce..b08a9640a7 100644
--- a/application/src/main/java/org/thingsboard/server/controller/DeviceController.java
+++ b/application/src/main/java/org/thingsboard/server/controller/DeviceController.java
@@ -15,6 +15,7 @@
*/
package org.thingsboard.server.controller;
+import com.google.common.util.concurrent.ListenableFuture;
import org.springframework.http.HttpStatus;
import org.springframework.security.access.prepost.PreAuthorize;
import org.springframework.web.bind.annotation.*;
@@ -22,6 +23,7 @@ import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.id.CustomerId;
import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.id.TenantId;
+import org.thingsboard.server.common.data.id.UUIDBased;
import org.thingsboard.server.common.data.page.TextPageData;
import org.thingsboard.server.common.data.page.TextPageLink;
import org.thingsboard.server.common.data.security.DeviceCredentials;
@@ -29,6 +31,11 @@ import org.thingsboard.server.dao.exception.IncorrectParameterException;
import org.thingsboard.server.dao.model.ModelConstants;
import org.thingsboard.server.exception.ThingsboardException;
import org.thingsboard.server.extensions.api.device.DeviceCredentialsUpdateNotificationMsg;
+import org.thingsboard.server.service.security.model.SecurityUser;
+
+import java.util.ArrayList;
+import java.util.List;
+import java.util.UUID;
@RestController
@RequestMapping("/api")
@@ -189,4 +196,30 @@ public class DeviceController extends BaseController {
throw handleException(e);
}
}
+
+ @PreAuthorize("hasAnyAuthority('TENANT_ADMIN', 'CUSTOMER_USER')")
+ @RequestMapping(value = "/devices", params = {"deviceIds"}, method = RequestMethod.GET)
+ @ResponseBody
+ public List getDevicesByIds(
+ @RequestParam("deviceIds") String[] strDeviceIds) throws ThingsboardException {
+ checkArrayParameter("deviceIds", strDeviceIds);
+ try {
+ SecurityUser user = getCurrentUser();
+ TenantId tenantId = user.getTenantId();
+ CustomerId customerId = user.getCustomerId();
+ List deviceIds = new ArrayList<>();
+ for (String strDeviceId : strDeviceIds) {
+ deviceIds.add(new DeviceId(toUUID(strDeviceId)));
+ }
+ ListenableFuture> devices;
+ if (customerId == null || customerId.isNullUid()) {
+ devices = deviceService.findDevicesByTenantIdAndIdsAsync(tenantId, deviceIds);
+ } else {
+ devices = deviceService.findDevicesByTenantIdCustomerIdAndIdsAsync(tenantId, customerId, deviceIds);
+ }
+ return checkNotNull(devices.get());
+ } catch (Exception e) {
+ throw handleException(e);
+ }
+ }
}
diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml
index d2cac80101..42b37db793 100644
--- a/application/src/main/resources/thingsboard.yml
+++ b/application/src/main/resources/thingsboard.yml
@@ -140,7 +140,7 @@ cassandra:
# Specify partitioning size for timestamp key-value storage. Example MINUTES, HOURS, DAYS, MONTHS
ts_key_value_partitioning: "${TS_KV_PARTITIONING:MONTHS}"
# Specify max data points per request
- max_limit_per_request: "${TS_KV_MAX_LIMIT_PER_REQUEST:86400}"
+ min_aggregation_step_ms: "${TS_KV_MIN_AGGREGATION_STEP_MS:100}"
# Actor system parameters
actors:
diff --git a/application/src/test/java/org/thingsboard/server/actors/DefaultActorServiceTest.java b/application/src/test/java/org/thingsboard/server/actors/DefaultActorServiceTest.java
index f6cb4e31da..2940a62a93 100644
--- a/application/src/test/java/org/thingsboard/server/actors/DefaultActorServiceTest.java
+++ b/application/src/test/java/org/thingsboard/server/actors/DefaultActorServiceTest.java
@@ -22,6 +22,7 @@ import static org.mockito.Mockito.when;
import java.util.*;
+import com.google.common.util.concurrent.Futures;
import org.thingsboard.server.actors.service.DefaultActorService;
import org.thingsboard.server.common.data.id.*;
import org.thingsboard.server.common.data.kv.TsKvEntry;
@@ -226,7 +227,9 @@ public class DefaultActorServiceTest {
when(pluginMock.getConfiguration()).thenReturn(pluginAdditionalInfo);
when(pluginMock.getClazz()).thenReturn(TelemetryStoragePlugin.class.getName());
- when(attributesService.findAll(deviceId, DataConstants.CLIENT_SCOPE)).thenReturn(Collections.emptyList());
+ when(attributesService.findAll(deviceId, DataConstants.CLIENT_SCOPE)).thenReturn(Futures.immediateFuture(Collections.emptyList()));
+ when(attributesService.findAll(deviceId, DataConstants.SHARED_SCOPE)).thenReturn(Futures.immediateFuture(Collections.emptyList()));
+ when(attributesService.findAll(deviceId, DataConstants.SERVER_SCOPE)).thenReturn(Futures.immediateFuture(Collections.emptyList()));
initActorSystem();
Thread.sleep(1000);
diff --git a/common/data/pom.xml b/common/data/pom.xml
index 19cfa515b6..2849541db5 100644
--- a/common/data/pom.xml
+++ b/common/data/pom.xml
@@ -20,7 +20,7 @@
4.0.0
org.thingsboard
- 1.1.1-SNAPSHOT
+ 1.2.0-SNAPSHOT
common
org.thingsboard.common
diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/kv/Aggregation.java b/common/data/src/main/java/org/thingsboard/server/common/data/kv/Aggregation.java
new file mode 100644
index 0000000000..479a49ac91
--- /dev/null
+++ b/common/data/src/main/java/org/thingsboard/server/common/data/kv/Aggregation.java
@@ -0,0 +1,25 @@
+/**
+ * Copyright © 2016-2017 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.common.data.kv;
+
+/**
+ * Created by ashvayka on 20.02.17.
+ */
+public enum Aggregation {
+
+ MIN, MAX, AVG, SUM, COUNT, NONE;
+
+}
diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/kv/BaseTsKvQuery.java b/common/data/src/main/java/org/thingsboard/server/common/data/kv/BaseTsKvQuery.java
index 7bb1a3fca5..e95496b48b 100644
--- a/common/data/src/main/java/org/thingsboard/server/common/data/kv/BaseTsKvQuery.java
+++ b/common/data/src/main/java/org/thingsboard/server/common/data/kv/BaseTsKvQuery.java
@@ -15,59 +15,29 @@
*/
package org.thingsboard.server.common.data.kv;
-import java.util.Optional;
+import lombok.Data;
+@Data
public class BaseTsKvQuery implements TsKvQuery {
- private String key;
- private Optional startTs;
- private Optional endTs;
- private Optional limit;
+ private final String key;
+ private final long startTs;
+ private final long endTs;
+ private final long interval;
+ private final int limit;
+ private final Aggregation aggregation;
- public BaseTsKvQuery(String key, Optional startTs, Optional endTs, Optional limit) {
+ public BaseTsKvQuery(String key, long startTs, long endTs, long interval, int limit, Aggregation aggregation) {
this.key = key;
this.startTs = startTs;
this.endTs = endTs;
+ this.interval = interval;
this.limit = limit;
- }
-
- public BaseTsKvQuery(String key, Long startTs, Long endTs, Integer limit) {
- this(key, Optional.ofNullable(startTs), Optional.ofNullable(endTs), Optional.ofNullable(limit));
- }
-
- public BaseTsKvQuery(String key, Long startTs, Integer limit) {
- this(key, startTs, null, limit);
- }
-
- public BaseTsKvQuery(String key, Long startTs, Long endTs) {
- this(key, startTs, endTs, null);
- }
-
- public BaseTsKvQuery(String key, Long startTs) {
- this(key, startTs, null, null);
+ this.aggregation = aggregation;
}
- public BaseTsKvQuery(String key, Integer limit) {
- this(key, null, null, limit);
+ public BaseTsKvQuery(String key, long startTs, long endTs) {
+ this(key, startTs, endTs, endTs-startTs, 1, Aggregation.AVG);
}
- @Override
- public String getKey() {
- return key;
- }
-
- @Override
- public Optional getStartTs() {
- return startTs;
- }
-
- @Override
- public Optional getEndTs() {
- return endTs;
- }
-
- @Override
- public Optional getLimit() {
- return limit;
- }
}
diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/kv/TsKvQuery.java b/common/data/src/main/java/org/thingsboard/server/common/data/kv/TsKvQuery.java
index 1303117d63..8d60f525f4 100644
--- a/common/data/src/main/java/org/thingsboard/server/common/data/kv/TsKvQuery.java
+++ b/common/data/src/main/java/org/thingsboard/server/common/data/kv/TsKvQuery.java
@@ -21,10 +21,14 @@ public interface TsKvQuery {
String getKey();
- Optional getStartTs();
+ long getStartTs();
- Optional getEndTs();
+ long getEndTs();
- Optional getLimit();
+ long getInterval();
+
+ int getLimit();
+
+ Aggregation getAggregation();
}
diff --git a/common/message/pom.xml b/common/message/pom.xml
index 5c8cd8cc17..3ff9fda7a0 100644
--- a/common/message/pom.xml
+++ b/common/message/pom.xml
@@ -20,7 +20,7 @@
4.0.0
org.thingsboard
- 1.1.1-SNAPSHOT
+ 1.2.0-SNAPSHOT
common
org.thingsboard.common
diff --git a/common/pom.xml b/common/pom.xml
index 01183e0353..d52152c78c 100644
--- a/common/pom.xml
+++ b/common/pom.xml
@@ -20,7 +20,7 @@
4.0.0
org.thingsboard
- 1.1.1-SNAPSHOT
+ 1.2.0-SNAPSHOT
thingsboard
org.thingsboard
diff --git a/common/transport/pom.xml b/common/transport/pom.xml
index 119310eba8..ca9e384fed 100644
--- a/common/transport/pom.xml
+++ b/common/transport/pom.xml
@@ -20,7 +20,7 @@
4.0.0
org.thingsboard
- 1.1.1-SNAPSHOT
+ 1.2.0-SNAPSHOT
common
org.thingsboard.common
diff --git a/dao/pom.xml b/dao/pom.xml
index 96e77a2130..e7a4286da5 100644
--- a/dao/pom.xml
+++ b/dao/pom.xml
@@ -20,7 +20,7 @@
4.0.0
org.thingsboard
- 1.1.1-SNAPSHOT
+ 1.2.0-SNAPSHOT
thingsboard
org.thingsboard
diff --git a/dao/src/main/java/org/thingsboard/server/dao/AbstractAsyncDao.java b/dao/src/main/java/org/thingsboard/server/dao/AbstractAsyncDao.java
new file mode 100644
index 0000000000..9b9368d45d
--- /dev/null
+++ b/dao/src/main/java/org/thingsboard/server/dao/AbstractAsyncDao.java
@@ -0,0 +1,42 @@
+/**
+ * Copyright © 2016-2017 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 javax.annotation.PostConstruct;
+import javax.annotation.PreDestroy;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+
+/**
+ * Created by ashvayka on 21.02.17.
+ */
+public abstract class AbstractAsyncDao extends AbstractDao {
+
+ protected ExecutorService readResultsProcessingExecutor;
+
+ @PostConstruct
+ public void startExecutor() {
+ readResultsProcessingExecutor = Executors.newCachedThreadPool();
+ }
+
+ @PreDestroy
+ public void stopExecutor() {
+ if (readResultsProcessingExecutor != null) {
+ readResultsProcessingExecutor.shutdownNow();
+ }
+ }
+
+}
diff --git a/dao/src/main/java/org/thingsboard/server/dao/AbstractModelDao.java b/dao/src/main/java/org/thingsboard/server/dao/AbstractModelDao.java
index 73a553767c..01346b0da9 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/AbstractModelDao.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/AbstractModelDao.java
@@ -16,17 +16,22 @@
package org.thingsboard.server.dao;
import com.datastax.driver.core.ResultSet;
+import com.datastax.driver.core.ResultSetFuture;
import com.datastax.driver.core.Statement;
import com.datastax.driver.core.querybuilder.QueryBuilder;
import com.datastax.driver.core.querybuilder.Select;
import com.datastax.driver.core.utils.UUIDs;
import com.datastax.driver.mapping.Mapper;
import com.datastax.driver.mapping.Result;
+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.thingsboard.server.dao.model.BaseEntity;
import org.thingsboard.server.dao.model.wrapper.EntityResultSet;
import org.thingsboard.server.dao.model.ModelConstants;
+import javax.annotation.Nullable;
import java.util.Collections;
import java.util.List;
import java.util.UUID;
@@ -59,6 +64,27 @@ public abstract class AbstractModelDao> extends Abstract
return list;
}
+ protected ListenableFuture> findListByStatementAsync(Statement statement) {
+ if (statement != null) {
+ statement.setConsistencyLevel(cluster.getDefaultReadConsistencyLevel());
+ ResultSetFuture resultSetFuture = getSession().executeAsync(statement);
+ ListenableFuture> result = Futures.transform(resultSetFuture, new Function>() {
+ @Nullable
+ @Override
+ public List apply(@Nullable ResultSet resultSet) {
+ Result result = getMapper().map(resultSet);
+ if (result != null) {
+ return result.all();
+ } else {
+ return Collections.emptyList();
+ }
+ }
+ });
+ return result;
+ }
+ return Futures.immediateFuture(Collections.emptyList());
+ }
+
protected T findOneByStatement(Statement statement) {
T object = null;
if (statement != null) {
@@ -72,6 +98,27 @@ public abstract class AbstractModelDao> extends Abstract
return object;
}
+ protected ListenableFuture findOneByStatementAsync(Statement statement) {
+ if (statement != null) {
+ statement.setConsistencyLevel(cluster.getDefaultReadConsistencyLevel());
+ ResultSetFuture resultSetFuture = getSession().executeAsync(statement);
+ ListenableFuture result = Futures.transform(resultSetFuture, new Function() {
+ @Nullable
+ @Override
+ public T apply(@Nullable ResultSet resultSet) {
+ Result result = getMapper().map(resultSet);
+ if (result != null) {
+ return result.one();
+ } else {
+ return null;
+ }
+ }
+ });
+ return result;
+ }
+ return Futures.immediateFuture(null);
+ }
+
protected Statement getSaveQuery(T dto) {
return getMapper().saveQuery(dto);
}
@@ -100,6 +147,14 @@ public abstract class AbstractModelDao> extends Abstract
return findOneByStatement(query);
}
+ public ListenableFuture findByIdAsync(UUID key) {
+ log.debug("Get entity by key {}", key);
+ Select.Where query = select().from(getColumnFamilyName()).where(eq(ModelConstants.ID_PROPERTY, key));
+ log.trace("Execute query {}", query);
+ return findOneByStatementAsync(query);
+ }
+
+
public ResultSet removeById(UUID key) {
Statement delete = QueryBuilder.delete().all().from(getColumnFamilyName()).where(eq(ModelConstants.ID_PROPERTY, key));
log.debug("Remove request: {}", delete.toString());
diff --git a/dao/src/main/java/org/thingsboard/server/dao/Dao.java b/dao/src/main/java/org/thingsboard/server/dao/Dao.java
index 7aa35e7c0e..2703cdc23a 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/Dao.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/Dao.java
@@ -16,6 +16,7 @@
package org.thingsboard.server.dao;
import com.datastax.driver.core.ResultSet;
+import com.google.common.util.concurrent.ListenableFuture;
import java.util.List;
import java.util.UUID;
@@ -26,6 +27,8 @@ public interface Dao {
T findById(UUID id);
+ ListenableFuture findByIdAsync(UUID id);
+
T save(T t);
ResultSet removeById(UUID id);
diff --git a/dao/src/main/java/org/thingsboard/server/dao/DaoUtil.java b/dao/src/main/java/org/thingsboard/server/dao/DaoUtil.java
index fa98b86b59..27499bb1bd 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/DaoUtil.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/DaoUtil.java
@@ -56,4 +56,12 @@ public abstract class DaoUtil {
return id;
}
+ public static List toUUIDs(List extends UUIDBased> idBasedIds) {
+ List ids = new ArrayList<>();
+ for (UUIDBased idBased : idBasedIds) {
+ ids.add(getId(idBased));
+ }
+ return ids;
+ }
+
}
diff --git a/dao/src/main/java/org/thingsboard/server/dao/attributes/AttributesDao.java b/dao/src/main/java/org/thingsboard/server/dao/attributes/AttributesDao.java
index ead2c044cc..ae58d4d9c6 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/attributes/AttributesDao.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/attributes/AttributesDao.java
@@ -15,23 +15,28 @@
*/
package org.thingsboard.server.dao.attributes;
+import com.datastax.driver.core.ResultSet;
import com.datastax.driver.core.ResultSetFuture;
+import com.google.common.util.concurrent.ListenableFuture;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.kv.AttributeKvEntry;
+import java.util.Collection;
import java.util.List;
-import java.util.UUID;
+import java.util.Optional;
/**
* @author Andrew Shvayka
*/
public interface AttributesDao {
- AttributeKvEntry find(EntityId entityId, String attributeType, String attributeKey);
+ ListenableFuture> find(EntityId entityId, String attributeType, String attributeKey);
- List findAll(EntityId entityId, String attributeType);
+ ListenableFuture> find(EntityId entityId, String attributeType, Collection attributeKey);
+
+ ListenableFuture> findAll(EntityId entityId, String attributeType);
ResultSetFuture save(EntityId entityId, String attributeType, AttributeKvEntry attribute);
- void removeAll(EntityId entityId, String scope, List keys);
+ ListenableFuture> removeAll(EntityId entityId, String scope, List keys);
}
diff --git a/dao/src/main/java/org/thingsboard/server/dao/attributes/AttributesService.java b/dao/src/main/java/org/thingsboard/server/dao/attributes/AttributesService.java
index 5a1fd70bc7..6bf9fb2bd5 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/attributes/AttributesService.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/attributes/AttributesService.java
@@ -23,18 +23,22 @@ import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.UUIDBased;
import org.thingsboard.server.common.data.kv.AttributeKvEntry;
+import java.util.Collection;
import java.util.List;
+import java.util.Optional;
/**
* @author Andrew Shvayka
*/
public interface AttributesService {
- AttributeKvEntry find(EntityId entityId, String scope, String attributeKey);
+ ListenableFuture> find(EntityId entityId, String scope, String attributeKey);
- List findAll(EntityId entityId, String scope);
+ ListenableFuture> find(EntityId entityId, String scope, Collection attributeKeys);
+
+ ListenableFuture> findAll(EntityId entityId, String scope);
ListenableFuture> save(EntityId entityId, String scope, List attributes);
- void removeAll(EntityId entityId, String scope, List attributeKeys);
+ ListenableFuture> removeAll(EntityId entityId, String scope, List attributeKeys);
}
diff --git a/dao/src/main/java/org/thingsboard/server/dao/attributes/BaseAttributesDao.java b/dao/src/main/java/org/thingsboard/server/dao/attributes/BaseAttributesDao.java
index 5148a121de..fd50f4d2ef 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/attributes/BaseAttributesDao.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/attributes/BaseAttributesDao.java
@@ -18,19 +18,24 @@ package org.thingsboard.server.dao.attributes;
import com.datastax.driver.core.*;
import com.datastax.driver.core.querybuilder.QueryBuilder;
import com.datastax.driver.core.querybuilder.Select;
+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.springframework.stereotype.Component;
import org.thingsboard.server.common.data.id.EntityId;
-import org.thingsboard.server.common.data.kv.DataType;
-import org.thingsboard.server.dao.AbstractDao;
+import org.thingsboard.server.dao.AbstractAsyncDao;
import org.thingsboard.server.dao.model.ModelConstants;
-import org.slf4j.Logger;
-import org.slf4j.LoggerFactory;
import org.thingsboard.server.common.data.kv.*;
import org.thingsboard.server.dao.timeseries.BaseTimeseriesDao;
+import javax.annotation.PostConstruct;
+import javax.annotation.PreDestroy;
import java.util.ArrayList;
+import java.util.Collection;
import java.util.List;
+import java.util.Optional;
+import java.util.stream.Collectors;
import static org.thingsboard.server.dao.model.ModelConstants.*;
import static com.datastax.driver.core.querybuilder.QueryBuilder.*;
@@ -40,29 +45,55 @@ import static com.datastax.driver.core.querybuilder.QueryBuilder.*;
*/
@Component
@Slf4j
-public class BaseAttributesDao extends AbstractDao implements AttributesDao {
-
+public class BaseAttributesDao extends AbstractAsyncDao implements AttributesDao {
+
private PreparedStatement saveStmt;
+ @PostConstruct
+ public void init() {
+ super.startExecutor();
+ }
+
+ @PreDestroy
+ public void stop() {
+ super.stopExecutor();
+ }
+
@Override
- public AttributeKvEntry find(EntityId entityId, String attributeType, String attributeKey) {
+ public ListenableFuture> find(EntityId entityId, String attributeType, String attributeKey) {
Select.Where select = select().from(ATTRIBUTES_KV_CF)
.where(eq(ENTITY_TYPE_COLUMN, entityId.getEntityType()))
.and(eq(ENTITY_ID_COLUMN, entityId.getId()))
.and(eq(ATTRIBUTE_TYPE_COLUMN, attributeType))
.and(eq(ATTRIBUTE_KEY_COLUMN, attributeKey));
log.trace("Generated query [{}] for entityId {} and key {}", select, entityId, attributeKey);
- return convertResultToAttributesKvEntry(attributeKey, executeRead(select).one());
+ return Futures.transform(executeAsyncRead(select), (Function super ResultSet, ? extends Optional>) input ->
+ Optional.ofNullable(convertResultToAttributesKvEntry(attributeKey, input.one()))
+ , readResultsProcessingExecutor);
}
@Override
- public List findAll(EntityId entityId, String attributeType) {
+ public ListenableFuture> find(EntityId entityId, String attributeType, Collection attributeKeys) {
+ List>> entries = new ArrayList<>();
+ attributeKeys.forEach(attributeKey -> entries.add(find(entityId, attributeType, attributeKey)));
+ return Futures.transform(Futures.allAsList(entries), (Function>, ? extends List>) input -> {
+ List result = new ArrayList<>();
+ input.stream().filter(opt -> opt.isPresent()).forEach(opt -> result.add(opt.get()));
+ return result;
+ }, readResultsProcessingExecutor);
+ }
+
+
+ @Override
+ public ListenableFuture> findAll(EntityId entityId, String attributeType) {
Select.Where select = select().from(ATTRIBUTES_KV_CF)
.where(eq(ENTITY_TYPE_COLUMN, entityId.getEntityType()))
.and(eq(ENTITY_ID_COLUMN, entityId.getId()))
.and(eq(ATTRIBUTE_TYPE_COLUMN, attributeType));
log.trace("Generated query [{}] for entityId {} and attributeType {}", select, entityId, attributeType);
- return convertResultToAttributesKvEntryList(executeRead(select));
+ return Futures.transform(executeAsyncRead(select), (Function super ResultSet, ? extends List>) input ->
+ convertResultToAttributesKvEntryList(input)
+ , readResultsProcessingExecutor);
}
@Override
@@ -93,20 +124,19 @@ public class BaseAttributesDao extends AbstractDao implements AttributesDao {
}
@Override
- public void removeAll(EntityId entityId, String attributeType, List keys) {
- for (String key : keys) {
- delete(entityId, attributeType, key);
- }
+ public ListenableFuture> removeAll(EntityId entityId, String attributeType, List keys) {
+ List futures = keys.stream().map(key -> delete(entityId, attributeType, key)).collect(Collectors.toList());
+ return Futures.allAsList(futures);
}
- private void delete(EntityId entityId, String attributeType, String key) {
+ private ResultSetFuture delete(EntityId entityId, String attributeType, String key) {
Statement delete = QueryBuilder.delete().all().from(ModelConstants.ATTRIBUTES_KV_CF)
.where(eq(ENTITY_TYPE_COLUMN, entityId.getEntityType()))
.and(eq(ENTITY_ID_COLUMN, entityId.getId()))
.and(eq(ATTRIBUTE_TYPE_COLUMN, attributeType))
.and(eq(ATTRIBUTE_KEY_COLUMN, key));
log.debug("Remove request: {}", delete.toString());
- getSession().execute(delete);
+ return getSession().executeAsync(delete);
}
private PreparedStatement getSaveStmt() {
@@ -150,5 +180,4 @@ public class BaseAttributesDao extends AbstractDao implements AttributesDao {
}
return entries;
}
-
}
diff --git a/dao/src/main/java/org/thingsboard/server/dao/attributes/BaseAttributesService.java b/dao/src/main/java/org/thingsboard/server/dao/attributes/BaseAttributesService.java
index d3a1cb3550..43612419d0 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/attributes/BaseAttributesService.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/attributes/BaseAttributesService.java
@@ -27,7 +27,9 @@ import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
import org.thingsboard.server.dao.service.Validator;
+import java.util.Collection;
import java.util.List;
+import java.util.Optional;
/**
* @author Andrew Shvayka
@@ -39,14 +41,21 @@ public class BaseAttributesService implements AttributesService {
private AttributesDao attributesDao;
@Override
- public AttributeKvEntry find(EntityId entityId, String scope, String attributeKey) {
+ public ListenableFuture> find(EntityId entityId, String scope, String attributeKey) {
validate(entityId, scope);
Validator.validateString(attributeKey, "Incorrect attribute key " + attributeKey);
return attributesDao.find(entityId, scope, attributeKey);
}
@Override
- public List findAll(EntityId entityId, String scope) {
+ public ListenableFuture> find(EntityId entityId, String scope, Collection attributeKeys) {
+ validate(entityId, scope);
+ attributeKeys.forEach(attributeKey -> Validator.validateString(attributeKey, "Incorrect attribute key " + attributeKey));
+ return attributesDao.find(entityId, scope, attributeKeys);
+ }
+
+ @Override
+ public ListenableFuture> findAll(EntityId entityId, String scope) {
validate(entityId, scope);
return attributesDao.findAll(entityId, scope);
}
@@ -56,16 +65,16 @@ public class BaseAttributesService implements AttributesService {
validate(entityId, scope);
attributes.forEach(attribute -> validate(attribute));
List futures = Lists.newArrayListWithExpectedSize(attributes.size());
- for(AttributeKvEntry attribute : attributes) {
+ for (AttributeKvEntry attribute : attributes) {
futures.add(attributesDao.save(entityId, scope, attribute));
}
return Futures.allAsList(futures);
}
@Override
- public void removeAll(EntityId entityId, String scope, List keys) {
+ public ListenableFuture> removeAll(EntityId entityId, String scope, List keys) {
validate(entityId, scope);
- attributesDao.removeAll(entityId, scope, keys);
+ return attributesDao.removeAll(entityId, scope, keys);
}
private static void validate(EntityId id, String scope) {
diff --git a/dao/src/main/java/org/thingsboard/server/dao/device/DeviceDao.java b/dao/src/main/java/org/thingsboard/server/dao/device/DeviceDao.java
index 58034161af..b8d395c8e2 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/device/DeviceDao.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/device/DeviceDao.java
@@ -19,6 +19,7 @@ import java.util.List;
import java.util.Optional;
import java.util.UUID;
+import com.google.common.util.concurrent.ListenableFuture;
import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.page.TextPageLink;
import org.thingsboard.server.dao.Dao;
@@ -46,7 +47,16 @@ public interface DeviceDao extends Dao {
* @return the list of device objects
*/
List findDevicesByTenantId(UUID tenantId, TextPageLink pageLink);
-
+
+ /**
+ * Find devices by tenantId and devices Ids.
+ *
+ * @param tenantId the tenantId
+ * @param deviceIds the device Ids
+ * @return the list of device objects
+ */
+ ListenableFuture> findDevicesByTenantIdAndIdsAsync(UUID tenantId, List deviceIds);
+
/**
* Find devices by tenantId, customerId and page link.
*
@@ -57,6 +67,16 @@ public interface DeviceDao extends Dao {
*/
List findDevicesByTenantIdAndCustomerId(UUID tenantId, UUID customerId, TextPageLink pageLink);
+ /**
+ * Find devices by tenantId, customerId and devices Ids.
+ *
+ * @param tenantId the tenantId
+ * @param customerId the customerId
+ * @param deviceIds the device Ids
+ * @return the list of device objects
+ */
+ ListenableFuture> findDevicesByTenantIdCustomerIdAndIdsAsync(UUID tenantId, UUID customerId, List deviceIds);
+
/**
* Find devices by tenantId and device name.
*
diff --git a/dao/src/main/java/org/thingsboard/server/dao/device/DeviceDaoImpl.java b/dao/src/main/java/org/thingsboard/server/dao/device/DeviceDaoImpl.java
index 540204d057..81fc0bcde1 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/device/DeviceDaoImpl.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/device/DeviceDaoImpl.java
@@ -16,12 +16,14 @@
package org.thingsboard.server.dao.device;
import static com.datastax.driver.core.querybuilder.QueryBuilder.eq;
+import static com.datastax.driver.core.querybuilder.QueryBuilder.in;
import static com.datastax.driver.core.querybuilder.QueryBuilder.select;
import static org.thingsboard.server.dao.model.ModelConstants.*;
import java.util.*;
import com.datastax.driver.core.querybuilder.Select;
+import com.google.common.util.concurrent.ListenableFuture;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Component;
import org.thingsboard.server.common.data.Device;
@@ -61,6 +63,16 @@ public class DeviceDaoImpl extends AbstractSearchTextDao implement
return deviceEntities;
}
+ @Override
+ public ListenableFuture> findDevicesByTenantIdAndIdsAsync(UUID tenantId, List deviceIds) {
+ log.debug("Try to find devices by tenantId [{}] and device Ids [{}]", tenantId, deviceIds);
+ Select select = select().from(getColumnFamilyName());
+ Select.Where query = select.where();
+ query.and(eq(DEVICE_TENANT_ID_PROPERTY, tenantId));
+ query.and(in(ID_PROPERTY, deviceIds));
+ return findListByStatementAsync(query);
+ }
+
@Override
public List findDevicesByTenantIdAndCustomerId(UUID tenantId, UUID customerId, TextPageLink pageLink) {
log.debug("Try to find devices by tenantId [{}], customerId[{}] and pageLink [{}]", tenantId, customerId, pageLink);
@@ -73,6 +85,17 @@ public class DeviceDaoImpl extends AbstractSearchTextDao implement
return deviceEntities;
}
+ @Override
+ public ListenableFuture> findDevicesByTenantIdCustomerIdAndIdsAsync(UUID tenantId, UUID customerId, List deviceIds) {
+ log.debug("Try to find devices by tenantId [{}], customerId [{}] and device Ids [{}]", tenantId, customerId, deviceIds);
+ Select select = select().from(getColumnFamilyName());
+ Select.Where query = select.where();
+ query.and(eq(DEVICE_TENANT_ID_PROPERTY, tenantId));
+ query.and(eq(DEVICE_CUSTOMER_ID_PROPERTY, customerId));
+ query.and(in(ID_PROPERTY, deviceIds));
+ return findListByStatementAsync(query);
+ }
+
@Override
public Optional findDevicesByTenantIdAndName(UUID tenantId, String deviceName) {
Select select = select().from(DEVICE_BY_TENANT_AND_NAME_VIEW_NAME);
diff --git a/dao/src/main/java/org/thingsboard/server/dao/device/DeviceService.java b/dao/src/main/java/org/thingsboard/server/dao/device/DeviceService.java
index 8d780b6aa1..35d34968cf 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/device/DeviceService.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/device/DeviceService.java
@@ -15,6 +15,7 @@
*/
package org.thingsboard.server.dao.device;
+import com.google.common.util.concurrent.ListenableFuture;
import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.id.CustomerId;
import org.thingsboard.server.common.data.id.DeviceId;
@@ -22,12 +23,15 @@ import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.page.TextPageData;
import org.thingsboard.server.common.data.page.TextPageLink;
+import java.util.List;
import java.util.Optional;
public interface DeviceService {
Device findDeviceById(DeviceId deviceId);
+ ListenableFuture findDeviceByIdAsync(DeviceId deviceId);
+
Optional findDeviceByTenantIdAndName(TenantId tenantId, String name);
Device saveDevice(Device device);
@@ -40,9 +44,13 @@ public interface DeviceService {
TextPageData findDevicesByTenantId(TenantId tenantId, TextPageLink pageLink);
+ ListenableFuture> findDevicesByTenantIdAndIdsAsync(TenantId tenantId, List deviceIds);
+
void deleteDevicesByTenantId(TenantId tenantId);
TextPageData findDevicesByTenantIdAndCustomerId(TenantId tenantId, CustomerId customerId, TextPageLink pageLink);
+ ListenableFuture> findDevicesByTenantIdCustomerIdAndIdsAsync(TenantId tenantId, CustomerId customerId, List deviceIds);
+
void unassignCustomerDevices(TenantId tenantId, CustomerId customerId);
}
diff --git a/dao/src/main/java/org/thingsboard/server/dao/device/DeviceServiceImpl.java b/dao/src/main/java/org/thingsboard/server/dao/device/DeviceServiceImpl.java
index 681188e6be..3d1ce31349 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/device/DeviceServiceImpl.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/device/DeviceServiceImpl.java
@@ -15,6 +15,9 @@
*/
package org.thingsboard.server.dao.device;
+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.apache.commons.lang3.RandomStringUtils;
import org.springframework.beans.factory.annotation.Autowired;
@@ -42,8 +45,10 @@ import java.util.Optional;
import static org.thingsboard.server.dao.DaoUtil.convertDataList;
import static org.thingsboard.server.dao.DaoUtil.getData;
+import static org.thingsboard.server.dao.DaoUtil.toUUIDs;
import static org.thingsboard.server.dao.model.ModelConstants.NULL_UUID;
import static org.thingsboard.server.dao.service.Validator.validateId;
+import static org.thingsboard.server.dao.service.Validator.validateIds;
import static org.thingsboard.server.dao.service.Validator.validatePageLink;
@Service
@@ -70,6 +75,14 @@ public class DeviceServiceImpl implements DeviceService {
return getData(deviceEntity);
}
+ @Override
+ public ListenableFuture findDeviceByIdAsync(DeviceId deviceId) {
+ log.trace("Executing findDeviceById [{}]", deviceId);
+ validateId(deviceId, "Incorrect deviceId " + deviceId);
+ ListenableFuture deviceEntity = deviceDao.findByIdAsync(deviceId.getId());
+ return Futures.transform(deviceEntity, (Function super DeviceEntity, ? extends Device>) input -> getData(input));
+ }
+
@Override
public Optional findDeviceByTenantIdAndName(TenantId tenantId, String name) {
log.trace("Executing findDeviceByTenantIdAndName [{}][{}]", tenantId, name);
@@ -132,6 +145,16 @@ public class DeviceServiceImpl implements DeviceService {
return new TextPageData(devices, pageLink);
}
+ @Override
+ public ListenableFuture> findDevicesByTenantIdAndIdsAsync(TenantId tenantId, List deviceIds) {
+ log.trace("Executing findDevicesByTenantIdAndIdsAsync, tenantId [{}], deviceIds [{}]", tenantId, deviceIds);
+ validateId(tenantId, "Incorrect tenantId " + tenantId);
+ validateIds(deviceIds, "Incorrect deviceIds " + deviceIds);
+ ListenableFuture> deviceEntities = deviceDao.findDevicesByTenantIdAndIdsAsync(tenantId.getId(), toUUIDs(deviceIds));
+ return Futures.transform(deviceEntities, (Function, List>) input -> convertDataList(input));
+ }
+
+
@Override
public void deleteDevicesByTenantId(TenantId tenantId) {
log.trace("Executing deleteDevicesByTenantId, tenantId [{}]", tenantId);
@@ -150,6 +173,17 @@ public class DeviceServiceImpl implements DeviceService {
return new TextPageData(devices, pageLink);
}
+ @Override
+ public ListenableFuture> findDevicesByTenantIdCustomerIdAndIdsAsync(TenantId tenantId, CustomerId customerId, List deviceIds) {
+ log.trace("Executing findDevicesByTenantIdCustomerIdAndIdsAsync, tenantId [{}], customerId [{}], deviceIds [{}]", tenantId, customerId, deviceIds);
+ validateId(tenantId, "Incorrect tenantId " + tenantId);
+ validateId(customerId, "Incorrect customerId " + customerId);
+ validateIds(deviceIds, "Incorrect deviceIds " + deviceIds);
+ ListenableFuture> deviceEntities = deviceDao.findDevicesByTenantIdCustomerIdAndIdsAsync(tenantId.getId(),
+ customerId.getId(), toUUIDs(deviceIds));
+ return Futures.transform(deviceEntities, (Function, List>) input -> convertDataList(input));
+ }
+
@Override
public void unassignCustomerDevices(TenantId tenantId, CustomerId customerId) {
log.trace("Executing unassignCustomerDevices, tenantId [{}], customerId [{}]", tenantId, customerId);
diff --git a/dao/src/main/java/org/thingsboard/server/dao/model/ModelConstants.java b/dao/src/main/java/org/thingsboard/server/dao/model/ModelConstants.java
index b68fb75b5a..d3ed5d1dd1 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/model/ModelConstants.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/model/ModelConstants.java
@@ -18,14 +18,16 @@ package org.thingsboard.server.dao.model;
import java.util.UUID;
import com.datastax.driver.core.utils.UUIDs;
+import org.apache.commons.lang3.ArrayUtils;
+import org.thingsboard.server.common.data.kv.Aggregation;
public class ModelConstants {
private ModelConstants() {
}
-
+
public static UUID NULL_UUID = UUIDs.startOf(0);
-
+
/**
* Generic constants.
*/
@@ -38,7 +40,7 @@ public class ModelConstants {
public static final String ALIAS_PROPERTY = "alias";
public static final String SEARCH_TEXT_PROPERTY = "search_text";
public static final String ADDITIONAL_INFO_PROPERTY = "additional_info";
-
+
/**
* Cassandra user constants.
*/
@@ -50,11 +52,11 @@ public class ModelConstants {
public static final String USER_FIRST_NAME_PROPERTY = "first_name";
public static final String USER_LAST_NAME_PROPERTY = "last_name";
public static final String USER_ADDITIONAL_INFO_PROPERTY = ADDITIONAL_INFO_PROPERTY;
-
+
public static final String USER_BY_EMAIL_COLUMN_FAMILY_NAME = "user_by_email";
public static final String USER_BY_TENANT_AND_SEARCH_TEXT_COLUMN_FAMILY_NAME = "user_by_tenant_and_search_text";
public static final String USER_BY_CUSTOMER_AND_SEARCH_TEXT_COLUMN_FAMILY_NAME = "user_by_customer_and_search_text";
-
+
/**
* Cassandra user_credentials constants.
*/
@@ -64,20 +66,20 @@ public class ModelConstants {
public static final String USER_CREDENTIALS_PASSWORD_PROPERTY = "password";
public static final String USER_CREDENTIALS_ACTIVATE_TOKEN_PROPERTY = "activate_token";
public static final String USER_CREDENTIALS_RESET_TOKEN_PROPERTY = "reset_token";
-
+
public static final String USER_CREDENTIALS_BY_USER_COLUMN_FAMILY_NAME = "user_credentials_by_user";
public static final String USER_CREDENTIALS_BY_ACTIVATE_TOKEN_COLUMN_FAMILY_NAME = "user_credentials_by_activate_token";
public static final String USER_CREDENTIALS_BY_RESET_TOKEN_COLUMN_FAMILY_NAME = "user_credentials_by_reset_token";
-
+
/**
* Cassandra admin_settings constants.
*/
public static final String ADMIN_SETTINGS_COLUMN_FAMILY_NAME = "admin_settings";
public static final String ADMIN_SETTINGS_KEY_PROPERTY = "key";
public static final String ADMIN_SETTINGS_JSON_VALUE_PROPERTY = "json_value";
-
+
public static final String ADMIN_SETTINGS_BY_KEY_COLUMN_FAMILY_NAME = "admin_settings_by_key";
-
+
/**
* Cassandra contact constants.
*/
@@ -97,9 +99,9 @@ public class ModelConstants {
public static final String TENANT_TITLE_PROPERTY = TITLE_PROPERTY;
public static final String TENANT_REGION_PROPERTY = "region";
public static final String TENANT_ADDITIONAL_INFO_PROPERTY = ADDITIONAL_INFO_PROPERTY;
-
+
public static final String TENANT_BY_REGION_AND_SEARCH_TEXT_COLUMN_FAMILY_NAME = "tenant_by_region_and_search_text";
-
+
/**
* Cassandra customer constants.
*/
@@ -107,9 +109,9 @@ public class ModelConstants {
public static final String CUSTOMER_TENANT_ID_PROPERTY = TENTANT_ID_PROPERTY;
public static final String CUSTOMER_TITLE_PROPERTY = TITLE_PROPERTY;
public static final String CUSTOMER_ADDITIONAL_INFO_PROPERTY = ADDITIONAL_INFO_PROPERTY;
-
+
public static final String CUSTOMER_BY_TENANT_AND_SEARCH_TEXT_COLUMN_FAMILY_NAME = "customer_by_tenant_and_search_text";
-
+
/**
* Cassandra device constants.
*/
@@ -118,12 +120,12 @@ public class ModelConstants {
public static final String DEVICE_CUSTOMER_ID_PROPERTY = CUSTOMER_ID_PROPERTY;
public static final String DEVICE_NAME_PROPERTY = "name";
public static final String DEVICE_ADDITIONAL_INFO_PROPERTY = ADDITIONAL_INFO_PROPERTY;
-
+
public static final String DEVICE_BY_TENANT_AND_SEARCH_TEXT_COLUMN_FAMILY_NAME = "device_by_tenant_and_search_text";
public static final String DEVICE_BY_CUSTOMER_AND_SEARCH_TEXT_COLUMN_FAMILY_NAME = "device_by_customer_and_search_text";
public static final String DEVICE_BY_TENANT_AND_NAME_VIEW_NAME = "device_by_tenant_and_name";
-
+
/**
* Cassandra device_credentials constants.
*/
@@ -132,7 +134,7 @@ public class ModelConstants {
public static final String DEVICE_CREDENTIALS_CREDENTIALS_TYPE_PROPERTY = "credentials_type";
public static final String DEVICE_CREDENTIALS_CREDENTIALS_ID_PROPERTY = "credentials_id";
public static final String DEVICE_CREDENTIALS_CREDENTIALS_VALUE_PROPERTY = "credentials_value";
-
+
public static final String DEVICE_CREDENTIALS_BY_DEVICE_COLUMN_FAMILY_NAME = "device_credentials_by_device";
public static final String DEVICE_CREDENTIALS_BY_CREDENTIALS_ID_COLUMN_FAMILY_NAME = "device_credentials_by_credentials_id";
@@ -203,9 +205,9 @@ public class ModelConstants {
public static final String COMPONENT_DESCRIPTOR_BY_SCOPE_TYPE_AND_SEARCH_TEXT_COLUMN_FAMILY_NAME = "component_desc_by_scope_type_search_text";
public static final String COMPONENT_DESCRIPTOR_BY_ID = "component_desc_by_id";
- /**
- * Cassandra rule metadata constants.
- */
+ /**
+ * Cassandra rule metadata constants.
+ */
public static final String RULE_COLUMN_FAMILY_NAME = "rule";
public static final String RULE_TENANT_ID_PROPERTY = TENTANT_ID_PROPERTY;
public static final String RULE_NAME_PROPERTY = "name";
@@ -259,4 +261,51 @@ public class ModelConstants {
public static final String STRING_VALUE_COLUMN = "str_v";
public static final String LONG_VALUE_COLUMN = "long_v";
public static final String DOUBLE_VALUE_COLUMN = "dbl_v";
+
+ public static final String[] NONE_AGGREGATION_COLUMNS = new String[]{LONG_VALUE_COLUMN, DOUBLE_VALUE_COLUMN, BOOLEAN_VALUE_COLUMN, STRING_VALUE_COLUMN, KEY_COLUMN, TS_COLUMN};
+
+ public static final String[] COUNT_AGGREGATION_COLUMNS = new String[]{count(LONG_VALUE_COLUMN), count(DOUBLE_VALUE_COLUMN), count(BOOLEAN_VALUE_COLUMN), count(STRING_VALUE_COLUMN)};
+
+ public static final String[] MIN_AGGREGATION_COLUMNS = ArrayUtils.addAll(COUNT_AGGREGATION_COLUMNS,
+ new String[]{min(LONG_VALUE_COLUMN), min(DOUBLE_VALUE_COLUMN), min(BOOLEAN_VALUE_COLUMN), min(STRING_VALUE_COLUMN)});
+ public static final String[] MAX_AGGREGATION_COLUMNS = ArrayUtils.addAll(COUNT_AGGREGATION_COLUMNS,
+ new String[]{max(LONG_VALUE_COLUMN), max(DOUBLE_VALUE_COLUMN), max(BOOLEAN_VALUE_COLUMN), max(STRING_VALUE_COLUMN)});
+ public static final String[] SUM_AGGREGATION_COLUMNS = ArrayUtils.addAll(COUNT_AGGREGATION_COLUMNS,
+ new String[]{sum(LONG_VALUE_COLUMN), sum(DOUBLE_VALUE_COLUMN)});
+ public static final String[] AVG_AGGREGATION_COLUMNS = SUM_AGGREGATION_COLUMNS;
+
+ public static String min(String s) {
+ return "min(" + s + ")";
+ }
+
+ public static String max(String s) {
+ return "max(" + s + ")";
+ }
+
+ public static String sum(String s) {
+ return "sum(" + s + ")";
+ }
+
+ public static String count(String s) {
+ return "count(" + s + ")";
+ }
+
+ public static String[] getFetchColumnNames(Aggregation aggregation) {
+ switch (aggregation) {
+ case NONE:
+ return NONE_AGGREGATION_COLUMNS;
+ case MIN:
+ return MIN_AGGREGATION_COLUMNS;
+ case MAX:
+ return MAX_AGGREGATION_COLUMNS;
+ case SUM:
+ return SUM_AGGREGATION_COLUMNS;
+ case COUNT:
+ return COUNT_AGGREGATION_COLUMNS;
+ case AVG:
+ return AVG_AGGREGATION_COLUMNS;
+ default:
+ throw new RuntimeException("Aggregation type: " + aggregation + " is not supported!");
+ }
+ }
}
diff --git a/dao/src/main/java/org/thingsboard/server/dao/service/Validator.java b/dao/src/main/java/org/thingsboard/server/dao/service/Validator.java
index 1976eb9b33..70e9860095 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/service/Validator.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/service/Validator.java
@@ -19,6 +19,7 @@ import org.thingsboard.server.common.data.id.UUIDBased;
import org.thingsboard.server.common.data.page.TextPageLink;
import org.thingsboard.server.dao.exception.IncorrectParameterException;
+import java.util.List;
import java.util.UUID;
public class Validator {
@@ -77,6 +78,23 @@ public class Validator {
}
}
+ /**
+ * This method validate list of UUIDBased ids. If at least one of the ids is null than throw
+ * IncorrectParameterException exception
+ *
+ * @param ids the list of ids
+ * @param errorMessage the error message for exception
+ */
+ public static void validateIds(List extends UUIDBased> ids, String errorMessage) {
+ if (ids == null || ids.isEmpty()) {
+ throw new IncorrectParameterException(errorMessage);
+ } else {
+ for (UUIDBased id : ids) {
+ validateId(id, errorMessage);
+ }
+ }
+ }
+
/**
* This method validate PageLink page link. If pageLink is invalid than throw
* IncorrectParameterException exception
diff --git a/dao/src/main/java/org/thingsboard/server/dao/timeseries/AggregatePartitionsFunction.java b/dao/src/main/java/org/thingsboard/server/dao/timeseries/AggregatePartitionsFunction.java
new file mode 100644
index 0000000000..f099eec004
--- /dev/null
+++ b/dao/src/main/java/org/thingsboard/server/dao/timeseries/AggregatePartitionsFunction.java
@@ -0,0 +1,193 @@
+/**
+ * Copyright © 2016-2017 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.driver.core.ResultSet;
+import com.datastax.driver.core.Row;
+import org.thingsboard.server.common.data.kv.*;
+
+import javax.annotation.Nullable;
+import java.util.List;
+import java.util.Optional;
+
+/**
+ * Created by ashvayka on 20.02.17.
+ */
+public class AggregatePartitionsFunction implements com.google.common.base.Function, Optional> {
+
+ private static final int LONG_CNT_POS = 0;
+ private static final int DOUBLE_CNT_POS = 1;
+ private static final int BOOL_CNT_POS = 2;
+ private static final int STR_CNT_POS = 3;
+ private static final int LONG_POS = 4;
+ private static final int DOUBLE_POS = 5;
+ private static final int BOOL_POS = 6;
+ private static final int STR_POS = 7;
+
+ private final Aggregation aggregation;
+ private final String key;
+ private final long ts;
+
+ public AggregatePartitionsFunction(Aggregation aggregation, String key, long ts) {
+ this.aggregation = aggregation;
+ this.key = key;
+ this.ts = ts;
+ }
+
+ @Nullable
+ @Override
+ public Optional apply(@Nullable List rsList) {
+ if (rsList == null || rsList.size() == 0) {
+ return Optional.empty();
+ }
+ long count = 0;
+ DataType dataType = null;
+
+ Boolean bValue = null;
+ String sValue = null;
+ Double dValue = null;
+ Long lValue = null;
+
+ for (ResultSet rs : rsList) {
+ for (Row row : rs.all()) {
+ long curCount;
+
+ Long curLValue = null;
+ Double curDValue = null;
+ Boolean curBValue = null;
+ String curSValue = null;
+
+ long longCount = row.getLong(LONG_CNT_POS);
+ long doubleCount = row.getLong(DOUBLE_CNT_POS);
+ long boolCount = row.getLong(BOOL_CNT_POS);
+ long strCount = row.getLong(STR_CNT_POS);
+
+ if (longCount > 0) {
+ dataType = DataType.LONG;
+ curCount = longCount;
+ curLValue = getLongValue(row);
+ } else if (doubleCount > 0) {
+ dataType = DataType.DOUBLE;
+ curCount = doubleCount;
+ curDValue = getDoubleValue(row);
+ } else if (boolCount > 0) {
+ dataType = DataType.BOOLEAN;
+ curCount = boolCount;
+ curBValue = getBooleanValue(row);
+ } else if (strCount > 0) {
+ dataType = DataType.STRING;
+ curCount = strCount;
+ curSValue = getStringValue(row);
+ } else {
+ continue;
+ }
+
+ if (aggregation == Aggregation.COUNT) {
+ count += curCount;
+ } else if (aggregation == Aggregation.AVG || aggregation == Aggregation.SUM) {
+ count += curCount;
+ if (curDValue != null) {
+ dValue = dValue == null ? curDValue : dValue + curDValue;
+ } else if (curLValue != null) {
+ lValue = lValue == null ? curLValue : lValue + curLValue;
+ }
+ } else if (aggregation == Aggregation.MIN) {
+ if (curDValue != null) {
+ dValue = dValue == null ? curDValue : Math.min(dValue, curDValue);
+ } else if (curLValue != null) {
+ lValue = lValue == null ? curLValue : Math.min(lValue, curLValue);
+ } else if (curBValue != null) {
+ bValue = bValue == null ? curBValue : bValue && curBValue;
+ } else if (curSValue != null) {
+ if (sValue == null || curSValue.compareTo(sValue) < 0) {
+ sValue = curSValue;
+ }
+ }
+ } else if (aggregation == Aggregation.MAX) {
+ if (curDValue != null) {
+ dValue = dValue == null ? curDValue : Math.max(dValue, curDValue);
+ } else if (curLValue != null) {
+ lValue = lValue == null ? curLValue : Math.max(lValue, curLValue);
+ } else if (curBValue != null) {
+ bValue = bValue == null ? curBValue : bValue || curBValue;
+ } else if (curSValue != null) {
+ if (sValue == null || curSValue.compareTo(sValue) > 0) {
+ sValue = curSValue;
+ }
+ }
+ }
+ }
+ }
+ if (dataType == null) {
+ return Optional.empty();
+ } else if (aggregation == Aggregation.COUNT) {
+ return Optional.of(new BasicTsKvEntry(ts, new LongDataEntry(key, (long) count)));
+ } else if (aggregation == Aggregation.AVG || aggregation == Aggregation.SUM) {
+ if (count == 0 || (dataType == DataType.DOUBLE && dValue == null) || (dataType == DataType.LONG && lValue == null)) {
+ return Optional.empty();
+ } else if (dataType == DataType.DOUBLE) {
+ return Optional.of(new BasicTsKvEntry(ts, new DoubleDataEntry(key, aggregation == Aggregation.SUM ? dValue : (dValue / count))));
+ } else if (dataType == DataType.LONG) {
+ return Optional.of(new BasicTsKvEntry(ts, new LongDataEntry(key, aggregation == Aggregation.SUM ? lValue : (lValue / count))));
+ }
+ } else if (aggregation == Aggregation.MIN || aggregation == Aggregation.MAX) {
+ if (dataType == DataType.DOUBLE) {
+ return Optional.of(new BasicTsKvEntry(ts, new DoubleDataEntry(key, dValue)));
+ } else if (dataType == DataType.LONG) {
+ return Optional.of(new BasicTsKvEntry(ts, new LongDataEntry(key, lValue)));
+ } else if (dataType == DataType.STRING) {
+ return Optional.of(new BasicTsKvEntry(ts, new StringDataEntry(key, sValue)));
+ } else {
+ return Optional.of(new BasicTsKvEntry(ts, new BooleanDataEntry(key, bValue)));
+ }
+ }
+ return null;
+ }
+
+ private Boolean getBooleanValue(Row row) {
+ if (aggregation == Aggregation.MIN || aggregation == Aggregation.MAX) {
+ return row.getBool(BOOL_POS);
+ } else {
+ return null;
+ }
+ }
+
+ private String getStringValue(Row row) {
+ if (aggregation == Aggregation.MIN || aggregation == Aggregation.MAX) {
+ return row.getString(STR_POS);
+ } else {
+ return null;
+ }
+ }
+
+ private Long getLongValue(Row row) {
+ if (aggregation == Aggregation.MIN || aggregation == Aggregation.MAX
+ || aggregation == Aggregation.SUM || aggregation == Aggregation.AVG) {
+ return row.getLong(LONG_POS);
+ } else {
+ return null;
+ }
+ }
+
+ private Double getDoubleValue(Row row) {
+ if (aggregation == Aggregation.MIN || aggregation == Aggregation.MAX
+ || aggregation == Aggregation.SUM || aggregation == Aggregation.AVG) {
+ return row.getDouble(DOUBLE_POS);
+ } else {
+ return null;
+ }
+ }
+}
diff --git a/dao/src/main/java/org/thingsboard/server/dao/timeseries/BaseTimeseriesDao.java b/dao/src/main/java/org/thingsboard/server/dao/timeseries/BaseTimeseriesDao.java
index 09c415c0eb..c7584fe511 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/timeseries/BaseTimeseriesDao.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/timeseries/BaseTimeseriesDao.java
@@ -18,15 +18,30 @@ package org.thingsboard.server.dao.timeseries;
import com.datastax.driver.core.*;
import com.datastax.driver.core.querybuilder.QueryBuilder;
import com.datastax.driver.core.querybuilder.Select;
+import com.google.common.base.Function;
+import com.google.common.util.concurrent.AsyncFunction;
+import com.google.common.util.concurrent.FutureCallback;
+import com.google.common.util.concurrent.Futures;
+import com.google.common.util.concurrent.ListenableFuture;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.stereotype.Component;
import org.thingsboard.server.common.data.kv.*;
import org.thingsboard.server.common.data.kv.DataType;
+import org.thingsboard.server.dao.AbstractAsyncDao;
import org.thingsboard.server.dao.AbstractDao;
import org.thingsboard.server.dao.model.ModelConstants;
+import javax.annotation.Nullable;
+import javax.annotation.PostConstruct;
+import javax.annotation.PreDestroy;
+import java.time.Instant;
+import java.time.LocalDateTime;
+import java.time.ZoneOffset;
import java.util.*;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.stream.Collectors;
import static com.datastax.driver.core.querybuilder.QueryBuilder.eq;
import static com.datastax.driver.core.querybuilder.QueryBuilder.select;
@@ -36,53 +51,186 @@ import static com.datastax.driver.core.querybuilder.QueryBuilder.select;
*/
@Component
@Slf4j
-public class BaseTimeseriesDao extends AbstractDao implements TimeseriesDao {
+public class BaseTimeseriesDao extends AbstractAsyncDao implements TimeseriesDao {
- @Value("${cassandra.query.max_limit_per_request}")
- protected Integer maxLimitPerRequest;
+ //@Value("${cassandra.query.min_aggregation_step_ms}")
+ //TODO:
+ private int minAggregationStepMs = 1000;
+
+ @Value("${cassandra.query.ts_key_value_partitioning}")
+ private String partitioning;
+
+ private TsPartitionDate tsFormat;
private PreparedStatement partitionInsertStmt;
private PreparedStatement[] latestInsertStmts;
private PreparedStatement[] saveStmts;
+ private PreparedStatement[] fetchStmts;
private PreparedStatement findLatestStmt;
private PreparedStatement findAllLatestStmt;
+ @PostConstruct
+ public void init() {
+ super.startExecutor();
+ getFetchStmt(Aggregation.NONE);
+ Optional partition = TsPartitionDate.parse(partitioning);
+ if (partition.isPresent()) {
+ tsFormat = partition.get();
+ } else {
+ log.warn("Incorrect configuration of partitioning {}", partitioning);
+ throw new RuntimeException("Failed to parse partitioning property: " + partitioning + "!");
+ }
+ }
+
+ @PreDestroy
+ public void stop() {
+ super.stopExecutor();
+ }
+
+ @Override
+ public long toPartitionTs(long ts) {
+ LocalDateTime time = LocalDateTime.ofInstant(Instant.ofEpochMilli(ts), ZoneOffset.UTC);
+ return tsFormat.truncatedTo(time).toInstant(ZoneOffset.UTC).toEpochMilli();
+ }
+
@Override
- public List find(String entityType, UUID entityId, TsKvQuery query, Optional minPartition, Optional maxPartition) {
- List rows = Collections.emptyList();
- Long[] parts = fetchPartitions(entityType, entityId, query.getKey(), minPartition, maxPartition);
- int partsLength = parts.length;
- if (parts != null && partsLength > 0) {
- int limit = maxLimitPerRequest;
- Optional lim = query.getLimit();
- if (lim.isPresent() && lim.get() < maxLimitPerRequest) {
- limit = lim.get();
+ public ListenableFuture> findAllAsync(String entityType, UUID entityId, List queries) {
+ List>> futures = queries.stream().map(query -> findAllAsync(entityType, entityId, query)).collect(Collectors.toList());
+ return Futures.transform(Futures.allAsList(futures), new Function>, List>() {
+ @Nullable
+ @Override
+ public List apply(@Nullable List> results) {
+ List result = new ArrayList();
+ results.forEach(r -> result.addAll(r));
+ return result;
}
+ }, readResultsProcessingExecutor);
+ }
- rows = new ArrayList<>(limit);
- int lastIdx = partsLength - 1;
- for (int i = 0; i < partsLength; i++) {
- int currentLimit;
- if (rows.size() >= limit) {
- break;
- } else {
- currentLimit = limit - rows.size();
- }
- Long partition = parts[i];
- Select.Where where = select().from(ModelConstants.TS_KV_CF).where(eq(ModelConstants.ENTITY_TYPE_COLUMN, entityType))
- .and(eq(ModelConstants.ENTITY_ID_COLUMN, entityId))
- .and(eq(ModelConstants.KEY_COLUMN, query.getKey()))
- .and(eq(ModelConstants.PARTITION_COLUMN, partition));
- if (i == 0 && query.getStartTs().isPresent()) {
- where.and(QueryBuilder.gt(ModelConstants.TS_COLUMN, query.getStartTs().get()));
- } else if (i == lastIdx && query.getEndTs().isPresent()) {
- where.and(QueryBuilder.lte(ModelConstants.TS_COLUMN, query.getEndTs().get()));
+
+ private ListenableFuture> findAllAsync(String entityType, UUID entityId, TsKvQuery query) {
+ if (query.getAggregation() == Aggregation.NONE) {
+ return findAllAsyncWithLimit(entityType, entityId, query);
+ } else {
+ long step = Math.max(query.getInterval(), minAggregationStepMs);
+ long stepTs = query.getStartTs();
+ List>> futures = new ArrayList<>();
+ while (stepTs < query.getEndTs()) {
+ long startTs = stepTs;
+ long endTs = stepTs + step;
+ TsKvQuery subQuery = new BaseTsKvQuery(query.getKey(), startTs, endTs, step, 1, query.getAggregation());
+ futures.add(findAndAggregateAsync(entityType, entityId, subQuery, toPartitionTs(startTs), toPartitionTs(endTs)));
+ stepTs = endTs;
+ }
+ ListenableFuture>> future = Futures.allAsList(futures);
+ return Futures.transform(future, new Function>, List>() {
+ @Nullable
+ @Override
+ public List apply(@Nullable List> input) {
+ return input.stream().filter(v -> v.isPresent()).map(v -> v.get()).collect(Collectors.toList());
}
- where.limit(currentLimit);
- rows.addAll(executeRead(where).all());
+ }, readResultsProcessingExecutor);
+ }
+ }
+
+ private ListenableFuture> findAllAsyncWithLimit(String entityType, UUID entityId, TsKvQuery query) {
+ long minPartition = toPartitionTs(query.getStartTs());
+ long maxPartition = toPartitionTs(query.getEndTs());
+
+ ResultSetFuture partitionsFuture = fetchPartitions(entityType, entityId, query.getKey(), minPartition, maxPartition);
+
+ final SimpleListenableFuture> resultFuture = new SimpleListenableFuture<>();
+ final ListenableFuture> partitionsListFuture = Futures.transform(partitionsFuture, getPartitionsArrayFunction(), readResultsProcessingExecutor);
+
+ Futures.addCallback(partitionsListFuture, new FutureCallback>() {
+ @Override
+ public void onSuccess(@Nullable List partitions) {
+ TsKvQueryCursor cursor = new TsKvQueryCursor(entityType, entityId, query, partitions);
+ findAllAsyncSequentiallyWithLimit(cursor, resultFuture);
}
+
+ @Override
+ public void onFailure(Throwable t) {
+ log.error("[{}][{}] Failed to fetch partitions for interval {}-{}", entityType, entityId, minPartition, maxPartition, t);
+ }
+ }, readResultsProcessingExecutor);
+
+ return resultFuture;
+ }
+
+ private void findAllAsyncSequentiallyWithLimit(final TsKvQueryCursor cursor, final SimpleListenableFuture> resultFuture) {
+ if (cursor.isFull() || !cursor.hasNextPartition()) {
+ resultFuture.set(cursor.getData());
+ } else {
+ PreparedStatement proto = getFetchStmt(Aggregation.NONE);
+ BoundStatement stmt = proto.bind();
+ stmt.setString(0, cursor.getEntityType());
+ stmt.setUUID(1, cursor.getEntityId());
+ stmt.setString(2, cursor.getKey());
+ stmt.setLong(3, cursor.getNextPartition());
+ stmt.setLong(4, cursor.getStartTs());
+ stmt.setLong(5, cursor.getEndTs());
+ stmt.setInt(6, cursor.getCurrentLimit());
+
+ Futures.addCallback(executeAsyncRead(stmt), new FutureCallback() {
+ @Override
+ public void onSuccess(@Nullable ResultSet result) {
+ cursor.addData(convertResultToTsKvEntryList(result.all()));
+ findAllAsyncSequentiallyWithLimit(cursor, resultFuture);
+ }
+
+ @Override
+ public void onFailure(Throwable t) {
+ log.error("[{}][{}] Failed to fetch data for query {}-{}", stmt, t);
+ }
+ }, readResultsProcessingExecutor);
}
- return convertResultToTsKvEntryList(rows);
+ }
+
+ private ListenableFuture> findAndAggregateAsync(String entityType, UUID entityId, TsKvQuery query, long minPartition, long maxPartition) {
+ final Aggregation aggregation = query.getAggregation();
+ final String key = query.getKey();
+ final long startTs = query.getStartTs();
+ final long endTs = query.getEndTs();
+ final long ts = startTs + (endTs - startTs) / 2;
+
+ ResultSetFuture partitionsFuture = fetchPartitions(entityType, entityId, key, minPartition, maxPartition);
+
+ ListenableFuture> partitionsListFuture = Futures.transform(partitionsFuture, getPartitionsArrayFunction(), readResultsProcessingExecutor);
+
+ ListenableFuture> aggregationChunks = Futures.transform(partitionsListFuture,
+ getFetchChunksAsyncFunction(entityType, entityId, key, aggregation, startTs, endTs), readResultsProcessingExecutor);
+
+ return Futures.transform(aggregationChunks, new AggregatePartitionsFunction(aggregation, key, ts), readResultsProcessingExecutor);
+ }
+
+ private Function> getPartitionsArrayFunction() {
+ return rows -> rows.all().stream()
+ .map(row -> row.getLong(ModelConstants.PARTITION_COLUMN)).collect(Collectors.toList());
+ }
+
+ private AsyncFunction, List> getFetchChunksAsyncFunction(String entityType, UUID entityId, String key, Aggregation aggregation, long startTs, long endTs) {
+ return partitions -> {
+ try {
+ PreparedStatement proto = getFetchStmt(aggregation);
+ List futures = new ArrayList<>(partitions.size());
+ for (Long partition : partitions) {
+ BoundStatement stmt = proto.bind();
+ stmt.setString(0, entityType);
+ stmt.setUUID(1, entityId);
+ stmt.setString(2, key);
+ stmt.setLong(3, partition);
+ stmt.setLong(4, startTs);
+ stmt.setLong(5, endTs);
+ log.debug("Generated query [{}] for entityType {} and entityId {}", stmt, entityType, entityId);
+ futures.add(executeAsyncRead(stmt));
+ }
+ return Futures.allAsList(futures);
+ } catch (Throwable e) {
+ log.error("Failed to fetch data", e);
+ throw e;
+ }
+ };
}
@Override
@@ -190,13 +338,12 @@ public class BaseTimeseriesDao extends AbstractDao implements TimeseriesDao {
* Select existing partitions from the table
* {@link ModelConstants#TS_KV_PARTITIONS_CF} for the given entity
*/
- private Long[] fetchPartitions(String entityType, UUID entityId, String key, Optional minPartition, Optional maxPartition) {
+ private ResultSetFuture fetchPartitions(String entityType, UUID entityId, String key, long minPartition, long maxPartition) {
Select.Where select = QueryBuilder.select(ModelConstants.PARTITION_COLUMN).from(ModelConstants.TS_KV_PARTITIONS_CF).where(eq(ModelConstants.ENTITY_TYPE_COLUMN, entityType))
.and(eq(ModelConstants.ENTITY_ID_COLUMN, entityId)).and(eq(ModelConstants.KEY_COLUMN, key));
- minPartition.ifPresent(startTs -> select.and(QueryBuilder.gte(ModelConstants.PARTITION_COLUMN, minPartition.get())));
- maxPartition.ifPresent(endTs -> select.and(QueryBuilder.lte(ModelConstants.PARTITION_COLUMN, maxPartition.get())));
- ResultSet resultSet = executeRead(select);
- return resultSet.all().stream().map(row -> row.getLong(ModelConstants.PARTITION_COLUMN)).toArray(Long[]::new);
+ select.and(QueryBuilder.gte(ModelConstants.PARTITION_COLUMN, minPartition));
+ select.and(QueryBuilder.lte(ModelConstants.PARTITION_COLUMN, maxPartition));
+ return executeAsyncRead(select);
}
private PreparedStatement getSaveStmt(DataType dataType) {
@@ -216,6 +363,30 @@ public class BaseTimeseriesDao extends AbstractDao implements TimeseriesDao {
return saveStmts[dataType.ordinal()];
}
+ private PreparedStatement getFetchStmt(Aggregation aggType) {
+ if (fetchStmts == null) {
+ fetchStmts = new PreparedStatement[Aggregation.values().length];
+ for (Aggregation type : Aggregation.values()) {
+ if (type == Aggregation.SUM && fetchStmts[Aggregation.AVG.ordinal()] != null) {
+ fetchStmts[type.ordinal()] = fetchStmts[Aggregation.AVG.ordinal()];
+ } else if (type == Aggregation.AVG && fetchStmts[Aggregation.SUM.ordinal()] != null) {
+ fetchStmts[type.ordinal()] = fetchStmts[Aggregation.SUM.ordinal()];
+ } else {
+ fetchStmts[type.ordinal()] = getSession().prepare("SELECT " +
+ String.join(", ", ModelConstants.getFetchColumnNames(type)) + " FROM " + ModelConstants.TS_KV_CF
+ + " WHERE " + ModelConstants.ENTITY_TYPE_COLUMN + " = ? "
+ + "AND " + ModelConstants.ENTITY_ID_COLUMN + " = ? "
+ + "AND " + ModelConstants.KEY_COLUMN + " = ? "
+ + "AND " + ModelConstants.PARTITION_COLUMN + " = ? "
+ + "AND " + ModelConstants.TS_COLUMN + " > ? "
+ + "AND " + ModelConstants.TS_COLUMN + " <= ?"
+ + (type == Aggregation.NONE ? " ORDER BY " + ModelConstants.TS_COLUMN + " DESC LIMIT ?" : ""));
+ }
+ }
+ }
+ return fetchStmts[aggType.ordinal()];
+ }
+
private PreparedStatement getLatestStmt(DataType dataType) {
if (latestInsertStmts == null) {
latestInsertStmts = new PreparedStatement[DataType.values().length];
diff --git a/dao/src/main/java/org/thingsboard/server/dao/timeseries/BaseTimeseriesService.java b/dao/src/main/java/org/thingsboard/server/dao/timeseries/BaseTimeseriesService.java
index 419b53447e..f27ed6e605 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/timeseries/BaseTimeseriesService.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/timeseries/BaseTimeseriesService.java
@@ -18,26 +18,31 @@ package org.thingsboard.server.dao.timeseries;
import com.datastax.driver.core.ResultSet;
import com.datastax.driver.core.ResultSetFuture;
import com.datastax.driver.core.Row;
+import com.google.common.base.Function;
import com.google.common.collect.Lists;
import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.ListenableFuture;
import lombok.extern.slf4j.Slf4j;
import org.thingsboard.server.common.data.id.UUIDBased;
+import org.thingsboard.server.common.data.kv.BaseTsKvQuery;
import org.thingsboard.server.common.data.kv.TsKvEntry;
import org.thingsboard.server.common.data.kv.TsKvQuery;
import org.thingsboard.server.dao.exception.IncorrectParameterException;
-import org.slf4j.Logger;
-import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.stereotype.Service;
import org.thingsboard.server.dao.service.Validator;
+import javax.annotation.Nullable;
import javax.annotation.PostConstruct;
+import javax.annotation.PreDestroy;
import java.time.Instant;
import java.time.LocalDateTime;
import java.time.ZoneOffset;
import java.util.*;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.stream.Collectors;
import static org.apache.commons.lang3.StringUtils.isBlank;
@@ -50,38 +55,14 @@ public class BaseTimeseriesService implements TimeseriesService {
public static final int INSERTS_PER_ENTRY = 3;
- @Value("${cassandra.query.ts_key_value_partitioning}")
- private String partitioning;
-
@Autowired
private TimeseriesDao timeseriesDao;
- private TsPartitionDate tsFormat;
-
- @PostConstruct
- public void init() {
- Optional partition = TsPartitionDate.parse(partitioning);
- if (partition.isPresent()) {
- tsFormat = partition.get();
- } else {
- log.warn("Incorrect configuration of partitioning {}", partitioning);
- throw new RuntimeException("Failed to parse partitioning property: " + partitioning + "!");
- }
- }
-
@Override
- public List find(String entityType, UUIDBased entityId, TsKvQuery query) {
+ public ListenableFuture> findAll(String entityType, UUIDBased entityId, List queries) {
validate(entityType, entityId);
- validate(query);
- return timeseriesDao.find(entityType, entityId.getId(), query, toPartitionTs(query.getStartTs()), toPartitionTs(query.getEndTs()));
- }
-
- private Optional toPartitionTs(Optional ts) {
- if (ts.isPresent()) {
- return Optional.of(toPartitionTs(ts.get()));
- } else {
- return Optional.empty();
- }
+ queries.forEach(query -> validate(query));
+ return timeseriesDao.findAllAsync(entityType, entityId.getId(), queries);
}
@Override
@@ -106,7 +87,7 @@ public class BaseTimeseriesService implements TimeseriesService {
throw new IncorrectParameterException("Key value entry can't be null");
}
UUID uid = entityId.getId();
- long partitionTs = toPartitionTs(tsKvEntry.getTs());
+ long partitionTs = timeseriesDao.toPartitionTs(tsKvEntry.getTs());
List futures = Lists.newArrayListWithExpectedSize(INSERTS_PER_ENTRY);
saveAndRegisterFutures(futures, entityType, tsKvEntry, uid, partitionTs);
@@ -122,7 +103,7 @@ public class BaseTimeseriesService implements TimeseriesService {
throw new IncorrectParameterException("Key value entry can't be null");
}
UUID uid = entityId.getId();
- long partitionTs = toPartitionTs(tsKvEntry.getTs());
+ long partitionTs = timeseriesDao.toPartitionTs(tsKvEntry.getTs());
saveAndRegisterFutures(futures, entityType, tsKvEntry, uid, partitionTs);
}
return Futures.allAsList(futures);
@@ -144,14 +125,6 @@ public class BaseTimeseriesService implements TimeseriesService {
futures.add(timeseriesDao.save(entityType, uid, partitionTs, tsKvEntry));
}
- private long toPartitionTs(long ts) {
- LocalDateTime time = LocalDateTime.ofInstant(Instant.ofEpochMilli(ts), ZoneOffset.UTC);
-
- LocalDateTime parititonTime = tsFormat.truncatedTo(time);
-
- return parititonTime.toInstant(ZoneOffset.UTC).toEpochMilli();
- }
-
private static void validate(String entityType, UUIDBased entityId) {
Validator.validateString(entityType, "Incorrect entityType " + entityType);
Validator.validateId(entityId, "Incorrect entityId " + entityId);
@@ -162,6 +135,8 @@ public class BaseTimeseriesService implements TimeseriesService {
throw new IncorrectParameterException("TsKvQuery can't be null");
} else if (isBlank(query.getKey())) {
throw new IncorrectParameterException("Incorrect TsKvQuery. Key can't be empty");
+ } else if (query.getAggregation() == null) {
+ throw new IncorrectParameterException("Incorrect TsKvQuery. Aggregation can't be empty");
}
}
}
diff --git a/dao/src/main/java/org/thingsboard/server/dao/timeseries/SimpleListenableFuture.java b/dao/src/main/java/org/thingsboard/server/dao/timeseries/SimpleListenableFuture.java
new file mode 100644
index 0000000000..3f3e0314c3
--- /dev/null
+++ b/dao/src/main/java/org/thingsboard/server/dao/timeseries/SimpleListenableFuture.java
@@ -0,0 +1,29 @@
+/**
+ * Copyright © 2016-2017 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.google.common.util.concurrent.AbstractFuture;
+
+/**
+ * Created by ashvayka on 21.02.17.
+ */
+public class SimpleListenableFuture extends AbstractFuture {
+
+ public boolean set(V value) {
+ return super.set(value);
+ }
+
+}
diff --git a/dao/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesDao.java b/dao/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesDao.java
index 294f57445c..177003ddf2 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesDao.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesDao.java
@@ -17,6 +17,7 @@ package org.thingsboard.server.dao.timeseries;
import com.datastax.driver.core.ResultSetFuture;
import com.datastax.driver.core.Row;
+import com.google.common.util.concurrent.ListenableFuture;
import org.thingsboard.server.common.data.kv.TsKvEntry;
import org.thingsboard.server.common.data.kv.TsKvQuery;
@@ -30,7 +31,9 @@ import java.util.UUID;
*/
public interface TimeseriesDao {
- List