Browse Source

Merge pull request #7651 from volodymyr-babak/bug/fix-edge-fetcher

[3.4.2][Bug] Updates to stability of synchronization between edge and cloud in case of many events simultaneously
pull/7687/head
Andrew Shvayka 4 years ago
committed by GitHub
parent
commit
c28526fb25
No known key found for this signature in database GPG Key ID: 4AEE18F83AFDEB23
  1. 50
      application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java
  2. 3
      application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeSessionState.java
  3. 2
      application/src/main/resources/thingsboard.yml
  4. 12
      application/src/test/java/org/thingsboard/server/edge/BaseDeviceEdgeTest.java
  5. 2
      application/src/test/java/org/thingsboard/server/edge/BaseTelemetryEdgeTest.java
  6. 6
      application/src/test/java/org/thingsboard/server/edge/imitator/EdgeImitator.java

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

@ -187,10 +187,7 @@ public final class EdgeGrpcSession implements Closeable {
public void startSyncProcess(TenantId tenantId, EdgeId edgeId, boolean fullSync) { public void startSyncProcess(TenantId tenantId, EdgeId edgeId, boolean fullSync) {
log.trace("[{}][{}] Staring edge sync process", tenantId, edgeId); log.trace("[{}][{}] Staring edge sync process", tenantId, edgeId);
syncCompleted = false; syncCompleted = false;
if (sessionState.getSendDownlinkMsgsFuture() != null && sessionState.getSendDownlinkMsgsFuture().isDone()) { interruptGeneralProcessingOnSync(tenantId, edgeId);
String errorMsg = String.format("[%s][%s] Sync process started. General processing interrupted!", tenantId, edgeId);
sessionState.getSendDownlinkMsgsFuture().setException(new RuntimeException(errorMsg));
}
doSync(new EdgeSyncCursor(ctx, edge, fullSync)); doSync(new EdgeSyncCursor(ctx, edge, fullSync));
} }
@ -265,10 +262,7 @@ public final class EdgeGrpcSession implements Closeable {
} }
if (sessionState.getPendingMsgsMap().isEmpty()) { if (sessionState.getPendingMsgsMap().isEmpty()) {
log.debug("[{}] Pending msgs map is empty. Stopping current iteration", edge.getRoutingKey()); log.debug("[{}] Pending msgs map is empty. Stopping current iteration", edge.getRoutingKey());
if (sessionState.getScheduledSendDownlinkTask() != null) { stopCurrentSendDownlinkMsgsTask(null);
sessionState.getScheduledSendDownlinkTask().cancel(false);
}
sessionState.getSendDownlinkMsgsFuture().set(null);
} }
} catch (Exception e) { } catch (Exception e) {
log.error("[{}] Can't process downlink response message [{}]", this.sessionId, msg, e); log.error("[{}] Can't process downlink response message [{}]", this.sessionId, msg, e);
@ -391,15 +385,14 @@ public final class EdgeGrpcSession implements Closeable {
} }
private ListenableFuture<Void> sendDownlinkMsgsPack(List<DownlinkMsg> downlinkMsgsPack) { private ListenableFuture<Void> sendDownlinkMsgsPack(List<DownlinkMsg> downlinkMsgsPack) {
if (sessionState.getSendDownlinkMsgsFuture() != null && !sessionState.getSendDownlinkMsgsFuture().isDone()) { interruptPreviousSendDownlinkMsgsTask();
String errorMsg = "[" + this.sessionId + "] Previous send downlink future was not properly completed, stopping it now";
log.error(errorMsg);
sessionState.getSendDownlinkMsgsFuture().setException(new RuntimeException(errorMsg));
}
sessionState.setSendDownlinkMsgsFuture(SettableFuture.create()); sessionState.setSendDownlinkMsgsFuture(SettableFuture.create());
sessionState.getPendingMsgsMap().clear(); sessionState.getPendingMsgsMap().clear();
downlinkMsgsPack.forEach(msg -> sessionState.getPendingMsgsMap().put(msg.getDownlinkMsgId(), msg)); downlinkMsgsPack.forEach(msg -> sessionState.getPendingMsgsMap().put(msg.getDownlinkMsgId(), msg));
scheduleDownlinkMsgsPackSend(1); scheduleDownlinkMsgsPackSend(1);
return sessionState.getSendDownlinkMsgsFuture(); return sessionState.getSendDownlinkMsgsFuture();
} }
@ -422,13 +415,13 @@ public final class EdgeGrpcSession implements Closeable {
} else { } else {
log.warn("[{}] Failed to deliver the batch after {} attempts. Next messages are going to be discarded {}", log.warn("[{}] Failed to deliver the batch after {} attempts. Next messages are going to be discarded {}",
this.sessionId, MAX_DOWNLINK_ATTEMPTS, copy); this.sessionId, MAX_DOWNLINK_ATTEMPTS, copy);
sessionState.getSendDownlinkMsgsFuture().set(null); stopCurrentSendDownlinkMsgsTask(null);
} }
} else { } else {
sessionState.getSendDownlinkMsgsFuture().set(null); stopCurrentSendDownlinkMsgsTask(null);
} }
} catch (Exception e) { } catch (Exception e) {
sessionState.getSendDownlinkMsgsFuture().setException(e); stopCurrentSendDownlinkMsgsTask(e);
} }
}; };
@ -673,4 +666,29 @@ public final class EdgeGrpcSession implements Closeable {
log.debug("[{}] Failed to close output stream: {}", sessionId, e.getMessage()); log.debug("[{}] Failed to close output stream: {}", sessionId, e.getMessage());
} }
} }
private void interruptPreviousSendDownlinkMsgsTask() {
String msg = String.format("[%s] Previous send downlink future was not properly completed, stopping it now!", this.sessionId);
stopCurrentSendDownlinkMsgsTask(new RuntimeException(msg));
}
private void interruptGeneralProcessingOnSync(TenantId tenantId, EdgeId edgeId) {
String msg = String.format("[%s][%s] Sync process started. General processing interrupted!", tenantId, edgeId);
stopCurrentSendDownlinkMsgsTask(new RuntimeException(msg));
}
public void stopCurrentSendDownlinkMsgsTask(Exception e) {
if (sessionState.getSendDownlinkMsgsFuture() != null && !sessionState.getSendDownlinkMsgsFuture().isDone()) {
if (e != null) {
log.warn(e.getMessage(), e);
sessionState.getSendDownlinkMsgsFuture().setException(e);
} else {
sessionState.getSendDownlinkMsgsFuture().set(null);
}
}
if (sessionState.getScheduledSendDownlinkTask() != null) {
sessionState.getScheduledSendDownlinkTask().cancel(true);
}
}
} }

