Skip to content
Merged
Show file tree
Hide file tree
Changes from 1 commit
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
2 changes: 1 addition & 1 deletion pulsar/consumer_impl.go
Original file line number Diff line number Diff line change
Expand Up @@ -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{
Expand Down
30 changes: 30 additions & 0 deletions pulsar/consumer_zero_queue_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -654,6 +654,36 @@ 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 := newTopicName()
err = createPartitionedTopic(topic, 1)
require.NoError(t, err)
topics, err := client.TopicPartitions(topic)
require.NoError(t, err)
require.Len(t, topics, 1)
topicName, err := internal.ParseTopicName(topics[0])
require.NoError(t, err)

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, topicName.Name, zeroConsumer.pc.topic)
Comment thread
Technoboy- marked this conversation as resolved.
Outdated
}

func TestZeroQueueConsumerGetLastMessageIDs(t *testing.T) {
client, err := NewClient(ClientOptions{
URL: lookupURL,
Expand Down
Loading