From e7f29798ece70e218f7233a7ec85f01e8706a062 Mon Sep 17 00:00:00 2001 From: guyinyou <36399867+guyinyou@users.noreply.github.com> Date: Fri, 31 Mar 2023 17:03:37 +0800 Subject: [PATCH] [ISSUE #6518] Fix bug that multi-threaded using bytebuffer Co-authored-by: guyinyou --- .../java/org/apache/rocketmq/store/DefaultMessageStore.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/store/src/main/java/org/apache/rocketmq/store/DefaultMessageStore.java b/store/src/main/java/org/apache/rocketmq/store/DefaultMessageStore.java index 117f204817..dc8e3efdbc 100644 --- a/store/src/main/java/org/apache/rocketmq/store/DefaultMessageStore.java +++ b/store/src/main/java/org/apache/rocketmq/store/DefaultMessageStore.java @@ -2885,7 +2885,7 @@ public class DefaultMessageStore implements MessageStore { BatchDispatchRequest task = batchDispatchRequestQueue.peek(); batchDispatchRequestExecutor.execute(() -> { try { - ByteBuffer tmpByteBuffer = task.byteBuffer.duplicate(); + ByteBuffer tmpByteBuffer = task.byteBuffer; tmpByteBuffer.position(task.position); tmpByteBuffer.limit(task.position + task.size); List dispatchRequestList = new ArrayList<>(); @@ -3018,7 +3018,7 @@ public class DefaultMessageStore implements MessageStore { return; } mappedPageHoldCount.getAndIncrement(); - BatchDispatchRequest task = new BatchDispatchRequest(byteBuffer, position, size, batchId++); + BatchDispatchRequest task = new BatchDispatchRequest(byteBuffer.duplicate(), position, size, batchId++); batchDispatchRequestQueue.offer(task); }