Browse Source

Code review changes - implemented grpc sending in async manner

pull/4571/head
Volodymyr Babak 5 years ago
parent
commit
47311663cd
  1. 4
      application/src/main/java/org/thingsboard/server/service/edge/EdgeContextComponent.java
  2. 60
      application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcService.java
  3. 280
      application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java
  4. 74
      application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeSyncCursor.java
  5. 2
      application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/AdminSettingsEdgeEventFetcher.java
  6. 49
      application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/CustomerEdgeEventFetcher.java
  7. 33
      application/src/main/java/org/thingsboard/server/service/executors/GrpcCallbackExecutorService.java
  8. 4
      application/src/main/resources/thingsboard.yml

4
application/src/main/java/org/thingsboard/server/service/edge/EdgeContextComponent.java

@ -48,6 +48,7 @@ import org.thingsboard.server.service.edge.rpc.processor.WidgetBundleEdgeProcess
import org.thingsboard.server.service.edge.rpc.processor.WidgetTypeEdgeProcessor;
import org.thingsboard.server.service.edge.rpc.sync.EdgeRequestsService;
import org.thingsboard.server.service.executors.DbCallbackExecutorService;
import org.thingsboard.server.service.executors.GrpcCallbackExecutorService;
@Component
@TbCoreComponent
@ -138,4 +139,7 @@ public class EdgeContextComponent {
@Autowired
private DbCallbackExecutorService dbCallbackExecutor;
@Autowired
private GrpcCallbackExecutorService grpcCallbackExecutorService;
}

60
application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcService.java

