Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -53,6 +53,7 @@
import org.apache.pulsar.client.api.ProducerConsumerBase;
import org.apache.pulsar.client.api.PulsarClientException;
import org.apache.pulsar.client.api.Reader;
import org.apache.pulsar.client.api.Schema;
import org.apache.pulsar.client.api.SizeUnit;
import org.apache.pulsar.client.impl.MessageImpl.SchemaState;
import org.apache.pulsar.client.impl.ProducerImpl.OpSendMsg;
Expand Down Expand Up @@ -550,6 +551,28 @@ public void testReaderChunkingConfiguration() throws Exception {
assertEquals(consumer.conf.getExpireTimeOfIncompleteChunkedMessageMillis(), 12);
}

@Test
public void testChunkSize() throws Exception {
final int maxMessageSize = 50;
final int payloadChunkSize = maxMessageSize - 32/* the default message metadata size for string schema */;
this.conf.setMaxMessageSize(maxMessageSize);

final Producer<String> producer = pulsarClient.newProducer(Schema.STRING)
.topic("my-property/my-ns/test-chunk-size")
.enableChunking(true)
.enableBatching(false)
.create();
for (int size = 1; size <= maxMessageSize; size++) {
final MessageId messageId = producer.send(createMessagePayload(size));
log.info("Send {} bytes to {}", size, messageId);
if (size <= payloadChunkSize) {
assertEquals(messageId.getClass(), MessageIdImpl.class);
} else {
assertEquals(messageId.getClass(), ChunkMessageIdImpl.class);
}
}
}

