diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/MessagesImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/MessagesImpl.java index cfd7758734494..e0a54c50fed0f 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/MessagesImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/MessagesImpl.java @@ -29,7 +29,7 @@ @NotThreadSafe public class MessagesImpl implements Messages { - private List> messageList; + private final List> messageList; private final int maxNumberOfMessages; private final long maxSizeOfMessages; @@ -80,6 +80,10 @@ public void clear() { this.messageList.clear(); } + List> getMessageList() { + return messageList; + } + @Override public Iterator> iterator() { return messageList.iterator(); diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/MultiTopicsConsumerImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/MultiTopicsConsumerImpl.java index d40fe0b0e43c6..986a13e446a8c 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/MultiTopicsConsumerImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/MultiTopicsConsumerImpl.java @@ -49,6 +49,7 @@ import java.util.stream.Collectors; import java.util.stream.IntStream; import org.apache.commons.lang3.tuple.Pair; +import org.apache.pulsar.client.api.BatchReceivePolicy; import org.apache.pulsar.client.api.Consumer; import org.apache.pulsar.client.api.ConsumerStats; import org.apache.pulsar.client.api.Message; @@ -238,19 +239,25 @@ private void startReceivingMessages(List> newConsumers) { newConsumers.forEach(consumer -> { consumer.increaseAvailablePermits(consumer.getConnectionHandler().cnx(), consumer.getCurrentReceiverQueueSize()); - internalPinnedExecutor.execute(() -> receiveMessageFromConsumer(consumer)); + internalPinnedExecutor.execute(() -> receiveMessageFromConsumer(consumer, true)); }); } } - private void receiveMessageFromConsumer(ConsumerImpl consumer) { - consumer.receiveAsync().thenAcceptAsync(message -> { + private void receiveMessageFromConsumer(ConsumerImpl consumer, boolean batchReceive) { + CompletableFuture>> messagesFuture; + if (batchReceive) { + messagesFuture = consumer.batchReceiveAsync().thenApply(msgs -> ((MessagesImpl) msgs).getMessageList()); + } else { + messagesFuture = consumer.receiveAsync().thenApply(Collections::singletonList); + } + messagesFuture.thenAcceptAsync(messages -> { if (log.isDebugEnabled()) { log.debug("[{}] [{}] Receive message from sub consumer:{}", topic, subscription, consumer.getTopic()); } // Process the message, add to the queue and trigger listener or async callback - messageReceived(consumer, message); + messages.forEach(msg -> messageReceived(consumer, msg)); int size = incomingMessages.size(); int maxReceiverQueueSize = getCurrentReceiverQueueSize(); @@ -268,7 +275,7 @@ private void receiveMessageFromConsumer(ConsumerImpl consumer) { } else { // Call receiveAsync() if the incoming queue is not full. Because this block is run with // thenAcceptAsync, there is no chance for recursion that would lead to stack overflow. - receiveMessageFromConsumer(consumer); + receiveMessageFromConsumer(consumer, messages.size() > 0); } }, internalPinnedExecutor).exceptionally(ex -> { if (ex instanceof PulsarClientException.AlreadyClosedException @@ -277,8 +284,8 @@ private void receiveMessageFromConsumer(ConsumerImpl consumer) { return null; } log.error("Receive operation failed on consumer {} - Retrying later", consumer, ex); - ((ScheduledExecutorService) client.getScheduledExecutorProvider()) - .schedule(() -> receiveMessageFromConsumer(consumer), 10, TimeUnit.SECONDS); + ((ScheduledExecutorService) client.getScheduledExecutorProvider().getExecutor()) + .schedule(() -> receiveMessageFromConsumer(consumer, true), 10, TimeUnit.SECONDS); return null; }); } @@ -321,7 +328,7 @@ private void resumeReceivingFromPausedConsumersIfNeeded() { } internalPinnedExecutor.execute(() -> { - receiveMessageFromConsumer(consumer); + receiveMessageFromConsumer(consumer, true); }); } } @@ -1045,11 +1052,8 @@ private void doSubscribeTopicPartitions(Schema schema, String partitionName = TopicName.get(topicName).getPartition(partitionIndex).toString(); CompletableFuture> subFuture = new CompletableFuture<>(); configurationData.setStartPaused(paused); - ConsumerImpl newConsumer = ConsumerImpl.newConsumerImpl(client, partitionName, - configurationData, client.externalExecutorProvider(), - partitionIndex, true, listener != null, subFuture, - startMessageId, schema, interceptors, - createIfDoesNotExist, startMessageRollbackDurationInSec); + ConsumerImpl newConsumer = createInternalConsumer(configurationData, partitionName, + partitionIndex, subFuture, createIfDoesNotExist, schema); synchronized (pauseMutex) { if (paused) { newConsumer.pause(); @@ -1075,10 +1079,8 @@ private void doSubscribeTopicPartitions(Schema schema, return existingValue; } else { internalConfig.setStartPaused(paused); - ConsumerImpl newConsumer = ConsumerImpl.newConsumerImpl(client, topicName, internalConfig, - client.externalExecutorProvider(), -1, - true, listener != null, subFuture, startMessageId, schema, interceptors, - createIfDoesNotExist, startMessageRollbackDurationInSec); + ConsumerImpl newConsumer = createInternalConsumer(internalConfig, topicName, + -1, subFuture, createIfDoesNotExist, schema); synchronized (pauseMutex) { if (paused) { @@ -1121,6 +1123,22 @@ private void doSubscribeTopicPartitions(Schema schema, }); } + private ConsumerImpl createInternalConsumer(ConsumerConfigurationData configurationData, String partitionName, + int partitionIndex, CompletableFuture> subFuture, + boolean createIfDoesNotExist, Schema schema) { + BatchReceivePolicy internalBatchReceivePolicy = BatchReceivePolicy.builder() + .maxNumMessages(Math.max(configurationData.getReceiverQueueSize() / 2, 1)) + .maxNumBytes(-1) + .timeout(1, TimeUnit.MILLISECONDS) + .build(); + configurationData.setBatchReceivePolicy(internalBatchReceivePolicy); + return ConsumerImpl.newConsumerImpl(client, partitionName, + configurationData, client.externalExecutorProvider(), + partitionIndex, true, listener != null, subFuture, + startMessageId, schema, interceptors, + createIfDoesNotExist, startMessageRollbackDurationInSec); + } + // handling failure during subscribe new topic, unsubscribe success created partitions private void handleSubscribeOneTopicError(String topicName, Throwable error, @@ -1378,11 +1396,8 @@ private CompletableFuture subscribeIncreasedTopicPartitions(String topicNa CompletableFuture> subFuture = new CompletableFuture<>(); ConsumerConfigurationData configurationData = getInternalConsumerConfig(); configurationData.setStartPaused(paused); - ConsumerImpl newConsumer = ConsumerImpl.newConsumerImpl( - client, partitionName, configurationData, - client.externalExecutorProvider(), - partitionIndex, true, listener != null, subFuture, startMessageId, schema, interceptors, - true /* createTopicIfDoesNotExist */, startMessageRollbackDurationInSec); + ConsumerImpl newConsumer = createInternalConsumer(configurationData, partitionName, + partitionIndex, subFuture, true, schema); synchronized (pauseMutex) { if (paused) { newConsumer.pause();