3
application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeSessionState.java

@ -19,6 +19,7 @@ import com.google.common.util.concurrent.SettableFuture;
import lombok.Data; import lombok.Data;
import org.thingsboard.server.gen.edge.v1.DownlinkMsg; import org.thingsboard.server.gen.edge.v1.DownlinkMsg;
import java.util.Collections;
import java.util.LinkedHashMap; import java.util.LinkedHashMap;
import java.util.Map; import java.util.Map;
import java.util.concurrent.ScheduledFuture; import java.util.concurrent.ScheduledFuture;
@ -26,7 +27,7 @@ import java.util.concurrent.ScheduledFuture;
@Data @Data
public class EdgeSessionState { public class EdgeSessionState {
private final Map<Integer, DownlinkMsg> pendingMsgsMap = new LinkedHashMap<>(); private final Map<Integer, DownlinkMsg> pendingMsgsMap = Collections.synchronizedMap(new LinkedHashMap<>());
private SettableFuture<Void> sendDownlinkMsgsFuture; private SettableFuture<Void> sendDownlinkMsgsFuture;
private ScheduledFuture<?> scheduledSendDownlinkTask; private ScheduledFuture<?> scheduledSendDownlinkTask;
} }

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

@ -930,7 +930,7 @@ edges:
storage: storage:
max_read_records_count: "${EDGES_STORAGE_MAX_READ_RECORDS_COUNT:50}" max_read_records_count: "${EDGES_STORAGE_MAX_READ_RECORDS_COUNT:50}"
no_read_records_sleep: "${EDGES_NO_READ_RECORDS_SLEEP:1000}" no_read_records_sleep: "${EDGES_NO_READ_RECORDS_SLEEP:1000}"
sleep_between_batches: "${EDGES_SLEEP_BETWEEN_BATCHES:1000}" sleep_between_batches: "${EDGES_SLEEP_BETWEEN_BATCHES:10000}"
scheduler_pool_size: "${EDGES_SCHEDULER_POOL_SIZE:1}" scheduler_pool_size: "${EDGES_SCHEDULER_POOL_SIZE:1}"
send_scheduler_pool_size: "${EDGES_SEND_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}" grpc_callback_thread_pool_size: "${EDGES_GRPC_CALLBACK_POOL_SIZE:1}"