@ -16,8 +16,8 @@
package org.thingsboard.server.service.edge.rpc;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.google.common.io.Resources;
import com.google.common.util.concurrent.FutureCallback;
import com.google.common.util.concurrent.Futures;
import io.grpc.Server;
import io.grpc.netty.NettyServerBuilder;
import io.grpc.stub.StreamObserver;
@ -46,7 +46,6 @@ import org.thingsboard.server.service.telemetry.TelemetrySubscriptionService;
import javax.annotation.Nullable;
import javax.annotation.PostConstruct;
import javax.annotation.PreDestroy;
import java.io.File;
import java.io.IOException;
import java.io.InputStream;
import java.util.Collections;
@ -93,6 +92,9 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i
@Value("${edges.scheduler_pool_size}")
private int schedulerPoolSize;
@Value("${edges.send_scheduler_pool_size}")
private int sendSchedulerPoolSize;
@Autowired
private EdgeContextComponent ctx;
@ -105,6 +107,8 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i
private ExecutorService syncExecutorService;
private ScheduledExecutorService sendScheduler;
@PostConstruct
public void init() {
log.info("Initializing Edge RPC service!");
@ -131,6 +135,7 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i
throw new RuntimeException("Failed to start Edge RPC server!");
}
this.scheduler = Executors.newScheduledThreadPool(schedulerPoolSize, ThingsBoardThreadFactory.forName("edge-scheduler"));
this.sendScheduler = Executors.newScheduledThreadPool(sendSchedulerPoolSize, ThingsBoardThreadFactory.forName("edge-send-scheduler"));
this.syncExecutorService = Executors.newFixedThreadPool(
Runtime.getRuntime().availableProcessors(), ThingsBoardThreadFactory.forName("edge-sync"));
log.info("Edge RPC service initialized!");
@ -152,6 +157,9 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i
if (scheduler != null) {
scheduler.shutdownNow();
}
if (sendScheduler != null) {
sendScheduler.shutdownNow();
}
if (syncExecutorService != null) {
syncExecutorService.shutdownNow();
}
@ -159,7 +167,7 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i
@Override
public StreamObserver<RequestMsg> handleMsgs(StreamObserver<ResponseMsg> outputStream) {
return new EdgeGrpcSession(ctx, outputStream, this::onEdgeConnect, this::onEdgeDisconnect, mapper, syncExecutorService).getInputStream();
return new EdgeGrpcSession(ctx, outputStream, this::onEdgeConnect, this::onEdgeDisconnect, mapper, syncExecutorService, sendScheduler).getInputStream();
}
@Override
@ -180,12 +188,12 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i
log.info("[{}] Closing and removing session for edge [{}]", tenantId, edgeId);
session.close();
sessions.remove(edgeId);
Lock lock = sessionNewEventsLocks.computeIfAbsent(edgeId, id -> new ReentrantLock());
lock.lock();
final Lock newEventLock = sessionNewEventsLocks.computeIfAbsent(edgeId, id -> new ReentrantLock());
newEventLock.lock();
try {
sessionNewEvents.remove(edgeId);
} finally {
lock.unlock();
newEventLock.unlock();
}
cancelScheduleEdgeEventsCheck(edgeId);
}
@ -194,27 +202,27 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i
@Override
public void onEdgeEvent(TenantId tenantId, EdgeId edgeId) {
log.trace("[{}] onEdgeEvent [{}]", tenantId, edgeId.getId());
Lock lock = sessionNewEventsLocks.computeIfAbsent(edgeId, id -> new ReentrantLock());
lock.lock();
final Lock newEventLock = sessionNewEventsLocks.computeIfAbsent(edgeId, id -> new ReentrantLock());
newEventLock.lock();
try {
if (Boolean.FALSE.equals(sessionNewEvents.get(edgeId))) {
log.trace("[{}] set session new events flag to true [{}]", tenantId, edgeId.getId());
sessionNewEvents.put(edgeId, true);
}
} finally {
lock.unlock();
newEventLock.unlock();
}
}
private void onEdgeConnect(EdgeId edgeId, EdgeGrpcSession edgeGrpcSession) {
log.info("[{}] edge [{}] connected successfully.", edgeGrpcSession.getSessionId(), edgeId);
sessions.put(edgeId, edgeGrpcSession);
Lock lock = sessionNewEventsLocks.computeIfAbsent(edgeId, id -> new ReentrantLock());
lock.lock();
final Lock newEventLock = sessionNewEventsLocks.computeIfAbsent(edgeId, id -> new ReentrantLock());
newEventLock.lock();
try {
sessionNewEvents.put(edgeId, true);
} finally {
lock.unlock();
newEventLock.unlock();
}
save(edgeId, DefaultDeviceStateService.ACTIVITY_STATE, true);
save(edgeId, DefaultDeviceStateService.LAST_CONNECT_TIME, System.currentTimeMillis());
@ -239,21 +247,33 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i
if (sessions.containsKey(edgeId)) {
ScheduledFuture<?> schedule = scheduler.schedule(() -> {
try {
Lock lock = sessionNewEventsLocks.computeIfAbsent(edgeId, id -> new ReentrantLock());
lock.lock();
final Lock newEventLock = sessionNewEventsLocks.computeIfAbsent(edgeId, id -> new ReentrantLock());
newEventLock.lock();
try {
if (Boolean.TRUE.equals(sessionNewEvents.get(edgeId))) {
log.trace("[{}] Set session new events flag to false", edgeId.getId());
sessionNewEvents.put(edgeId, false);
session.processEdgeEvents();
Futures.addCallback(session.processEdgeEvents(), new FutureCallback<>() {
@Override
public void onSuccess(Void result) {
scheduleEdgeEventsCheck(session);
}
@Override
public void onFailure(Throwable t) {
log.warn("[{}] Failed to process edge events for edge [{}]!", tenantId, session.getEdge().getId().getId(), t);
scheduleEdgeEventsCheck(session);
}
}, ctx.getGrpcCallbackExecutorService());
} else {
scheduleEdgeEventsCheck(session);
}
} finally {
lock.unlock();
newEventLock.unlock();
}
} catch (Exception e) {
log.warn("[{}] Failed to process edge events for edge [{}]!", tenantId, session.getEdge().getId().getId(), e);
}
scheduleEdgeEventsCheck(session);
}, ctx.getEdgeEventStorageSettings().getNoRecordsSleepInterval(), TimeUnit.MILLISECONDS);
sessionEdgeEventChecks.put(edgeId, schedule);
log.trace("[{}] Check edge event scheduled for edge [{}]", tenantId, edgeId.getId());
@ -277,12 +297,12 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i
private void onEdgeDisconnect(EdgeId edgeId) {
log.info("[{}] edge disconnected!", edgeId);
sessions.remove(edgeId);
Lock lock = sessionNewEventsLocks.computeIfAbsent(edgeId, id -> new ReentrantLock());
lock.lock();
final Lock newEventLock = sessionNewEventsLocks.computeIfAbsent(edgeId, id -> new ReentrantLock());
newEventLock.lock();
try {
sessionNewEvents.remove(edgeId);
} finally {
lock.unlock();
newEventLock.unlock();
}
save(edgeId, DefaultDeviceStateService.ACTIVITY_STATE, false);
save(edgeId, DefaultDeviceStateService.LAST_DISCONNECT_TIME, System.currentTimeMillis());

280
application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java

@ -20,7 +20,7 @@ import com.fasterxml.jackson.databind.ObjectMapper;
import com.google.common.util.concurrent.FutureCallback;
import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.ListenableFuture;
import com.google.common.util.concurrent.MoreExecutors;
import com.google.common.util.concurrent.SettableFuture;
import io.grpc.stub.StreamObserver;
import lombok.Data;
import lombok.extern.slf4j.Slf4j;
@ -31,9 +31,7 @@ import org.thingsboard.server.common.data.EdgeUtils;
import org.thingsboard.server.common.data.edge.Edge;
import org.thingsboard.server.common.data.edge.EdgeEvent;
import org.thingsboard.server.common.data.edge.EdgeEventActionType;
import org.thingsboard.server.common.data.edge.EdgeEventType;
import org.thingsboard.server.common.data.id.EdgeId;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.kv.AttributeKvEntry;
import org.thingsboard.server.common.data.kv.BaseAttributeKvEntry;
@ -69,32 +67,21 @@ import org.thingsboard.server.gen.edge.v1.UplinkResponseMsg;
import org.thingsboard.server.gen.edge.v1.UserCredentialsRequestMsg;
import org.thingsboard.server.gen.edge.v1.WidgetBundleTypesRequestMsg;
import org.thingsboard.server.service.edge.EdgeContextComponent;
import org.thingsboard.server.service.edge.rpc.fetch.AdminSettingsEdgeEventFetcher;
import org.thingsboard.server.service.edge.rpc.fetch.AssetsEdgeEventFetcher;
import org.thingsboard.server.service.edge.rpc.fetch.CustomerUsersEdgeEventFetcher;
import org.thingsboard.server.service.edge.rpc.fetch.DashboardsEdgeEventFetcher;
import org.thingsboard.server.service.edge.rpc.fetch.DeviceProfilesEdgeEventFetcher;
import org.thingsboard.server.service.edge.rpc.fetch.EdgeEventFetcher;
import org.thingsboard.server.service.edge.rpc.fetch.GeneralEdgeEventFetcher;
import org.thingsboard.server.service.edge.rpc.fetch.RuleChainsEdgeEventFetcher;
import org.thingsboard.server.service.edge.rpc.fetch.SystemWidgetsBundlesEdgeEventFetcher;
import org.thingsboard.server.service.edge.rpc.fetch.TenantAdminUsersEdgeEventFetcher;
import org.thingsboard.server.service.edge.rpc.fetch.TenantWidgetsBundlesEdgeEventFetcher;
import java.io.Closeable;
import java.util.ArrayList;
import java.util.Collection;
import java.util.Collections;
import java.util.HashMap;
import java.util.Iterator;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.Objects;
import java.util.Optional;
import java.util.UUID;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.ScheduledFuture;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.locks.ReentrantLock;
import java.util.function.BiConsumer;
@ -116,6 +103,8 @@ public final class EdgeGrpcSession implements Closeable {
private final ObjectMapper mapper;
private final Map<Integer, DownlinkMsg> pendingMsgsMap = new LinkedHashMap<>();
private SettableFuture<Void> pendingFuture;
private ScheduledFuture<?> sendSchedule;
private EdgeContextComponent ctx;
private Edge edge;
@ -126,10 +115,10 @@ public final class EdgeGrpcSession implements Closeable {
private ExecutorService syncExecutorService;
private CountDownLatch latch;
private ScheduledExecutorService sendScheduler;
EdgeGrpcSession(EdgeContextComponent ctx, StreamObserver<ResponseMsg> outputStream, BiConsumer<EdgeId, EdgeGrpcSession> sessionOpenListener,
Consumer<EdgeId> sessionCloseListener, ObjectMapper mapper, ExecutorService syncExecutorService) {
Consumer<EdgeId> sessionCloseListener, ObjectMapper mapper, ExecutorService syncExecutorService, ScheduledExecutorService sendScheduler) {
this.sessionId = UUID.randomUUID();
this.ctx = ctx;
this.outputStream = outputStream;
@ -137,6 +126,7 @@ public final class EdgeGrpcSession implements Closeable {
this.sessionCloseListener = sessionCloseListener;
this.mapper = mapper;
this.syncExecutorService = syncExecutorService;
this.sendScheduler = sendScheduler;
initInputStream();
}
@ -155,19 +145,21 @@ public final class EdgeGrpcSession implements Closeable {
connected = true;
}
}
if (connected && requestMsg.getMsgType().equals(RequestMsgType.SYNC_REQUEST_RPC_MESSAGE)) {
if (requestMsg.getSyncRequestMsg().getSyncRequired()) {
startSyncProcess(edge.getTenantId(), edge.getId());
} else {
syncCompleted = true;
}
}
if (connected) {
if (requestMsg.getMsgType().equals(RequestMsgType.UPLINK_RPC_MESSAGE) && requestMsg.hasUplinkMsg()) {
onUplinkMsg(requestMsg.getUplinkMsg());
if (requestMsg.getMsgType().equals(RequestMsgType.SYNC_REQUEST_RPC_MESSAGE)) {
if (requestMsg.hasSyncRequestMsg() && requestMsg.getSyncRequestMsg().getSyncRequired()) {
startSyncProcess(edge.getTenantId(), edge.getId());
} else {
syncCompleted = true;
}
}
if (requestMsg.getMsgType().equals(RequestMsgType.UPLINK_RPC_MESSAGE) && requestMsg.hasDownlinkResponseMsg()) {
onDownlinkResponse(requestMsg.getDownlinkResponseMsg());
if (requestMsg.getMsgType().equals(RequestMsgType.UPLINK_RPC_MESSAGE)) {
if (requestMsg.hasUplinkMsg()) {
onUplinkMsg(requestMsg.getUplinkMsg());
}
if (requestMsg.hasDownlinkResponseMsg()) {
onDownlinkResponse(requestMsg.getDownlinkResponseMsg());
}
}
}
}
@ -204,39 +196,47 @@ public final class EdgeGrpcSession implements Closeable {
syncCompleted = false;
syncExecutorService.submit(() -> {
try {
startProcessingEdgeEvents(new SystemWidgetsBundlesEdgeEventFetcher(ctx.getWidgetsBundleService()));
startProcessingEdgeEvents(new TenantWidgetsBundlesEdgeEventFetcher(ctx.getWidgetsBundleService()));
startProcessingEdgeEvents(new DeviceProfilesEdgeEventFetcher(ctx.getDeviceProfileService()));
startProcessingEdgeEvents(new RuleChainsEdgeEventFetcher(ctx.getRuleChainService()));
startProcessingEdgeEvents(new TenantAdminUsersEdgeEventFetcher(ctx.getUserService()));
if (edge.getCustomerId() != null && !EntityId.NULL_UUID.equals(edge.getCustomerId().getId())) {
EdgeEvent customerEdgeEvent = EdgeEventUtils.constructEdgeEvent(edge.getTenantId(), edge.getId(),
EdgeEventType.CUSTOMER, EdgeEventActionType.ADDED, edge.getCustomerId(), null);
DownlinkMsg customerDownlinkMsg = convertToDownlinkMsg(customerEdgeEvent);
sendDownlinkMsgsPack(Collections.singletonList(customerDownlinkMsg));
startProcessingEdgeEvents(new CustomerUsersEdgeEventFetcher(ctx.getUserService(), edge.getCustomerId()));
}
startProcessingEdgeEvents(new AdminSettingsEdgeEventFetcher(ctx.getAdminSettingsService()));
startProcessingEdgeEvents(new AssetsEdgeEventFetcher(ctx.getAssetService()));
startProcessingEdgeEvents(new DashboardsEdgeEventFetcher(ctx.getDashboardService()));
DownlinkMsg syncCompleteDownlinkMsg = DownlinkMsg.newBuilder()
.setDownlinkMsgId(EdgeUtils.nextPositiveInt())
.setSyncCompletedMsg(SyncCompletedMsg.newBuilder().build())
.build();
sendDownlinkMsgsPack(Collections.singletonList(syncCompleteDownlinkMsg));
syncCompleted = true;
doSync(new EdgeSyncCursor(ctx, edge));
} catch (Exception e) {
log.error("[{}][{}] Exception during sync process", edge.getTenantId(), edge.getId(), e);
}
});
}
private void doSync(EdgeSyncCursor cursor) {
if (cursor.hasNext()) {
log.info("[{}][{}] starting sync process, cursor current idx = {}", edge.getTenantId(), edge.getId(), cursor.getCurrentIdx());
ListenableFuture<UUID> uuidListenableFuture = startProcessingEdgeEvents(cursor.getNext());
Futures.addCallback(uuidListenableFuture, new FutureCallback<>() {
@Override
public void onSuccess(@Nullable UUID result) {
doSync(cursor);
}
@Override
public void onFailure(Throwable t) {
log.error("[{}][{}] Exception during sync process", edge.getTenantId(), edge.getId(), t);
}
}, ctx.getGrpcCallbackExecutorService());
} else {
DownlinkMsg syncCompleteDownlinkMsg = DownlinkMsg.newBuilder()
.setDownlinkMsgId(EdgeUtils.nextPositiveInt())
.setSyncCompletedMsg(SyncCompletedMsg.newBuilder().build())
.build();
Futures.addCallback(sendDownlinkMsgsPack(Collections.singletonList(syncCompleteDownlinkMsg)), new FutureCallback<Void>() {
@Override
public void onSuccess(Void result) {
syncCompleted = true;
}
@Override
public void onFailure(Throwable t) {
log.error("[{}][{}] Exception during sending sync complete", edge.getTenantId(), edge.getId(), t);
}
}, ctx.getGrpcCallbackExecutorService());
}
}
private void onUplinkMsg(UplinkMsg uplinkMsg) {
ListenableFuture<List<Void>> future = processUplinkMsg(uplinkMsg);
Futures.addCallback(future, new FutureCallback<>() {
@ -260,7 +260,7 @@ public final class EdgeGrpcSession implements Closeable {
.setUplinkResponseMsg(uplinkResponseMsg)
.build());
}
}, MoreExecutors.directExecutor());
}, ctx.getGrpcCallbackExecutorService());
}
private void onDownlinkResponse(DownlinkResponseMsg msg) {
@ -271,7 +271,13 @@ public final class EdgeGrpcSession implements Closeable {
} else {
log.error("[{}] Msg processing failed! Error msg: {}", edge.getRoutingKey(), msg.getErrorMsg());
}
latch.countDown();
if (pendingMsgsMap.size() == 0) {
log.debug("[{}] Pending msgs map is empty. Stopping current iteration {}", edge.getRoutingKey(), msg);
if (sendSchedule != null) {
sendSchedule.cancel(false);
}
pendingFuture.set(null);
}
} catch (Exception e) {
log.error("[{}] Can't process downlink response message [{}]", this.sessionId, msg, e);
}
@ -305,76 +311,136 @@ public final class EdgeGrpcSession implements Closeable {
sendDownlinkMsg(edgeConfigMsg);
}
void processEdgeEvents() throws Exception {
log.trace("[{}] processHandleMessages started", this.sessionId);
ListenableFuture<Void> processEdgeEvents() throws Exception {
SettableFuture<Void> result = SettableFuture.create();
log.trace("[{}] starting processing edge events", this.sessionId);
if (isConnected() && isSyncCompleted()) {
Long queueStartTs = getQueueStartTs().get();
GeneralEdgeEventFetcher fetcher = new GeneralEdgeEventFetcher(
queueStartTs,
ctx.getEdgeEventService());
UUID ifOffset = startProcessingEdgeEvents(fetcher);
if (ifOffset != null) {
Long newStartTs = Uuids.unixTimestamp(ifOffset);
updateQueueStartTs(newStartTs);
log.debug("[{}] queue offset was updated [{}][{}]", this.sessionId, ifOffset, newStartTs);
}
ListenableFuture<UUID> ifOffsetFuture = startProcessingEdgeEvents(fetcher);
Futures.addCallback(ifOffsetFuture, new FutureCallback<>() {
@Override
public void onSuccess(@Nullable UUID ifOffset) {
if (ifOffset != null) {
Long newStartTs = Uuids.unixTimestamp(ifOffset);
ListenableFuture<List<Void>> updateFuture = updateQueueStartTs(newStartTs);
Futures.addCallback(updateFuture, new FutureCallback<>() {
@Override
public void onSuccess(@Nullable List<Void> list) {
log.debug("[{}] queue offset was updated [{}][{}]", sessionId, ifOffset, newStartTs);
result.set(null);
}
@Override
public void onFailure(Throwable t) {
log.error("[{}] Failed to update queue offset [{}]", sessionId, ifOffset, t);
result.setException(t);
}
}, ctx.getGrpcCallbackExecutorService());
} else {
log.trace("[{}] ifOffset is null. Skipping iteration without db update", sessionId);
result.set(null);
}
}
@Override
public void onFailure(Throwable t) {
log.error("[{}] Failed to process events", sessionId, t);
result.setException(t);
}
}, ctx.getGrpcCallbackExecutorService());
} else {
log.trace("[{}] edge is not connected or sync is not completed. Skipping iteration", sessionId);
result.set(null);
}
log.trace("[{}] processHandleMessages finished", this.sessionId);
return result;
}
private UUID startProcessingEdgeEvents(EdgeEventFetcher fetcher) throws Exception {
private ListenableFuture<UUID> startProcessingEdgeEvents(EdgeEventFetcher fetcher) {
SettableFuture<UUID> result = SettableFuture.create();
PageLink pageLink = fetcher.getPageLink(ctx.getEdgeEventStorageSettings().getMaxReadRecordsCount());
PageData<EdgeEvent> pageData;
UUID ifOffset = null;
boolean success;
do {
pageData = fetcher.fetchEdgeEvents(edge.getTenantId(), edge, pageLink);
processEdgeEvents(fetcher, pageLink, result);
return result;
}
private void processEdgeEvents(EdgeEventFetcher fetcher, PageLink pageLink, SettableFuture<UUID> result) {
try {
PageData<EdgeEvent> pageData = fetcher.fetchEdgeEvents(edge.getTenantId(), edge, pageLink);
if (isConnected() && !pageData.getData().isEmpty()) {
log.trace("[{}] [{}] event(s) are going to be processed.", this.sessionId, pageData.getData().size());
List<DownlinkMsg> downlinkMsgsPack = convertToDownlinkMsgsPack(pageData.getData());
success = sendDownlinkMsgsPack(downlinkMsgsPack);
ifOffset = pageData.getData().get(pageData.getData().size() - 1).getUuidId();
if (success) {
pageLink = pageLink.nextPageLink();
}
Futures.addCallback(sendDownlinkMsgsPack(downlinkMsgsPack), new FutureCallback<Void>() {
@Override
public void onSuccess(@Nullable Void tmp) {
if (isConnected() && pageData.hasNext()) {
processEdgeEvents(fetcher, pageLink.nextPageLink(), result);
} else {
UUID ifOffset = pageData.getData().get(pageData.getData().size() - 1).getUuidId();
result.set(ifOffset);
}
}
@Override
public void onFailure(Throwable t) {
log.error("[{}] Failed to send downlink msgs pack", sessionId, t);
result.setException(t);
}
}, ctx.getGrpcCallbackExecutorService());
} else {
log.trace("[{}] no event(s) found. Stop processing edge events", this.sessionId);
result.set(null);
}
} while (isConnected() && pageData.hasNext());
return ifOffset;
} catch (Exception e) {
log.error("[{}] Failed to fetch edge events", this.sessionId, e);
result.setException(e);
}
}
private boolean sendDownlinkMsgsPack(List<DownlinkMsg> downlinkMsgsPack) throws InterruptedException {
private ListenableFuture<Void> sendDownlinkMsgsPack(List<DownlinkMsg> downlinkMsgsPack) {
SettableFuture<Void> result = SettableFuture.create();
downlinkMsgsPackLock.lock();
try {
boolean success;
pendingMsgsMap.clear();
downlinkMsgsPack.forEach(msg -> pendingMsgsMap.put(msg.getDownlinkMsgId(), msg));
do {
log.trace("[{}] [{}] downlink msg(s) are going to be send.", this.sessionId, pendingMsgsMap.values().size());
latch = new CountDownLatch(pendingMsgsMap.values().size());
List<DownlinkMsg> copy = new ArrayList<>(pendingMsgsMap.values());
for (DownlinkMsg downlinkMsg : copy) {
sendDownlinkMsg(ResponseMsg.newBuilder()
.setDownlinkMsg(downlinkMsg)
.build());
}
success = latch.await(10, TimeUnit.SECONDS);
if (!success || pendingMsgsMap.values().size() > 0) {
log.warn("[{}] Failed to deliver the batch: {}", this.sessionId, pendingMsgsMap.values());
}
if (isConnected() && (!success || pendingMsgsMap.values().size() > 0)) {
try {
Thread.sleep(ctx.getEdgeEventStorageSettings().getSleepIntervalBetweenBatches());
} catch (InterruptedException e) {
log.error("[{}] Error during sleep between batches", this.sessionId, e);
}
}
} while (isConnected() && (!success || pendingMsgsMap.values().size() > 0));
return success;
scheduleDownlinkMsgsPackSend(true, result);
} finally {
downlinkMsgsPackLock.unlock();
}
return result;
}
private void scheduleDownlinkMsgsPackSend(boolean firstRun, SettableFuture<Void> futureResult) {
pendingFuture = futureResult;
Runnable runnable = () -> {
try {
if (isConnected() && pendingMsgsMap.values().size() > 0) {
if (!firstRun) {
log.warn("[{}] Failed to deliver the batch: {}", this.sessionId, pendingMsgsMap.values());
}
log.trace("[{}] [{}] downlink msg(s) are going to be send.", this.sessionId, pendingMsgsMap.values().size());
List<DownlinkMsg> copy = new ArrayList<>(pendingMsgsMap.values());
for (DownlinkMsg downlinkMsg : copy) {
sendDownlinkMsg(ResponseMsg.newBuilder()
.setDownlinkMsg(downlinkMsg)
.build());
}
scheduleDownlinkMsgsPackSend(false, futureResult);
} else {
futureResult.set(null);
}
} catch (Exception e) {
futureResult.setException(e);
}
};
if (firstRun) {
syncExecutorService.submit(runnable);
} else {
sendSchedule = sendScheduler.schedule(runnable, ctx.getEdgeEventStorageSettings().getSleepIntervalBetweenBatches(), TimeUnit.MILLISECONDS);
}
}
private DownlinkMsg convertToDownlinkMsg(EdgeEvent edgeEvent) {
@ -439,14 +505,14 @@ public final class EdgeGrpcSession implements Closeable {
} else {
return 0L;
}
}, ctx.getDbCallbackExecutor());
}, ctx.getGrpcCallbackExecutorService());
}
private void updateQueueStartTs(Long newStartTs) {
private ListenableFuture<List<Void>> updateQueueStartTs(Long newStartTs) {
log.trace("[{}] updating QueueStartTs [{}][{}]", this.sessionId, edge.getId(), newStartTs);
newStartTs = ++newStartTs; // increments ts by 1 - next edge event search starts from current offset + 1
List<AttributeKvEntry> attributes = Collections.singletonList(new BaseAttributeKvEntry(new LongDataEntry(QUEUE_START_TS_ATTR_KEY, newStartTs), System.currentTimeMillis()));
ctx.getAttributesService().save(edge.getTenantId(), edge.getId(), DataConstants.SERVER_SCOPE, attributes);
return ctx.getAttributesService().save(edge.getTenantId(), edge.getId(), DataConstants.SERVER_SCOPE, attributes);
}
private DownlinkMsg processEntityMessage(EdgeEvent edgeEvent, EdgeEventActionType action) {

74
application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeSyncCursor.java

@ -0,0 +1,74 @@
/**
* 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;
import org.thingsboard.server.common.data.edge.Edge;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.service.edge.EdgeContextComponent;
import org.thingsboard.server.service.edge.rpc.fetch.AdminSettingsEdgeEventFetcher;
import org.thingsboard.server.service.edge.rpc.fetch.AssetsEdgeEventFetcher;
import org.thingsboard.server.service.edge.rpc.fetch.CustomerEdgeEventFetcher;
import org.thingsboard.server.service.edge.rpc.fetch.CustomerUsersEdgeEventFetcher;
import org.thingsboard.server.service.edge.rpc.fetch.DashboardsEdgeEventFetcher;
import org.thingsboard.server.service.edge.rpc.fetch.DeviceProfilesEdgeEventFetcher;
import org.thingsboard.server.service.edge.rpc.fetch.EdgeEventFetcher;
import org.thingsboard.server.service.edge.rpc.fetch.RuleChainsEdgeEventFetcher;
import org.thingsboard.server.service.edge.rpc.fetch.SystemWidgetsBundlesEdgeEventFetcher;
import org.thingsboard.server.service.edge.rpc.fetch.TenantAdminUsersEdgeEventFetcher;
import org.thingsboard.server.service.edge.rpc.fetch.TenantWidgetsBundlesEdgeEventFetcher;
import java.util.LinkedList;
import java.util.List;
public class EdgeSyncCursor {
List<EdgeEventFetcher> fetchers = new LinkedList<>();
int currentIdx = 0;
public EdgeSyncCursor(EdgeContextComponent ctx, Edge edge) {
fetchers.add(new SystemWidgetsBundlesEdgeEventFetcher(ctx.getWidgetsBundleService()));
fetchers.add(new TenantWidgetsBundlesEdgeEventFetcher(ctx.getWidgetsBundleService()));
fetchers.add(new DeviceProfilesEdgeEventFetcher(ctx.getDeviceProfileService()));
fetchers.add(new RuleChainsEdgeEventFetcher(ctx.getRuleChainService()));
fetchers.add(new TenantAdminUsersEdgeEventFetcher(ctx.getUserService()));
if (edge.getCustomerId() != null && !EntityId.NULL_UUID.equals(edge.getCustomerId().getId())) {
fetchers.add(new CustomerEdgeEventFetcher());
fetchers.add(new CustomerUsersEdgeEventFetcher(ctx.getUserService(), edge.getCustomerId()));
}
fetchers.add(new AdminSettingsEdgeEventFetcher(ctx.getAdminSettingsService()));
fetchers.add(new AssetsEdgeEventFetcher(ctx.getAssetService()));
fetchers.add(new DashboardsEdgeEventFetcher(ctx.getDashboardService()));
}
public boolean hasNext() {
return fetchers.size() > currentIdx;
}
public EdgeEventFetcher getNext() {
if (hasNext()) {
EdgeEventFetcher edgeEventFetcher = fetchers.get(currentIdx);
currentIdx++;
return edgeEventFetcher;
} else {
throw new IndexOutOfBoundsException();
}
}
public int getCurrentIdx() {
return currentIdx;
}
}

2
application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/AdminSettingsEdgeEventFetcher.java

@ -56,7 +56,7 @@ public class AdminSettingsEdgeEventFetcher implements EdgeEventFetcher {
@Override
public PageLink getPageLink(int pageSize) {
return new PageLink(pageSize);
return null;
}
@Override

49
application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/CustomerEdgeEventFetcher.java

@ -0,0 +1,49 @@
/**
* 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 lombok.AllArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.thingsboard.server.common.data.edge.Edge;
import org.thingsboard.server.common.data.edge.EdgeEvent;
import org.thingsboard.server.common.data.edge.EdgeEventActionType;
import org.thingsboard.server.common.data.edge.EdgeEventType;
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.service.edge.rpc.EdgeEventUtils;
import java.util.ArrayList;
import java.util.List;
@AllArgsConstructor
@Slf4j
public class CustomerEdgeEventFetcher implements EdgeEventFetcher {
@Override
public PageLink getPageLink(int pageSize) {
return null;
}
@Override
public PageData<EdgeEvent> fetchEdgeEvents(TenantId tenantId, Edge edge, PageLink pageLink) {
List<EdgeEvent> result = new ArrayList<>();
result.add(EdgeEventUtils.constructEdgeEvent(edge.getTenantId(), edge.getId(),
EdgeEventType.CUSTOMER, EdgeEventActionType.ADDED, edge.getCustomerId(), null));
// @voba - returns PageData object to be in sync with other fetchers
return new PageData<>(result, 1, result.size(), false);
}
}

33
application/src/main/java/org/thingsboard/server/service/executors/GrpcCallbackExecutorService.java

@ -0,0 +1,33 @@
/**
* 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.executors;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.stereotype.Component;
import org.thingsboard.common.util.AbstractListeningExecutor;
@Component
public class GrpcCallbackExecutorService extends AbstractListeningExecutor {
@Value("${edges.grpc_callback_thread_pool_size}")
private int grpcCallbackExecutorThreadPoolSize;
@Override
protected int getThreadPollSize() {
return grpcCallbackExecutorThreadPoolSize;
}
}

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

@ -707,7 +707,9 @@ edges:
max_read_records_count: "${EDGES_STORAGE_MAX_READ_RECORDS_COUNT:50}"
no_read_records_sleep: "${EDGES_NO_READ_RECORDS_SLEEP:1000}"
sleep_between_batches: "${EDGES_SLEEP_BETWEEN_BATCHES:1000}"
scheduler_pool_size: "${EDGES_SCHEDULER_POOL_SIZE:4}"
scheduler_pool_size: "${EDGES_SCHEDULER_POOL_SIZE:1}"
send_scheduler_pool_size: "${EDGES_SEND_SCHEDULER_POOL_SIZE:1}"
grpc_callback_thread_pool_size: "${EDGES_GRPC_CALLBACK_POOL_SIZE:1}"
edge_events_ttl: "${EDGES_EDGE_EVENTS_TTL:0}"
state:
persistToTelemetry: "${EDGES_PERSIST_STATE_TO_TELEMETRY:false}"

Loading…
Cancel
Save