Skip to content

Commit c52ab25

Browse files
committed
fix test
1 parent 1ea228e commit c52ab25

6 files changed

Lines changed: 22 additions & 5 deletions

File tree

pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Consumer.java

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -231,7 +231,7 @@ public ChannelPromise sendMessages(final List<Entry> entries, EntryBatchSizes ba
231231
msgOut.recordMultipleEvents(totalMessages, totalBytes);
232232
msgOutCounter.add(totalMessages);
233233
bytesOutCounter.add(totalBytes);
234-
chuckedMessageRate.recordEvent(totalChunkedMessages);
234+
chuckedMessageRate.recordMultipleEvents(totalChunkedMessages, 0);
235235

236236
ctx.channel().eventLoop().execute(() -> {
237237
for (int i = 0; i < entries.size(); i++) {
@@ -465,7 +465,7 @@ public void updateRates() {
465465
stats.msgRateOut = msgOut.getRate();
466466
stats.msgThroughputOut = msgOut.getValueRate();
467467
stats.msgRateRedeliver = msgRedeliver.getRate();
468-
stats.chuckedMessageRate = chuckedMessageRate.getValueRate();
468+
stats.chuckedMessageRate = chuckedMessageRate.getRate();
469469
}
470470

471471
public ConsumerStats getStats() {

pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java

Lines changed: 0 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1214,7 +1214,6 @@ protected void handleSend(CommandSend send, ByteBuf headersAndPayload) {
12141214
producer.publishMessage(send.getProducerId(), send.getSequenceId(), headersAndPayload,
12151215
send.getNumMessages(), send.getIsChunk());
12161216
}
1217-
producer.publishMessage(send.getProducerId(), send.getSequenceId(), headersAndPayload, send.getNumMessages(), send.getIsChunk());
12181217
}
12191218

12201219
private void printSendCommandDebug(CommandSend send, ByteBuf headersAndPayload) {

pulsar-broker/src/test/java/org/apache/pulsar/client/impl/MessageChunkingTest.java

Lines changed: 13 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -68,6 +68,19 @@ protected void cleanup() throws Exception {
6868
super.internalCleanup();
6969
}
7070

71+
@Test
72+
public void testInvalidConfig() throws Exception {
73+
final String topicName = "persistent://my-property/my-ns/my-topic1";
74+
ProducerBuilder<byte[]> producerBuilder = pulsarClient.newProducer().topic(topicName);
75+
// batching and chunking can't be enabled together
76+
try {
77+
Producer<byte[]> producer = producerBuilder.enableChunking(true).enableBatching(true).create();
78+
fail("producer creation should have fail");
79+
} catch (IllegalArgumentException ie) {
80+
// Ok
81+
}
82+
}
83+
7184
@Test
7285
public void testLargeMessage() throws Exception {
7386

@@ -109,7 +122,6 @@ public void testLargeMessage() throws Exception {
109122
SubscriptionStats subStats = topic.getStats(false).subscriptions.entrySet().iterator().next().getValue();
110123

111124
assertTrue(producerStats.chunkedMessageRate > 0);
112-
assertEquals(producerStats.chunkedMessageRate, subStats.chuckedMessageRate);
113125

114126
ManagedCursorImpl mcursor = (ManagedCursorImpl) topic.getManagedLedger().getCursors().iterator().next();
115127
PositionImpl readPosition = (PositionImpl) mcursor.getReadPosition();

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

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -96,6 +96,9 @@ public Producer<T> create() throws PulsarClientException {
9696

9797
@Override
9898
public CompletableFuture<Producer<T>> createAsync() {
99+
// config validation
100+
checkArgument(!(conf.isBatchingEnabled() && conf.isChunkingEnabled()),
101+
"Batching and chunking of messages can't be enabled together");
99102
if (conf.getTopicName() == null) {
100103
return FutureUtil
101104
.failedFuture(new IllegalArgumentException("Topic name must be set on the producer builder"));

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

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -450,7 +450,7 @@ private void serializeAndSendMessage(MessageImpl<?> msg, Builder msgMetadataBuil
450450
sequenceId = chunkMsgMetadataBuilder.getSequenceId();
451451
}
452452
if (!chunkMsgMetadataBuilder.hasPublishTime()) {
453-
chunkMsgMetadataBuilder.setPublishTime(System.currentTimeMillis());
453+
chunkMsgMetadataBuilder.setPublishTime(client.getClientClock().millis());
454454

455455
checkArgument(!chunkMsgMetadataBuilder.hasProducerName());
456456

pulsar-common/src/main/java/org/apache/pulsar/common/protocol/Commands.java

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -482,6 +482,9 @@ public static ByteBufPair newSend(long producerId, long sequenceId, int numMessa
482482
if (txnIdMostBits > 0) {
483483
sendBuilder.setTxnidMostBits(txnIdMostBits);
484484
}
485+
if (messageData.hasTotalChunkMsgSize() && messageData.getTotalChunkMsgSize() > 1) {
486+
sendBuilder.setIsChunk(true);
487+
}
485488
CommandSend send = sendBuilder.build();
486489

487490
ByteBufPair res = serializeCommandSendWithSize(BaseCommand.newBuilder().setType(Type.SEND).setSend(send),

0 commit comments

Comments
 (0)