12
application/src/test/java/org/thingsboard/server/edge/BaseDeviceEdgeTest.java

@ -439,9 +439,9 @@ abstract public class BaseDeviceEdgeTest extends AbstractEdgeTest {
Assert.assertTrue(edgeImitator.waitForResponses()); Assert.assertTrue(edgeImitator.waitForResponses());
Assert.assertTrue(edgeImitator.waitForMessages()); Assert.assertTrue(edgeImitator.waitForMessages());
AbstractMessage latestMessage = edgeImitator.getMessageFromTail(2); Optional<DeviceUpdateMsg> deviceUpdateMsgOpt = edgeImitator.findMessageByType(DeviceUpdateMsg.class);
Assert.assertTrue(latestMessage instanceof DeviceUpdateMsg); Assert.assertTrue(deviceUpdateMsgOpt.isPresent());
DeviceUpdateMsg latestDeviceUpdateMsg = (DeviceUpdateMsg) latestMessage; DeviceUpdateMsg latestDeviceUpdateMsg = deviceUpdateMsgOpt.get();
Assert.assertNotEquals(deviceOnCloudName, latestDeviceUpdateMsg.getName()); Assert.assertNotEquals(deviceOnCloudName, latestDeviceUpdateMsg.getName());
Assert.assertEquals(deviceOnCloudName, latestDeviceUpdateMsg.getConflictName()); Assert.assertEquals(deviceOnCloudName, latestDeviceUpdateMsg.getConflictName());
@ -453,9 +453,9 @@ abstract public class BaseDeviceEdgeTest extends AbstractEdgeTest {
Assert.assertNotNull(device); Assert.assertNotNull(device);
Assert.assertNotEquals(deviceOnCloudName, device.getName()); Assert.assertNotEquals(deviceOnCloudName, device.getName());
latestMessage = edgeImitator.getLatestMessage(); Optional<DeviceCredentialsRequestMsg> deviceCredentialsUpdateMsgOpt = edgeImitator.findMessageByType(DeviceCredentialsRequestMsg.class);
Assert.assertTrue(latestMessage instanceof DeviceCredentialsRequestMsg); Assert.assertTrue(deviceCredentialsUpdateMsgOpt.isPresent());
DeviceCredentialsRequestMsg latestDeviceCredentialsRequestMsg = (DeviceCredentialsRequestMsg) latestMessage; DeviceCredentialsRequestMsg latestDeviceCredentialsRequestMsg = deviceCredentialsUpdateMsgOpt.get();
Assert.assertEquals(uuid.getMostSignificantBits(), latestDeviceCredentialsRequestMsg.getDeviceIdMSB()); Assert.assertEquals(uuid.getMostSignificantBits(), latestDeviceCredentialsRequestMsg.getDeviceIdMSB());
Assert.assertEquals(uuid.getLeastSignificantBits(), latestDeviceCredentialsRequestMsg.getDeviceIdLSB()); Assert.assertEquals(uuid.getLeastSignificantBits(), latestDeviceCredentialsRequestMsg.getDeviceIdLSB());

2
application/src/test/java/org/thingsboard/server/edge/BaseTelemetryEdgeTest.java

@ -34,7 +34,7 @@ abstract public class BaseTelemetryEdgeTest extends AbstractEdgeTest {
@Test @Test
public void testTimeseriesWithFailures() throws Exception { public void testTimeseriesWithFailures() throws Exception {
int numberOfTimeseriesToSend = 1000; int numberOfTimeseriesToSend = 333;
Device device = findDeviceByName("Edge Device 1"); Device device = findDeviceByName("Edge Device 1");

6
application/src/test/java/org/thingsboard/server/edge/imitator/EdgeImitator.java

@ -365,11 +365,7 @@ public class EdgeImitator {
} }
public AbstractMessage getLatestMessage() { public AbstractMessage getLatestMessage() {
return getMessageFromTail(1); return downlinkMsgs.get(downlinkMsgs.size() - 1);
}
public AbstractMessage getMessageFromTail(int offset) {
return downlinkMsgs.get(downlinkMsgs.size() - offset);
} }
public void ignoreType(Class<? extends AbstractMessage> type) { public void ignoreType(Class<? extends AbstractMessage> type) {

Loading…
Cancel
Save