mirror of
https://github.com/apache/rocketmq.git
synced 2026-09-24 16:04:00 +08:00
Combine the process of decoding byteBuffer into preCheckMessageAndReturnSize method
This commit is contained in:
@@ -3020,28 +3020,7 @@ public class DefaultMessageStore implements MessageStore {
|
||||
for (int readSize = 0; readSize < result.getSize() && reputFromOffset < DefaultMessageStore.this.getConfirmOffset() && doNext; ) {
|
||||
ByteBuffer byteBuffer = result.getByteBuffer();
|
||||
|
||||
byteBuffer.mark();
|
||||
|
||||
int totalSize = byteBuffer.getInt();
|
||||
if (reputFromOffset + totalSize > DefaultMessageStore.this.getConfirmOffset()) {
|
||||
doNext = false;
|
||||
break;
|
||||
}
|
||||
|
||||
int magicCode = byteBuffer.getInt();
|
||||
switch (magicCode) {
|
||||
case MessageDecoder.MESSAGE_MAGIC_CODE:
|
||||
case MessageDecoder.MESSAGE_MAGIC_CODE_V2:
|
||||
break;
|
||||
case MessageDecoder.BLANK_MAGIC_CODE:
|
||||
totalSize = 0;
|
||||
break;
|
||||
default:
|
||||
totalSize = -1;
|
||||
doNext = false;
|
||||
}
|
||||
|
||||
byteBuffer.reset();
|
||||
int totalSize = preCheckMessageAndReturnSize(byteBuffer);
|
||||
|
||||
if (totalSize > 0) {
|
||||
if (batchDispatchRequestStart == -1) {
|
||||
@@ -3058,9 +3037,9 @@ public class DefaultMessageStore implements MessageStore {
|
||||
this.reputFromOffset += totalSize;
|
||||
readSize += totalSize;
|
||||
} else {
|
||||
doNext = false;
|
||||
if (totalSize == 0) {
|
||||
this.reputFromOffset = DefaultMessageStore.this.commitLog.rollNextFile(this.reputFromOffset);
|
||||
readSize = result.getSize();
|
||||
}
|
||||
this.createBatchDispatchRequest(byteBuffer, batchDispatchRequestStart, batchDispatchRequestSize);
|
||||
batchDispatchRequestStart = -1;
|
||||
@@ -3083,6 +3062,35 @@ public class DefaultMessageStore implements MessageStore {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* pre-check the message and returns the message size
|
||||
*
|
||||
* @return 0 Come to the end of file // >0 Normal messages // -1 Message checksum failure
|
||||
*/
|
||||
public int preCheckMessageAndReturnSize(ByteBuffer byteBuffer) {
|
||||
byteBuffer.mark();
|
||||
|
||||
int totalSize = byteBuffer.getInt();
|
||||
if (reputFromOffset + totalSize > DefaultMessageStore.this.getConfirmOffset()) {
|
||||
return -1;
|
||||
}
|
||||
|
||||
int magicCode = byteBuffer.getInt();
|
||||
switch (magicCode) {
|
||||
case MessageDecoder.MESSAGE_MAGIC_CODE:
|
||||
case MessageDecoder.MESSAGE_MAGIC_CODE_V2:
|
||||
break;
|
||||
case MessageDecoder.BLANK_MAGIC_CODE:
|
||||
return 0;
|
||||
default:
|
||||
return -1;
|
||||
}
|
||||
|
||||
byteBuffer.reset();
|
||||
|
||||
return totalSize;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void shutdown() {
|
||||
for (int i = 0; i < 50 && this.isCommitLogAvailable(); i++) {
|
||||
|
||||
Reference in New Issue
Block a user