functional_test.go 4.0 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179
  1. package sarama
  2. import (
  3. "fmt"
  4. "net"
  5. "os"
  6. "sync"
  7. "testing"
  8. "time"
  9. )
  10. const (
  11. TestBatchSize = 1000
  12. )
  13. var (
  14. kafkaIsAvailable, kafkaShouldBeAvailable bool
  15. kafkaAddr string
  16. )
  17. func init() {
  18. kafkaAddr = os.Getenv("KAFKA_ADDR")
  19. if kafkaAddr == "" {
  20. kafkaAddr = "localhost:6667"
  21. }
  22. c, err := net.Dial("tcp", kafkaAddr)
  23. if err == nil {
  24. kafkaIsAvailable = true
  25. c.Close()
  26. }
  27. kafkaShouldBeAvailable = os.Getenv("CI") != ""
  28. }
  29. func checkKafkaAvailability(t *testing.T) {
  30. if !kafkaIsAvailable {
  31. if kafkaShouldBeAvailable {
  32. t.Fatalf("Kafka broker is not available on %s. Set KAFKA_ADDR to connect to Kafka on a different location.", kafkaAddr)
  33. } else {
  34. t.Skipf("Kafka broker is not available on %s. Set KAFKA_ADDR to connect to Kafka on a different location.", kafkaAddr)
  35. }
  36. }
  37. }
  38. func TestFuncProducing(t *testing.T) {
  39. config := NewProducerConfig()
  40. testProducingMessages(t, config)
  41. }
  42. func TestFuncProducingGzip(t *testing.T) {
  43. config := NewProducerConfig()
  44. config.Compression = CompressionGZIP
  45. testProducingMessages(t, config)
  46. }
  47. func TestFuncProducingSnappy(t *testing.T) {
  48. config := NewProducerConfig()
  49. config.Compression = CompressionSnappy
  50. testProducingMessages(t, config)
  51. }
  52. func TestFuncProducingNoResponse(t *testing.T) {
  53. config := NewProducerConfig()
  54. config.RequiredAcks = NoResponse
  55. testProducingMessages(t, config)
  56. }
  57. func TestFuncProducingFlushing(t *testing.T) {
  58. config := NewProducerConfig()
  59. config.FlushMsgCount = TestBatchSize / 8
  60. config.FlushFrequency = 250 * time.Millisecond
  61. testProducingMessages(t, config)
  62. }
  63. func TestFuncMultiPartitionProduce(t *testing.T) {
  64. checkKafkaAvailability(t)
  65. client, err := NewClient("functional_test", []string{kafkaAddr}, nil)
  66. if err != nil {
  67. t.Fatal(err)
  68. }
  69. defer safeClose(t, client)
  70. config := NewProducerConfig()
  71. config.FlushFrequency = 50 * time.Millisecond
  72. config.FlushMsgCount = 200
  73. config.ChannelBufferSize = 20
  74. config.AckSuccesses = true
  75. producer, err := NewProducer(client, config)
  76. if err != nil {
  77. t.Fatal(err)
  78. }
  79. var wg sync.WaitGroup
  80. wg.Add(TestBatchSize)
  81. for i := 1; i <= TestBatchSize; i++ {
  82. go func(i int, w *sync.WaitGroup) {
  83. defer w.Done()
  84. msg := &MessageToSend{Topic: "multi_partition", Key: nil, Value: StringEncoder(fmt.Sprintf("hur %d", i))}
  85. producer.Input() <- msg
  86. select {
  87. case ret := <-producer.Errors():
  88. t.Fatal(ret.Err)
  89. case <-producer.Successes():
  90. }
  91. }(i, &wg)
  92. }
  93. wg.Wait()
  94. if err := producer.Close(); err != nil {
  95. t.Error(err)
  96. }
  97. }
  98. func testProducingMessages(t *testing.T, config *ProducerConfig) {
  99. checkKafkaAvailability(t)
  100. client, err := NewClient("functional_test", []string{kafkaAddr}, nil)
  101. if err != nil {
  102. t.Fatal(err)
  103. }
  104. defer safeClose(t, client)
  105. consumerConfig := NewConsumerConfig()
  106. consumerConfig.OffsetMethod = OffsetMethodNewest
  107. consumer, err := NewConsumer(client, "single_partition", 0, "functional_test", consumerConfig)
  108. if err != nil {
  109. t.Fatal(err)
  110. }
  111. defer safeClose(t, consumer)
  112. config.AckSuccesses = true
  113. producer, err := NewProducer(client, config)
  114. if err != nil {
  115. t.Fatal(err)
  116. }
  117. expectedResponses := TestBatchSize
  118. for i := 1; i <= TestBatchSize; {
  119. msg := &MessageToSend{Topic: "single_partition", Key: nil, Value: StringEncoder(fmt.Sprintf("testing %d", i))}
  120. select {
  121. case producer.Input() <- msg:
  122. i++
  123. case ret := <-producer.Errors():
  124. t.Fatal(ret.Err)
  125. case <-producer.Successes():
  126. expectedResponses--
  127. }
  128. }
  129. for expectedResponses > 0 {
  130. select {
  131. case ret := <-producer.Errors():
  132. t.Fatal(ret.Err)
  133. case <-producer.Successes():
  134. expectedResponses--
  135. }
  136. }
  137. err = producer.Close()
  138. if err != nil {
  139. t.Error(err)
  140. }
  141. events := consumer.Events()
  142. for i := 1; i <= TestBatchSize; i++ {
  143. select {
  144. case <-time.After(10 * time.Second):
  145. t.Fatal("Not received any more events in the last 10 seconds.")
  146. case event := <-events:
  147. if string(event.Value) != fmt.Sprintf("testing %d", i) {
  148. t.Fatalf("Unexpected message with index %d: %s", i, event.Value)
  149. }
  150. }
  151. }
  152. }