未验证 提交 bf9a9019 编写于 作者: R Rajan Dhabalia 提交者: GitHub

[pulsar-client] Fix message corruption on OOM for batch messages (#5443)

* [pulsar-client] Fix message corruption on OOM for batch messages

* remove comments

* Address comments: index in local-var + remove lastSerializedMessageIndex var
上级 68236287
...@@ -82,11 +82,31 @@ class BatchMessageContainerImpl extends AbstractBatchMessageContainer { ...@@ -82,11 +82,31 @@ class BatchMessageContainerImpl extends AbstractBatchMessageContainer {
} }
private ByteBuf getCompressedBatchMetadataAndPayload() { private ByteBuf getCompressedBatchMetadataAndPayload() {
for (MessageImpl<?> msg : messages) { int batchWriteIndex = batchedMessageMetadataAndPayload.writerIndex();
int batchReadIndex = batchedMessageMetadataAndPayload.readerIndex();
for (int i = 0, n = messages.size(); i < n; i++) {
MessageImpl<?> msg = messages.get(i);
PulsarApi.MessageMetadata.Builder msgBuilder = msg.getMessageBuilder(); PulsarApi.MessageMetadata.Builder msgBuilder = msg.getMessageBuilder();
batchedMessageMetadataAndPayload = Commands.serializeSingleMessageInBatchWithPayload(msgBuilder, msg.getDataBuffer().markReaderIndex();
msg.getDataBuffer(), batchedMessageMetadataAndPayload); try {
msgBuilder.recycle(); batchedMessageMetadataAndPayload = Commands.serializeSingleMessageInBatchWithPayload(msgBuilder,
msg.getDataBuffer(), batchedMessageMetadataAndPayload);
} catch (Throwable th) {
// serializing batch message can corrupt the index of message and batch-message. Reset the index so,
// next iteration doesn't send corrupt message to broker.
for (int j = 0; j <= i; j++) {
MessageImpl<?> previousMsg = messages.get(j);
previousMsg.getDataBuffer().resetReaderIndex();
}
batchedMessageMetadataAndPayload.writerIndex(batchWriteIndex);
batchedMessageMetadataAndPayload.readerIndex(batchReadIndex);
throw new RuntimeException(th);
}
}
// Recycle messages only once they serialized successfully in batch
for (MessageImpl<?> msg : messages) {
msg.getMessageBuilder().recycle();
} }
int uncompressedSize = batchedMessageMetadataAndPayload.readableBytes(); int uncompressedSize = batchedMessageMetadataAndPayload.readableBytes();
ByteBuf compressedPayload = compressor.encode(batchedMessageMetadataAndPayload); ByteBuf compressedPayload = compressor.encode(batchedMessageMetadataAndPayload);
......
...@@ -1286,7 +1286,7 @@ public class ProducerImpl<T> extends ProducerBase<T> implements TimerTask, Conne ...@@ -1286,7 +1286,7 @@ public class ProducerImpl<T> extends ProducerBase<T> implements TimerTask, Conne
private void batchMessageAndSend() { private void batchMessageAndSend() {
if (log.isDebugEnabled()) { if (log.isDebugEnabled()) {
log.debug("[{}] [{}] Batching the messages from the batch container with {} messages", topic, producerName, log.debug("[{}] [{}] Batching the messages from the batch container with {} messages", topic, producerName,
batchMessageContainer.getNumMessagesInBatch()); batchMessageContainer.getNumMessagesInBatch());
} }
if (!batchMessageContainer.isEmpty()) { if (!batchMessageContainer.isEmpty()) {
try { try {
......
Markdown is supported
0% .
You are about to add 0 people to the discussion. Proceed with caution.
先完成此消息的编辑!
想要评论请 注册