123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386 |
- package sarama
- import (
- "fmt"
- "sync"
- "testing"
- "time"
- )
- func TestDefaultConsumerConfigValidates(t *testing.T) {
- config := NewConsumerConfig()
- if err := config.Validate(); err != nil {
- t.Error(err)
- }
- }
- func TestDefaultPartitionConsumerConfigValidates(t *testing.T) {
- config := NewPartitionConsumerConfig()
- if err := config.Validate(); err != nil {
- t.Error(err)
- }
- }
- func TestConsumerOffsetManual(t *testing.T) {
- seedBroker := newMockBroker(t, 1)
- leader := newMockBroker(t, 2)
- metadataResponse := new(MetadataResponse)
- metadataResponse.AddBroker(leader.Addr(), leader.BrokerID())
- metadataResponse.AddTopicPartition("my_topic", 0, leader.BrokerID(), nil, nil, ErrNoError)
- seedBroker.Returns(metadataResponse)
- for i := 0; i <= 10; i++ {
- fetchResponse := new(FetchResponse)
- fetchResponse.AddMessage("my_topic", 0, nil, ByteEncoder([]byte{0x00, 0x0E}), int64(i+1234))
- leader.Returns(fetchResponse)
- }
- client, err := NewClient("client_id", []string{seedBroker.Addr()}, nil)
- if err != nil {
- t.Fatal(err)
- }
- master, err := NewConsumer(client, nil)
- if err != nil {
- t.Fatal(err)
- }
- config := NewPartitionConsumerConfig()
- config.OffsetMethod = OffsetMethodManual
- config.OffsetValue = 1234
- consumer, err := master.ConsumePartition("my_topic", 0, config)
- if err != nil {
- t.Fatal(err)
- }
- seedBroker.Close()
- for i := 0; i < 10; i++ {
- select {
- case message := <-consumer.Messages():
- if message.Offset != int64(i+1234) {
- t.Error("Incorrect message offset!")
- }
- case err := <-consumer.Errors():
- t.Error(err)
- }
- }
- safeClose(t, consumer)
- safeClose(t, client)
- leader.Close()
- }
- func TestConsumerLatestOffset(t *testing.T) {
- seedBroker := newMockBroker(t, 1)
- leader := newMockBroker(t, 2)
- metadataResponse := new(MetadataResponse)
- metadataResponse.AddBroker(leader.Addr(), leader.BrokerID())
- metadataResponse.AddTopicPartition("my_topic", 0, leader.BrokerID(), nil, nil, ErrNoError)
- seedBroker.Returns(metadataResponse)
- offsetResponse := new(OffsetResponse)
- offsetResponse.AddTopicPartition("my_topic", 0, 0x010101)
- leader.Returns(offsetResponse)
- fetchResponse := new(FetchResponse)
- fetchResponse.AddMessage("my_topic", 0, nil, ByteEncoder([]byte{0x00, 0x0E}), 0x010101)
- leader.Returns(fetchResponse)
- client, err := NewClient("client_id", []string{seedBroker.Addr()}, nil)
- if err != nil {
- t.Fatal(err)
- }
- seedBroker.Close()
- master, err := NewConsumer(client, nil)
- if err != nil {
- t.Fatal(err)
- }
- config := NewPartitionConsumerConfig()
- config.OffsetMethod = OffsetMethodNewest
- consumer, err := master.ConsumePartition("my_topic", 0, config)
- if err != nil {
- t.Fatal(err)
- }
- leader.Close()
- safeClose(t, consumer)
- safeClose(t, client)
-
- if consumer.offset != 0x010102 {
- t.Error("Latest offset not fetched correctly:", consumer.offset)
- }
- }
- func TestConsumerFunnyOffsets(t *testing.T) {
-
-
-
- seedBroker := newMockBroker(t, 1)
- leader := newMockBroker(t, 2)
- metadataResponse := new(MetadataResponse)
- metadataResponse.AddBroker(leader.Addr(), leader.BrokerID())
- metadataResponse.AddTopicPartition("my_topic", 0, leader.BrokerID(), nil, nil, ErrNoError)
- seedBroker.Returns(metadataResponse)
- fetchResponse := new(FetchResponse)
- fetchResponse.AddMessage("my_topic", 0, nil, ByteEncoder([]byte{0x00, 0x0E}), int64(1))
- fetchResponse.AddMessage("my_topic", 0, nil, ByteEncoder([]byte{0x00, 0x0E}), int64(3))
- leader.Returns(fetchResponse)
- fetchResponse = new(FetchResponse)
- fetchResponse.AddMessage("my_topic", 0, nil, ByteEncoder([]byte{0x00, 0x0E}), int64(5))
- leader.Returns(fetchResponse)
- client, err := NewClient("client_id", []string{seedBroker.Addr()}, nil)
- if err != nil {
- t.Fatal(err)
- }
- master, err := NewConsumer(client, nil)
- if err != nil {
- t.Fatal(err)
- }
- config := NewPartitionConsumerConfig()
- config.OffsetMethod = OffsetMethodManual
- config.OffsetValue = 2
- consumer, err := master.ConsumePartition("my_topic", 0, config)
- message := <-consumer.Messages()
- if message.Offset != 3 {
- t.Error("Incorrect message offset!")
- }
- leader.Close()
- seedBroker.Close()
- safeClose(t, consumer)
- safeClose(t, client)
- }
- func TestConsumerRebalancingMultiplePartitions(t *testing.T) {
-
- seedBroker := newMockBroker(t, 1)
- leader0 := newMockBroker(t, 2)
- leader1 := newMockBroker(t, 3)
- metadataResponse := new(MetadataResponse)
- metadataResponse.AddBroker(leader0.Addr(), leader0.BrokerID())
- metadataResponse.AddBroker(leader1.Addr(), leader1.BrokerID())
- metadataResponse.AddTopicPartition("my_topic", 0, leader0.BrokerID(), nil, nil, ErrNoError)
- metadataResponse.AddTopicPartition("my_topic", 1, leader1.BrokerID(), nil, nil, ErrNoError)
- seedBroker.Returns(metadataResponse)
-
- client, err := NewClient("client_id", []string{seedBroker.Addr()}, nil)
- if err != nil {
- t.Fatal(err)
- }
- master, err := NewConsumer(client, nil)
- if err != nil {
- t.Fatal(err)
- }
- config := NewPartitionConsumerConfig()
- config.OffsetMethod = OffsetMethodManual
- config.OffsetValue = 0
-
- var wg sync.WaitGroup
- for i := 0; i < 2; i++ {
- consumer, err := master.ConsumePartition("my_topic", int32(i), config)
- if err != nil {
- t.Error(err)
- }
- go func(c *PartitionConsumer) {
- for err := range c.Errors() {
- t.Error(err)
- }
- }(consumer)
- wg.Add(1)
- go func(partition int32, c *PartitionConsumer) {
- for i := 0; i < 10; i++ {
- message := <-consumer.Messages()
- if message.Offset != int64(i) {
- t.Error("Incorrect message offset!", i, partition, message.Offset)
- }
- if message.Partition != partition {
- t.Error("Incorrect message partition!")
- }
- }
- safeClose(t, consumer)
- wg.Done()
- }(int32(i), consumer)
- }
-
- fetchResponse := new(FetchResponse)
- for i := 0; i < 4; i++ {
- fetchResponse.AddMessage("my_topic", 0, nil, ByteEncoder([]byte{0x00, 0x0E}), int64(i))
- }
- leader0.Returns(fetchResponse)
-
- fetchResponse = new(FetchResponse)
- fetchResponse.AddError("my_topic", 0, ErrNotLeaderForPartition)
- leader0.Returns(fetchResponse)
-
- metadataResponse = new(MetadataResponse)
- metadataResponse.AddTopicPartition("my_topic", 0, leader1.BrokerID(), nil, nil, ErrNoError)
- metadataResponse.AddTopicPartition("my_topic", 1, leader1.BrokerID(), nil, nil, ErrNoError)
- seedBroker.Returns(metadataResponse)
- time.Sleep(5 * time.Millisecond)
-
- fetchResponse = new(FetchResponse)
- for i := 0; i < 5; i++ {
- fetchResponse.AddMessage("my_topic", 1, nil, ByteEncoder([]byte{0x00, 0x0E}), int64(i))
- }
- leader1.Returns(fetchResponse)
-
- fetchResponse = new(FetchResponse)
- for i := 0; i < 3; i++ {
- fetchResponse.AddMessage("my_topic", 0, nil, ByteEncoder([]byte{0x00, 0x0E}), int64(i+4))
- fetchResponse.AddMessage("my_topic", 1, nil, ByteEncoder([]byte{0x00, 0x0E}), int64(i+5))
- }
- leader1.Returns(fetchResponse)
-
- fetchResponse = new(FetchResponse)
- for i := 0; i < 3; i++ {
- fetchResponse.AddMessage("my_topic", 0, nil, ByteEncoder([]byte{0x00, 0x0E}), int64(i+7))
- }
- fetchResponse.AddError("my_topic", 1, ErrNotLeaderForPartition)
- leader1.Returns(fetchResponse)
-
- metadataResponse = new(MetadataResponse)
- metadataResponse.AddTopicPartition("my_topic", 0, leader1.BrokerID(), nil, nil, ErrNoError)
- metadataResponse.AddTopicPartition("my_topic", 1, leader0.BrokerID(), nil, nil, ErrNoError)
- seedBroker.Returns(metadataResponse)
- time.Sleep(5 * time.Millisecond)
-
- fetchResponse = new(FetchResponse)
- fetchResponse.AddMessage("my_topic", 1, nil, ByteEncoder([]byte{0x00, 0x0E}), int64(8))
- fetchResponse.AddMessage("my_topic", 1, nil, ByteEncoder([]byte{0x00, 0x0E}), int64(9))
- leader0.Returns(fetchResponse)
-
- fetchResponse = new(FetchResponse)
- fetchResponse.AddMessage("my_topic", 1, nil, ByteEncoder([]byte{0x00, 0x0E}), int64(10))
- leader0.Returns(fetchResponse)
-
- fetchResponse = new(FetchResponse)
- fetchResponse.AddMessage("my_topic", 0, nil, ByteEncoder([]byte{0x00, 0x0E}), int64(10))
- leader1.Returns(fetchResponse)
- wg.Wait()
- leader1.Close()
- leader0.Close()
- seedBroker.Close()
- safeClose(t, client)
- }
- func ExampleConsumerWithSelect() {
- client, err := NewClient("my_client", []string{"localhost:9092"}, nil)
- if err != nil {
- panic(err)
- } else {
- fmt.Println("> connected")
- }
- defer client.Close()
- master, err := NewConsumer(client, nil)
- if err != nil {
- panic(err)
- } else {
- fmt.Println("> master consumer ready")
- }
- consumer, err := master.ConsumePartition("my_topic", 0, nil)
- if err != nil {
- panic(err)
- } else {
- fmt.Println("> consumer ready")
- }
- defer consumer.Close()
- msgCount := 0
- consumerLoop:
- for {
- select {
- case err := <-consumer.Errors():
- panic(err)
- case <-consumer.Messages():
- msgCount++
- case <-time.After(5 * time.Second):
- fmt.Println("> timed out")
- break consumerLoop
- }
- }
- fmt.Println("Got", msgCount, "messages.")
- }
- func ExampleConsumerWithGoroutines() {
- client, err := NewClient("my_client", []string{"localhost:9092"}, nil)
- if err != nil {
- panic(err)
- } else {
- fmt.Println("> connected")
- }
- defer client.Close()
- master, err := NewConsumer(client, nil)
- if err != nil {
- panic(err)
- } else {
- fmt.Println("> master consumer ready")
- }
- consumer, err := master.ConsumePartition("my_topic", 0, nil)
- if err != nil {
- panic(err)
- } else {
- fmt.Println("> consumer ready")
- }
- defer consumer.Close()
- var (
- wg sync.WaitGroup
- msgCount int
- )
- wg.Add(1)
- go func() {
- defer wg.Done()
- for message := range consumer.Messages() {
- fmt.Printf("Consumed message with offset %d", message.Offset)
- msgCount++
- }
- }()
- wg.Add(1)
- go func() {
- defer wg.Done()
- for err := range consumer.Errors() {
- fmt.Println(err)
- }
- }()
- wg.Wait()
- fmt.Println("Got", msgCount, "messages.")
- }
|