diff --git a/pulsar/consumer_impl.go b/pulsar/consumer_impl.go index ffca996d1d..909cb595e3 100644 --- a/pulsar/consumer_impl.go +++ b/pulsar/consumer_impl.go @@ -284,7 +284,7 @@ func newInternalConsumer(client *client, options ConsumerOptions, topic string, } if len(partitions) == 1 && options.EnableZeroQueueConsumer { - return newZeroConsumer(client, options, topic, messageCh, dlq, rlq, disableForceTopicCreation) + return newZeroConsumer(client, options, partitions[0], messageCh, dlq, rlq, disableForceTopicCreation) } consumer := &consumer{ diff --git a/pulsar/consumer_zero_queue_test.go b/pulsar/consumer_zero_queue_test.go index 8fe7985154..83e6b759d3 100644 --- a/pulsar/consumer_zero_queue_test.go +++ b/pulsar/consumer_zero_queue_test.go @@ -654,6 +654,34 @@ func TestSpecifiedPartitionZeroQueueConsumer(t *testing.T) { assert.Nil(t, err) } +func TestZeroQueueConsumerUsesOnlyPartitionForSinglePartitionTopic(t *testing.T) { + client, err := NewClient(ClientOptions{ + URL: lookupURL, + }) + + require.NoError(t, err) + defer client.Close() + + topic := "persistent://public/default/" + newTopicName() + err = createPartitionedTopic(topic, 1) + require.NoError(t, err) + topics, err := client.TopicPartitions(topic) + require.NoError(t, err) + require.Len(t, topics, 1) + + consumer, err := client.Subscribe(ConsumerOptions{ + Topic: topic, + SubscriptionName: "my-sub", + EnableZeroQueueConsumer: true, + }) + require.NoError(t, err) + defer consumer.Close() + + zeroConsumer, ok := consumer.(*zeroQueueConsumer) + require.True(t, ok) + assert.Equal(t, topics[0], zeroConsumer.pc.topic) +} + func TestZeroQueueConsumerGetLastMessageIDs(t *testing.T) { client, err := NewClient(ClientOptions{ URL: lookupURL,