private String createMessagePayload(int size) {
StringBuilder str = new StringBuilder();
Random rand = new Random();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -441,7 +441,7 @@ public void sendAsync(Message<?> message, SendCallback callback) {
MessageImpl<?> msg = (MessageImpl<?>) message;
MessageMetadata msgMetadata = msg.getMessageBuilder();
ByteBuf payload = msg.getDataBuffer();
int uncompressedSize = payload.readableBytes();
final int uncompressedSize = payload.readableBytes();

if (!canEnqueueRequest(callback, message.getSequenceId(), uncompressedSize)) {
return;
Expand Down Expand Up @@ -488,6 +488,10 @@ public void sendAsync(Message<?> message, SendCallback callback) {
return;
}

// Update the message metadata before computing the payload chunk size to avoid a large message cannot be split
// into chunks.
final long sequenceId = updateMessageMetadata(msgMetadata, uncompressedSize);
Comment thread
BewareMyPower marked this conversation as resolved.

// send in chunks
int totalChunks;
int payloadChunkSize;
Expand Down Expand Up @@ -526,13 +530,6 @@ public void sendAsync(Message<?> message, SendCallback callback) {
try {
synchronized (this) {
int readStartIndex = 0;
long sequenceId;
if (!msgMetadata.hasSequenceId()) {
sequenceId = msgIdGeneratorUpdater.getAndIncrement(this);
msgMetadata.setSequenceId(sequenceId);
} else {
sequenceId = msgMetadata.getSequenceId();
}
String uuid = totalChunks > 1 ? String.format("%s-%d", producerName, sequenceId) : null;
ChunkedMessageCtx chunkedMessageCtx = totalChunks > 1 ? ChunkedMessageCtx.get(totalChunks) : null;
byte[] schemaVersion = totalChunks > 1 && msg.getMessageBuilder().hasSchemaVersion()
Expand All @@ -554,7 +551,7 @@ public void sendAsync(Message<?> message, SendCallback callback) {
}
serializeAndSendMessage(msg, payload, sequenceId, uuid, chunkId, totalChunks,
readStartIndex, payloadChunkSize, compressedPayload, compressed,
compressedPayload.readableBytes(), uncompressedSize, callback, chunkedMessageCtx);
compressedPayload.readableBytes(), callback, chunkedMessageCtx);
readStartIndex = ((chunkId + 1) * payloadChunkSize);
}
}
Expand All @@ -567,6 +564,38 @@ public void sendAsync(Message<?> message, SendCallback callback) {
}
}

/**
* Update the message metadata except those fields that will be updated for chunks later.
*
* @param msgMetadata
* @param uncompressedSize
* @return the sequence id
*/
private long updateMessageMetadata(final MessageMetadata msgMetadata, final int uncompressedSize) {
final long sequenceId;
if (!msgMetadata.hasSequenceId()) {
sequenceId = msgIdGeneratorUpdater.getAndIncrement(this);
msgMetadata.setSequenceId(sequenceId);
} else {
sequenceId = msgMetadata.getSequenceId();
}

if (!msgMetadata.hasPublishTime()) {
msgMetadata.setPublishTime(client.getClientClock().millis());

checkArgument(!msgMetadata.hasProducerName());

msgMetadata.setProducerName(producerName);

if (conf.getCompressionType() != CompressionType.NONE) {
msgMetadata
.setCompression(CompressionCodecProvider.convertToWireProtocol(conf.getCompressionType()));
}
msgMetadata.setUncompressedSize(uncompressedSize);
}
return sequenceId;
}

@Override
public int getNumOfPartitions() {
return 0;
Expand All @@ -583,7 +612,6 @@ private void serializeAndSendMessage(MessageImpl<?> msg,
ByteBuf compressedPayload,
boolean compressed,
int compressedPayloadSize,
int uncompressedSize,
SendCallback callback,
ChunkedMessageCtx chunkedMessageCtx) throws IOException {
ByteBuf chunkPayload = compressedPayload;
Expand All @@ -603,19 +631,6 @@ private void serializeAndSendMessage(MessageImpl<?> msg,
.setNumChunksFromMsg(totalChunks)
.setTotalChunkMsgSize(compressedPayloadSize);
}
if (!msgMetadata.hasPublishTime()) {
msgMetadata.setPublishTime(client.getClientClock().millis());

checkArgument(!msgMetadata.hasProducerName());

msgMetadata.setProducerName(producerName);

if (conf.getCompressionType() != CompressionType.NONE) {
msgMetadata
.setCompression(CompressionCodecProvider.convertToWireProtocol(conf.getCompressionType()));
}
msgMetadata.setUncompressedSize(uncompressedSize);
}

if (canAddToBatch(msg) && totalChunks <= 1) {
if (canAddToCurrentBatch(msg)) {
Expand Down Expand Up @@ -1483,7 +1498,11 @@ public int getMessageHeaderAndPayloadSize() {
cmdHeader.markReaderIndex();
int totalSize = cmdHeader.readInt();
int cmdSize = cmdHeader.readInt();
int msgHeadersAndPayloadSize = totalSize - cmdSize - 4;
// The totalSize includes:
// | cmdLength | cmdSize | magic and checksum | msgMetadataLength | msgMetadata |
// | --------- | ------- | ------------------ | ----------------- | ----------- |
// | 4 | | 6 | 4 | |
int msgHeadersAndPayloadSize = totalSize - 4 - cmdSize - 6 - 4;
cmdHeader.resetReaderIndex();
return msgHeadersAndPayloadSize;
}
Expand Down Expand Up @@ -2214,8 +2233,8 @@ private void recoverProcessOpSendMsgFrom(ClientCnx cnx, MessageImpl from, long e
/**
* Check if final message size for non-batch and non-chunked messages is larger than max message size.
*/
public boolean isMessageSizeExceeded(OpSendMsg op) {
if (op.msg != null && op.totalChunks <= 1) {
private boolean isMessageSizeExceeded(OpSendMsg op) {
if (op.msg != null && !conf.isChunkingEnabled()) {
int messageSize = op.getMessageHeaderAndPayloadSize();
if (messageSize > ClientCnx.getMaxMessageSize()) {
releaseSemaphoreForSendOp(op);
Expand Down