Skip to content

Commit fcf1555

Browse files
authored
[improve][client] Avoid recomputing entry-bucket hashes (#26383)
1 parent 0897616 commit fcf1555

3 files changed

Lines changed: 89 additions & 58 deletions

File tree

pulsar-client/src/main/java/org/apache/pulsar/client/impl/BatchMessageContainerImpl.java

Lines changed: 26 additions & 24 deletions
Original file line numberDiff line numberDiff line change
@@ -26,7 +26,6 @@
2626
import java.util.ArrayList;
2727
import java.util.Arrays;
2828
import java.util.List;
29-
import java.util.function.ToIntFunction;
3029
import lombok.CustomLog;
3130
import lombok.Getter;
3231
import lombok.Setter;
@@ -64,29 +63,17 @@ class BatchMessageContainerImpl extends AbstractBatchMessageContainer {
6463
// keep track of callbacks for individual messages being published in a batch
6564
protected SendCallback firstCallback;
6665

67-
// PIP-486: when set, createOpSendMsg() stamps entry_hash_min/max on the batch metadata from the
68-
// min/max entry-bucket hash of the batch's messages. Null for non-scalable producers (no stamping,
69-
// so the path is byte-identical to before).
70-
private ToIntFunction<MessageImpl<?>> entryBucketHashFn;
66+
// PIP-486: createOpSendMsg() stamps the min/max entry-bucket hash maintained by the three-argument add.
67+
private int minEntryBucketHash = Integer.MAX_VALUE;
68+
private int maxEntryBucketHash = Integer.MIN_VALUE;
7169

72-
void setEntryBucketHashFn(ToIntFunction<MessageImpl<?>> entryBucketHashFn) {
73-
this.entryBucketHashFn = entryBucketHashFn;
74-
}
75-
76-
/** PIP-486: stamp the effective entry-bucket hash range (min/max over the batch's messages). */
70+
/** PIP-486: stamp the entry-bucket hash range maintained while messages are added to the batch. */
7771
private void stampEntryBucketRange() {
78-
if (entryBucketHashFn == null || messages.isEmpty()) {
72+
if (minEntryBucketHash > maxEntryBucketHash) {
7973
return;
8074
}
81-
int min = Integer.MAX_VALUE;
82-
int max = Integer.MIN_VALUE;
83-
for (int i = 0; i < messages.size(); i++) {
84-
int h = entryBucketHashFn.applyAsInt(messages.get(i));
85-
min = Math.min(min, h);
86-
max = Math.max(max, h);
87-
}
88-
messageMetadata.setEntryHashMin(min);
89-
messageMetadata.setEntryHashMax(max);
75+
messageMetadata.setEntryHashMin(minEntryBucketHash);
76+
messageMetadata.setEntryHashMax(maxEntryBucketHash);
9077
}
9178

9279
protected final ByteBufAllocator allocator;
@@ -113,10 +100,10 @@ public BatchMessageContainerImpl(ProducerImpl<?> producer) {
113100

114101
@Override
115102
public boolean add(MessageImpl<?> msg, SendCallback callback) {
116-
log.debug().attr("topic", topicName)
117-
.attr("producerName", () -> producer != null ? producer.getProducerName() : null)
118-
.attr("numMessagesInBatch", numMessagesInBatch)
119-
.log("add message to batch");
103+
log.debug().attr("topic", topicName)
104+
.attr("producerName", () -> producer != null ? producer.getProducerName() : null)
105+
.attr("numMessagesInBatch", numMessagesInBatch)
106+
.log("add message to batch");
120107

121108
if (++numMessagesInBatch == 1) {
122109
try {
@@ -166,6 +153,19 @@ public boolean add(MessageImpl<?> msg, SendCallback callback) {
166153
return isBatchFull();
167154
}
168155

156+
/**
157+
* Adds a message whose entry-bucket hash has already been computed, avoiding a second hash
158+
* calculation when the send operation is created.
159+
*/
160+
boolean add(MessageImpl<?> msg, SendCallback callback, int entryBucketHash) {
161+
boolean isBatchFull = add(msg, callback);
162+
if (!isEmpty()) {
163+
minEntryBucketHash = Math.min(minEntryBucketHash, entryBucketHash);
164+
maxEntryBucketHash = Math.max(maxEntryBucketHash, entryBucketHash);
165+
}
166+
return isBatchFull;
167+
}
168+
169169
protected ByteBuf getCompressedBatchMetadataAndPayload() {
170170
return getCompressedBatchMetadataAndPayload(true);
171171
}
@@ -252,6 +252,8 @@ public void clear() {
252252
currentBatchSizeBytes = 0;
253253
lowestSequenceId = -1L;
254254
highestSequenceId = -1L;
255+
minEntryBucketHash = Integer.MAX_VALUE;
256+
maxEntryBucketHash = Integer.MIN_VALUE;
255257
batchedMessageMetadataAndPayload = null;
256258
currentTxnidMostBits = -1L;
257259
currentTxnidLeastBits = -1L;

pulsar-client/src/main/java/org/apache/pulsar/client/impl/EntryBucketBatchContainer.java

Lines changed: 5 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -54,13 +54,11 @@ class EntryBucketBatchContainer extends AbstractBatchMessageContainer {
5454

5555
@Override
5656
public boolean add(MessageImpl<?> msg, SendCallback callback) {
57-
int bucket = bucketOf(splits, entryBucketHash(msg));
58-
final BatchMessageContainerImpl batchMessageContainer = batches.computeIfAbsent(bucket, __ -> {
59-
BatchMessageContainerImpl c = new BatchMessageContainerImpl(producer);
60-
c.setEntryBucketHashFn(this::entryBucketHash);
61-
return c;
62-
});
63-
batchMessageContainer.add(msg, callback);
57+
int hashCode = entryBucketHash(msg);
58+
int bucket = bucketOf(splits, hashCode);
59+
final BatchMessageContainerImpl batchMessageContainer =
60+
batches.computeIfAbsent(bucket, __ -> new BatchMessageContainerImpl(producer));
61+
batchMessageContainer.add(msg, callback, hashCode);
6462
// `add` fails iff the container was empty and `msg` (the first message) failed; then the
6563
// container is cleared and there is nothing to count.
6664
if (!batchMessageContainer.isEmpty()) {

pulsar-client/src/test/java/org/apache/pulsar/client/impl/BatchMessageContainerImplTest.java

Lines changed: 58 additions & 27 deletions
Original file line numberDiff line numberDiff line change
@@ -145,27 +145,7 @@ public void recoveryAfterOom() {
145145

146146
@Test
147147
public void testMessagesSize() throws Exception {
148-
ProducerImpl<?> producer = mock(ProducerImpl.class);
149-
150-
final ProducerConfigurationData producerConfigurationData = new ProducerConfigurationData();
151-
producerConfigurationData.setCompressionType(CompressionType.NONE);
152-
PulsarClientImpl pulsarClient = mock(PulsarClientImpl.class);
153-
ConnectionPool connectionPool = mock(ConnectionPool.class);
154-
when(pulsarClient.getCnxPool()).thenReturn(connectionPool);
155-
MemoryLimitController memoryLimitController = mock(MemoryLimitController.class);
156-
when(pulsarClient.getMemoryLimitController()).thenReturn(memoryLimitController);
157-
try {
158-
Field clientFiled = HandlerState.class.getDeclaredField("client");
159-
clientFiled.setAccessible(true);
160-
clientFiled.set(producer, pulsarClient);
161-
} catch (Exception e){
162-
fail(e.getMessage());
163-
}
164-
165-
ByteBuffer payload = ByteBuffer.wrap("payload".getBytes(StandardCharsets.UTF_8));
166-
167-
when(producer.getConfiguration()).thenReturn(producerConfigurationData);
168-
when(producer.encryptMessage(any(), any())).thenReturn(ByteBufAllocator.DEFAULT.buffer().writeBytes(payload));
148+
ProducerImpl<?> producer = createTestProducer();
169149

170150
final int initNum = 32;
171151
BatchMessageContainerImpl batchMessageContainer = new BatchMessageContainerImpl(producer);
@@ -181,16 +161,67 @@ public void testMessagesSize() throws Exception {
181161
assertEquals(batchMessageContainer.getMaxMessagesNum(), 200);
182162
}
183163

164+
@Test
165+
public void testEntryBucketHashRangeIsStampedWhenCreatingSendOperation() throws Exception {
166+
BatchMessageContainerImpl batchMessageContainer = new BatchMessageContainerImpl(createTestProducer());
167+
ArrayList<MessageImpl<?>> messages = new ArrayList<>();
168+
try {
169+
MessageImpl<?> singleMessage = createMessage(1);
170+
messages.add(singleMessage);
171+
batchMessageContainer.add(singleMessage, null, 0x3000);
172+
batchMessageContainer.createOpSendMsg();
173+
assertEquals(batchMessageContainer.messageMetadata.getEntryHashMin(), 0x3000);
174+
assertEquals(batchMessageContainer.messageMetadata.getEntryHashMax(), 0x3000);
175+
176+
batchMessageContainer.clear();
177+
MessageImpl<?> firstMessage = createMessage(2);
178+
MessageImpl<?> secondMessage = createMessage(3);
179+
MessageImpl<?> thirdMessage = createMessage(4);
180+
messages.add(firstMessage);
181+
messages.add(secondMessage);
182+
messages.add(thirdMessage);
183+
batchMessageContainer.add(firstMessage, null, 0x2000);
184+
batchMessageContainer.add(secondMessage, null, 0x1000);
185+
batchMessageContainer.add(thirdMessage, null, 0x1800);
186+
batchMessageContainer.createOpSendMsg();
187+
assertEquals(batchMessageContainer.messageMetadata.getEntryHashMin(), 0x1000);
188+
assertEquals(batchMessageContainer.messageMetadata.getEntryHashMax(), 0x2000);
189+
} finally {
190+
batchMessageContainer.discard(null);
191+
messages.forEach(ReferenceCountUtil::safeRelease);
192+
}
193+
}
194+
195+
private MessageImpl<?> createMessage(long sequenceId) {
196+
MessageMetadata messageMetadata = new MessageMetadata();
197+
messageMetadata.setSequenceId(sequenceId);
198+
messageMetadata.setProducerName("producer");
199+
messageMetadata.setPublishTime(System.currentTimeMillis());
200+
ByteBuffer payload = ByteBuffer.wrap("payload".getBytes(StandardCharsets.UTF_8));
201+
return MessageImpl.create(messageMetadata, payload, Schema.BYTES, null);
202+
}
203+
204+
private ProducerImpl<?> createTestProducer() throws Exception {
205+
ProducerImpl<?> producer = mock(ProducerImpl.class);
206+
ProducerConfigurationData producerConfigurationData = new ProducerConfigurationData();
207+
producerConfigurationData.setCompressionType(CompressionType.NONE);
208+
PulsarClientImpl pulsarClient = mock(PulsarClientImpl.class);
209+
when(pulsarClient.getCnxPool()).thenReturn(mock(ConnectionPool.class));
210+
when(pulsarClient.getMemoryLimitController()).thenReturn(mock(MemoryLimitController.class));
211+
Field clientField = HandlerState.class.getDeclaredField("client");
212+
clientField.setAccessible(true);
213+
clientField.set(producer, pulsarClient);
214+
when(producer.getConfiguration()).thenReturn(producerConfigurationData);
215+
when(producer.encryptMessage(any(), any())).thenAnswer(__ -> ByteBufAllocator.DEFAULT.buffer()
216+
.writeBytes("payload".getBytes(StandardCharsets.UTF_8)));
217+
return producer;
218+
}
219+
184220
private void addMessagesAndCreateOpSendMsg(BatchMessageContainerImpl batchMessageContainer, int num)
185221
throws Exception{
186222
ArrayList<MessageImpl<?>> messages = new ArrayList<>();
187223
for (int i = 0; i < num; ++i) {
188-
MessageMetadata messageMetadata = new MessageMetadata();
189-
messageMetadata.setSequenceId(i);
190-
messageMetadata.setProducerName("producer");
191-
messageMetadata.setPublishTime(System.currentTimeMillis());
192-
ByteBuffer payload = ByteBuffer.wrap("payload".getBytes(StandardCharsets.UTF_8));
193-
MessageImpl<?> message = MessageImpl.create(messageMetadata, payload, Schema.BYTES, null);
224+
MessageImpl<?> message = createMessage(i);
194225
messages.add(message);
195226
batchMessageContainer.add(message, null);
196227
}

0 commit comments

Comments
 (0)