From 58e418578c897b80306d6b2fe2057ec5579d8c2d Mon Sep 17 00:00:00 2001 From: Andrii Shvaika Date: Mon, 18 Jul 2022 15:53:02 +0300 Subject: [PATCH] Fix reordering of chunked entities from git --- .../DefaultGitVersionControlQueueService.java | 23 +++++++++++-------- common/cluster-api/src/main/proto/queue.proto | 8 +++---- .../DefaultClusterVersionControlService.java | 10 ++++---- 3 files changed, 22 insertions(+), 19 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/sync/vc/DefaultGitVersionControlQueueService.java b/application/src/main/java/org/thingsboard/server/service/sync/vc/DefaultGitVersionControlQueueService.java index 0c098d8281..413980d119 100644 --- a/application/src/main/java/org/thingsboard/server/service/sync/vc/DefaultGitVersionControlQueueService.java +++ b/application/src/main/java/org/thingsboard/server/service/sync/vc/DefaultGitVersionControlQueueService.java @@ -76,11 +76,13 @@ import org.thingsboard.server.service.sync.vc.data.VoidGitRequest; import java.util.ArrayList; import java.util.Collections; +import java.util.Comparator; import java.util.HashMap; import java.util.LinkedHashMap; import java.util.List; import java.util.Map; import java.util.Optional; +import java.util.TreeMap; import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.TimeUnit; @@ -102,7 +104,7 @@ public class DefaultGitVersionControlQueueService implements GitVersionControlQu private final SchedulerComponent scheduler; private final Map> pendingRequestMap = new HashMap<>(); - private final Map> chunkedMsgs = new ConcurrentHashMap<>(); + private final Map> chunkedMsgs = new ConcurrentHashMap<>(); @Value("${queue.vc.request-timeout:60000}") private int requestTimeout; @@ -286,7 +288,7 @@ public class DefaultGitVersionControlQueueService implements GitVersionControlQu @SuppressWarnings("rawtypes") public ListenableFuture getEntity(TenantId tenantId, String versionId, EntityId entityId) { EntityContentGitRequest request = new EntityContentGitRequest(tenantId, versionId, entityId); - chunkedMsgs.put(request.getRequestId(), new LinkedHashMap<>()); + chunkedMsgs.put(request.getRequestId(), new HashMap<>()); registerAndSend(request, builder -> builder.setEntityContentRequest(EntityContentRequestMsg.newBuilder() .setVersionId(versionId) .setEntityType(entityId.getEntityType().name()) @@ -328,7 +330,7 @@ public class DefaultGitVersionControlQueueService implements GitVersionControlQu @SuppressWarnings("rawtypes") public ListenableFuture> getEntities(TenantId tenantId, String versionId, EntityType entityType, int offset, int limit) { EntitiesContentGitRequest request = new EntitiesContentGitRequest(tenantId, versionId, entityType); - chunkedMsgs.put(request.getRequestId(), new LinkedHashMap<>()); + chunkedMsgs.put(request.getRequestId(), new HashMap<>()); registerAndSend(request, builder -> builder.setEntitiesContentRequest(EntitiesContentRequestMsg.newBuilder() .setVersionId(versionId) .setEntityType(entityType.name()) @@ -412,10 +414,10 @@ public class DefaultGitVersionControlQueueService implements GitVersionControlQu ((ListVersionsGitRequest) request).getFuture().set(toPageData(listVersionsResponse)); } else if (vcResponseMsg.hasEntityContentResponse()) { TransportProtos.EntityContentResponseMsg responseMsg = vcResponseMsg.getEntityContentResponse(); - log.trace("[{}] received chunk {} for 'getEntity'", responseMsg.getChunkedMsgId(), responseMsg.getChunkIndex()); - var joined = joinChunks(requestId, responseMsg, 1); + log.trace("Received chunk {} for 'getEntity'", responseMsg.getChunkIndex()); + var joined = joinChunks(requestId, responseMsg, 0, 1); if (joined.isPresent()) { - log.trace("[{}] collected all chunks for 'getEntity'", responseMsg.getChunkedMsgId()); + log.trace("Collected all chunks for 'getEntity'"); ((EntityContentGitRequest) request).getFuture().set(joined.get().get(0)); } else { completed = false; @@ -424,7 +426,7 @@ public class DefaultGitVersionControlQueueService implements GitVersionControlQu TransportProtos.EntitiesContentResponseMsg responseMsg = vcResponseMsg.getEntitiesContentResponse(); TransportProtos.EntityContentResponseMsg item = responseMsg.getItem(); if (responseMsg.getItemsCount() > 0) { - var joined = joinChunks(requestId, item, responseMsg.getItemsCount()); + var joined = joinChunks(requestId, item, responseMsg.getItemIdx(), responseMsg.getItemsCount()); if (joined.isPresent()) { ((EntitiesContentGitRequest) request).getFuture().set(joined.get()); } else { @@ -459,16 +461,17 @@ public class DefaultGitVersionControlQueueService implements GitVersionControlQu } @SuppressWarnings("rawtypes") - private Optional> joinChunks(UUID requestId, TransportProtos.EntityContentResponseMsg responseMsg, int expectedMsgCount) { + private Optional> joinChunks(UUID requestId, TransportProtos.EntityContentResponseMsg responseMsg, int itemIdx, int expectedMsgCount) { var chunksMap = chunkedMsgs.get(requestId); if (chunksMap == null) { return Optional.empty(); } - String[] msgChunks = chunksMap.computeIfAbsent(responseMsg.getChunkedMsgId(), id -> new String[responseMsg.getChunksCount()]); + String[] msgChunks = chunksMap.computeIfAbsent(itemIdx, id -> new String[responseMsg.getChunksCount()]); msgChunks[responseMsg.getChunkIndex()] = responseMsg.getData(); if (chunksMap.size() == expectedMsgCount && chunksMap.values().stream() .allMatch(chunks -> CollectionsUtil.countNonNull(chunks) == chunks.length)) { - return Optional.of(chunksMap.values().stream() + return Optional.of(chunksMap.entrySet().stream() + .sorted(Comparator.comparingInt(Map.Entry::getKey)).map(Map.Entry::getValue) .map(chunks -> String.join("", chunks)) .map(this::toData) .collect(Collectors.toList())); diff --git a/common/cluster-api/src/main/proto/queue.proto b/common/cluster-api/src/main/proto/queue.proto index 3995040080..1714e109a3 100644 --- a/common/cluster-api/src/main/proto/queue.proto +++ b/common/cluster-api/src/main/proto/queue.proto @@ -803,9 +803,8 @@ message EntityContentRequestMsg { message EntityContentResponseMsg { string data = 1; - string chunkedMsgId = 2; - int32 chunkIndex = 3; - int32 chunksCount = 4; + int32 chunkIndex = 2; + int32 chunksCount = 3; } message EntitiesContentRequestMsg { @@ -817,7 +816,8 @@ message EntitiesContentRequestMsg { message EntitiesContentResponseMsg { EntityContentResponseMsg item = 1; - int32 itemsCount = 2; + int32 itemIdx = 2; + int32 itemsCount = 3; } message VersionsDiffRequestMsg { diff --git a/common/version-control/src/main/java/org/thingsboard/server/service/sync/vc/DefaultClusterVersionControlService.java b/common/version-control/src/main/java/org/thingsboard/server/service/sync/vc/DefaultClusterVersionControlService.java index 4d642be521..eb7a3d20e4 100644 --- a/common/version-control/src/main/java/org/thingsboard/server/service/sync/vc/DefaultClusterVersionControlService.java +++ b/common/version-control/src/main/java/org/thingsboard/server/service/sync/vc/DefaultClusterVersionControlService.java @@ -274,20 +274,20 @@ public class DefaultClusterVersionControlService extends TbApplicationEventListe var ids = vcService.listEntitiesAtVersion(ctx.getTenantId(), request.getVersionId(), path) .stream().skip(request.getOffset()).limit(request.getLimit()).collect(Collectors.toList()); if (!ids.isEmpty()) { - for (VersionedEntityInfo info : ids) { + for (int i = 0; i < ids.size(); i++){ + VersionedEntityInfo info = ids.get(i); var data = vcService.getFileContentAtCommit(ctx.getTenantId(), getRelativePath(info.getExternalId().getEntityType(), info.getExternalId().getId().toString()), request.getVersionId()); - Iterable dataChunks = StringUtils.split(data, msgChunkSize); - String chunkedMsgId = UUID.randomUUID().toString(); int chunksCount = Iterables.size(dataChunks); AtomicInteger chunkIndex = new AtomicInteger(); + int itemIdx = i; dataChunks.forEach(chunk -> { EntitiesContentResponseMsg.Builder response = EntitiesContentResponseMsg.newBuilder() .setItemsCount(ids.size()) + .setItemIdx(itemIdx) .setItem(EntityContentResponseMsg.newBuilder() .setData(chunk) - .setChunkedMsgId(chunkedMsgId) .setChunksCount(chunksCount) .setChunkIndex(chunkIndex.getAndIncrement()) .build()); @@ -313,7 +313,7 @@ public class DefaultClusterVersionControlService extends TbApplicationEventListe dataChunks.forEach(chunk -> { log.trace("[{}] sending chunk {} for 'getEntity'", chunkedMsgId, chunkIndex.get()); reply(ctx, Optional.empty(), builder -> builder.setEntityContentResponse(EntityContentResponseMsg.newBuilder() - .setData(chunk).setChunkedMsgId(chunkedMsgId).setChunksCount(chunksCount) + .setData(chunk).setChunksCount(chunksCount) .setChunkIndex(chunkIndex.getAndIncrement()))); }); }