Browse Source

Introduced base pageable edge event fetcher

pull/4918/head
Volodymyr Babak 5 years ago
parent
commit
d3152874fc
  1. 3
      application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java
  2. 7
      application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/AdminSettingsEdgeEventFetcher.java
  3. 7
      application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/AssetsEdgeEventFetcher.java
  4. 26
      application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/BasePageableEdgeEventFetcher.java
  5. 10
      application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/BaseUsersEdgeEventFetcher.java
  6. 7
      application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/BaseWidgetsBundlesEdgeEventFetcher.java
  7. 7
      application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/DashboardsEdgeEventFetcher.java
  8. 7
      application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/DeviceProfilesEdgeEventFetcher.java
  9. 4
      application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/EdgeEventFetcher.java
  10. 5
      application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/GeneralEdgeEventFetcher.java
  11. 7
      application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/RuleChainsEdgeEventFetcher.java
  12. 2
      application/src/test/java/org/thingsboard/server/edge/imitator/EdgeImitator.java

3
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<EdgeEvent> pageData;
UUID ifOffset = null;
boolean success = true;

7
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<EdgeEvent> fetchEdgeEvents(TenantId tenantId, EdgeId edgeId, PageLink pageLink) {

7
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<EdgeEvent> fetchEdgeEvents(TenantId tenantId, EdgeId edgeId, PageLink pageLink) {
log.trace("[{}] start fetching edge events [{}]", tenantId, edgeId);

26
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);
}
}

10
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<EdgeEvent> fetchEdgeEvents(TenantId tenantId, EdgeId edgeId, PageLink pageLink) {

7
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<EdgeEvent> fetchEdgeEvents(TenantId tenantId, EdgeId edgeId, PageLink pageLink) {

7
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<EdgeEvent> fetchEdgeEvents(TenantId tenantId, EdgeId edgeId, PageLink pageLink) {
log.trace("[{}] start fetching edge events [{}]", tenantId, edgeId);

7
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<EdgeEvent> fetchEdgeEvents(TenantId tenantId, EdgeId edgeId, PageLink pageLink) {
log.trace("[{}] start fetching edge events [{}]", tenantId, edgeId);

4
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<EdgeEvent> fetchEdgeEvents(TenantId tenantId, EdgeId edgeId, PageLink pageLink);
}

5
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),

7
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<EdgeEvent> fetchEdgeEvents(TenantId tenantId, EdgeId edgeId, PageLink pageLink) {
log.trace("[{}] start fetching edge events [{}]", tenantId, edgeId);

2
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,

Loading…
Cancel
Save