diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java index da0b1a7736..71f47e7a69 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java @@ -330,7 +330,6 @@ public final class EdgeGrpcSession implements Closeable { if (isConnected() && isSyncCompleted()) { Long queueStartTs = getQueueStartTs().get(); GeneralEdgeEventFetcher fetcher = new GeneralEdgeEventFetcher( - ctx.getEdgeEventStorageSettings().getMaxReadRecordsCount(), queueStartTs, ctx.getEdgeEventService()); UUID ifOffset = startProcessingEdgeEvents(fetcher); @@ -343,7 +342,7 @@ public final class EdgeGrpcSession implements Closeable { } private UUID startProcessingEdgeEvents(EdgeEventFetcher fetcher) throws InterruptedException { - PageLink pageLink = fetcher.getPageLink(); + PageLink pageLink = fetcher.getPageLink(ctx.getEdgeEventStorageSettings().getMaxReadRecordsCount()); PageData pageData; UUID ifOffset = null; boolean success = true; diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/AdminSettingsEdgeEventFetcher.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/AdminSettingsEdgeEventFetcher.java index b232333d5d..be5344acd4 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/AdminSettingsEdgeEventFetcher.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/AdminSettingsEdgeEventFetcher.java @@ -25,12 +25,7 @@ import org.thingsboard.server.common.data.page.PageLink; @AllArgsConstructor @Slf4j -public class AdminSettingsEdgeEventFetcher implements EdgeEventFetcher { - - @Override - public PageLink getPageLink() { - return new PageLink(DEFAULT_LIMIT); - } +public class AdminSettingsEdgeEventFetcher extends BasePageableEdgeEventFetcher { @Override public PageData fetchEdgeEvents(TenantId tenantId, EdgeId edgeId, PageLink pageLink) { diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/AssetsEdgeEventFetcher.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/AssetsEdgeEventFetcher.java index e2eb2e8d61..30bbcd0aa1 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/AssetsEdgeEventFetcher.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/AssetsEdgeEventFetcher.java @@ -33,15 +33,10 @@ import java.util.List; @AllArgsConstructor @Slf4j -public class AssetsEdgeEventFetcher implements EdgeEventFetcher { +public class AssetsEdgeEventFetcher extends BasePageableEdgeEventFetcher { private final AssetService assetService; - @Override - public PageLink getPageLink() { - return new PageLink(DEFAULT_LIMIT); - } - @Override public PageData fetchEdgeEvents(TenantId tenantId, EdgeId edgeId, PageLink pageLink) { log.trace("[{}] start fetching edge events [{}]", tenantId, edgeId); diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/BasePageableEdgeEventFetcher.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/BasePageableEdgeEventFetcher.java new file mode 100644 index 0000000000..3c2f7f82fd --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/BasePageableEdgeEventFetcher.java @@ -0,0 +1,26 @@ +/** + * Copyright © 2016-2021 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.service.edge.rpc.fetch; + +import org.thingsboard.server.common.data.page.PageLink; + +public abstract class BasePageableEdgeEventFetcher implements EdgeEventFetcher { + + @Override + public PageLink getPageLink(int pageSize) { + return new PageLink(pageSize); + } +} diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/BaseUsersEdgeEventFetcher.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/BaseUsersEdgeEventFetcher.java index 883f451636..dbb98864d3 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/BaseUsersEdgeEventFetcher.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/BaseUsersEdgeEventFetcher.java @@ -15,7 +15,6 @@ */ package org.thingsboard.server.service.edge.rpc.fetch; -import lombok.AllArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.thingsboard.server.common.data.User; import org.thingsboard.server.common.data.edge.EdgeEvent; @@ -25,20 +24,13 @@ import org.thingsboard.server.common.data.id.EdgeId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.page.PageLink; -import org.thingsboard.server.common.data.widget.WidgetsBundle; -import org.thingsboard.server.dao.user.UserService; import org.thingsboard.server.service.edge.rpc.EdgeEventUtils; import java.util.ArrayList; import java.util.List; @Slf4j -public abstract class BaseUsersEdgeEventFetcher implements EdgeEventFetcher { - - @Override - public PageLink getPageLink() { - return new PageLink(DEFAULT_LIMIT); - } +public abstract class BaseUsersEdgeEventFetcher extends BasePageableEdgeEventFetcher { @Override public PageData fetchEdgeEvents(TenantId tenantId, EdgeId edgeId, PageLink pageLink) { diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/BaseWidgetsBundlesEdgeEventFetcher.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/BaseWidgetsBundlesEdgeEventFetcher.java index 53c6e579d3..0d02b6dd34 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/BaseWidgetsBundlesEdgeEventFetcher.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/BaseWidgetsBundlesEdgeEventFetcher.java @@ -30,12 +30,7 @@ import java.util.ArrayList; import java.util.List; @Slf4j -public abstract class BaseWidgetsBundlesEdgeEventFetcher implements EdgeEventFetcher { - - @Override - public PageLink getPageLink() { - return new PageLink(DEFAULT_LIMIT); - } +public abstract class BaseWidgetsBundlesEdgeEventFetcher extends BasePageableEdgeEventFetcher { @Override public PageData fetchEdgeEvents(TenantId tenantId, EdgeId edgeId, PageLink pageLink) { diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/DashboardsEdgeEventFetcher.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/DashboardsEdgeEventFetcher.java index 515a356850..751c6690c1 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/DashboardsEdgeEventFetcher.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/DashboardsEdgeEventFetcher.java @@ -33,15 +33,10 @@ import java.util.List; @AllArgsConstructor @Slf4j -public class DashboardsEdgeEventFetcher implements EdgeEventFetcher { +public class DashboardsEdgeEventFetcher extends BasePageableEdgeEventFetcher { private final DashboardService dashboardService; - @Override - public PageLink getPageLink() { - return new PageLink(DEFAULT_LIMIT); - } - @Override public PageData fetchEdgeEvents(TenantId tenantId, EdgeId edgeId, PageLink pageLink) { log.trace("[{}] start fetching edge events [{}]", tenantId, edgeId); diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/DeviceProfilesEdgeEventFetcher.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/DeviceProfilesEdgeEventFetcher.java index fb0af06b75..011761cf11 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/DeviceProfilesEdgeEventFetcher.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/DeviceProfilesEdgeEventFetcher.java @@ -33,15 +33,10 @@ import java.util.List; @AllArgsConstructor @Slf4j -public class DeviceProfilesEdgeEventFetcher implements EdgeEventFetcher { +public class DeviceProfilesEdgeEventFetcher extends BasePageableEdgeEventFetcher { private final DeviceProfileService deviceProfileService; - @Override - public PageLink getPageLink() { - return new PageLink(DEFAULT_LIMIT); - } - @Override public PageData fetchEdgeEvents(TenantId tenantId, EdgeId edgeId, PageLink pageLink) { log.trace("[{}] start fetching edge events [{}]", tenantId, edgeId); diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/EdgeEventFetcher.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/EdgeEventFetcher.java index e283b2c043..d36b072c3d 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/EdgeEventFetcher.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/EdgeEventFetcher.java @@ -23,9 +23,7 @@ import org.thingsboard.server.common.data.page.PageLink; public interface EdgeEventFetcher { - final int DEFAULT_LIMIT = 100; - - PageLink getPageLink(); + PageLink getPageLink(int pageSize); PageData fetchEdgeEvents(TenantId tenantId, EdgeId edgeId, PageLink pageLink); } diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/GeneralEdgeEventFetcher.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/GeneralEdgeEventFetcher.java index a87a4ebd68..39d0b482ed 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/GeneralEdgeEventFetcher.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/GeneralEdgeEventFetcher.java @@ -28,14 +28,13 @@ import org.thingsboard.server.dao.edge.EdgeEventService; @AllArgsConstructor public class GeneralEdgeEventFetcher implements EdgeEventFetcher { - private final int maxReadRecordsCount; private final Long queueStartTs; private final EdgeEventService edgeEventService; @Override - public PageLink getPageLink() { + public PageLink getPageLink(int pageSize) { return new TimePageLink( - maxReadRecordsCount, + pageSize, 0, null, new SortOrder("createdTime", SortOrder.Direction.ASC), diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/RuleChainsEdgeEventFetcher.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/RuleChainsEdgeEventFetcher.java index 2aca092856..3da1dcc47c 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/RuleChainsEdgeEventFetcher.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/RuleChainsEdgeEventFetcher.java @@ -33,15 +33,10 @@ import java.util.List; @Slf4j @AllArgsConstructor -public class RuleChainsEdgeEventFetcher implements EdgeEventFetcher { +public class RuleChainsEdgeEventFetcher extends BasePageableEdgeEventFetcher { private final RuleChainService ruleChainService; - @Override - public PageLink getPageLink() { - return new PageLink(DEFAULT_LIMIT); - } - @Override public PageData fetchEdgeEvents(TenantId tenantId, EdgeId edgeId, PageLink pageLink) { log.trace("[{}] start fetching edge events [{}]", tenantId, edgeId); diff --git a/application/src/test/java/org/thingsboard/server/edge/imitator/EdgeImitator.java b/application/src/test/java/org/thingsboard/server/edge/imitator/EdgeImitator.java index 7625647144..6fff7d31f6 100644 --- a/application/src/test/java/org/thingsboard/server/edge/imitator/EdgeImitator.java +++ b/application/src/test/java/org/thingsboard/server/edge/imitator/EdgeImitator.java @@ -5,7 +5,7 @@ * 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 + * 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,