From 991af5e2042ddd0a5519b251a7b693d3bb4ef7ee Mon Sep 17 00:00:00 2001 From: zhouxiang Date: Fri, 18 Nov 2022 17:44:39 +0800 Subject: [PATCH] [ISSUE #5542] Fix ConsumerProcessor lockBatchMQ future allOf data race issue --- .../proxy/processor/ConsumerProcessor.java | 20 +++++-------------- 1 file changed, 5 insertions(+), 15 deletions(-) diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/processor/ConsumerProcessor.java b/proxy/src/main/java/org/apache/rocketmq/proxy/processor/ConsumerProcessor.java index f37238f693..5fec0cd3a5 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/processor/ConsumerProcessor.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/processor/ConsumerProcessor.java @@ -313,19 +313,15 @@ public class ConsumerProcessor extends AbstractProcessor { Set successSet = new CopyOnWriteArraySet<>(); Set addressableMessageQueueSet = buildAddressableSet(mqSet); Map> messageQueueSetMap = buildAddressableMapByBrokerName(addressableMessageQueueSet); - List>> futureList = new ArrayList<>(); + List> futureList = new ArrayList<>(); messageQueueSetMap.forEach((k, v) -> { LockBatchRequestBody requestBody = new LockBatchRequestBody(); requestBody.setConsumerGroup(consumerGroup); requestBody.setClientId(clientId); requestBody.setMqSet(v.stream().map(AddressableMessageQueue::getMessageQueue).collect(Collectors.toSet())); - CompletableFuture> future0 = new CompletableFuture<>(); - try { - future0 = serviceManager.getMessageService().lockBatchMQ(ctx, v.get(0), requestBody, timeoutMillis); - future0.thenAccept(successSet::addAll); - } catch (Throwable t) { - future0.completeExceptionally(t); - } + CompletableFuture future0 = serviceManager.getMessageService() + .lockBatchMQ(ctx, v.get(0), requestBody, timeoutMillis) + .thenAccept(successSet::addAll); futureList.add(FutureUtils.addExecutor(future0, this.executor)); }); CompletableFuture.allOf(futureList.toArray(new CompletableFuture[0])).whenComplete((v, t) -> { @@ -348,13 +344,7 @@ public class ConsumerProcessor extends AbstractProcessor { requestBody.setConsumerGroup(consumerGroup); requestBody.setClientId(clientId); requestBody.setMqSet(v.stream().map(AddressableMessageQueue::getMessageQueue).collect(Collectors.toSet())); - CompletableFuture future0 = new CompletableFuture<>(); - try { - future0 = serviceManager.getMessageService().unlockBatchMQ(ctx, v.get(0), requestBody, timeoutMillis); - future0.complete(null); - } catch (Throwable t) { - future0.completeExceptionally(t); - } + CompletableFuture future0 = serviceManager.getMessageService().unlockBatchMQ(ctx, v.get(0), requestBody, timeoutMillis); futureList.add(FutureUtils.addExecutor(future0, this.executor)); }); CompletableFuture.allOf(futureList.toArray(new CompletableFuture[0])).whenComplete((v, t) -> {