functional_test.go 2.0 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091
  1. package sarama
  2. import (
  3. "fmt"
  4. "net"
  5. "os"
  6. "testing"
  7. "time"
  8. )
  9. const (
  10. TestBatchSize = 1000
  11. )
  12. var (
  13. kafkaIsAvailable, kafkaShouldBeAvailable bool
  14. kafkaAddr string
  15. )
  16. func init() {
  17. kafkaAddr = os.Getenv("KAFKA_ADDR")
  18. if kafkaAddr == "" {
  19. kafkaAddr = "localhost:9092"
  20. }
  21. c, err := net.Dial("tcp", kafkaAddr)
  22. if err == nil {
  23. kafkaIsAvailable = true
  24. c.Close()
  25. }
  26. kafkaShouldBeAvailable = os.Getenv("CI") != ""
  27. }
  28. func checkKafkaAvailability(t *testing.T) {
  29. if !kafkaIsAvailable {
  30. if kafkaShouldBeAvailable {
  31. t.Fatalf("Kafka broker is not available on %s. Set KAFKA_ADDR to connect to Kafka on a different location.", kafkaAddr)
  32. } else {
  33. t.Skipf("Kafka broker is not available on %s. Set KAFKA_ADDR to connect to Kafka on a different location.", kafkaAddr)
  34. }
  35. }
  36. }
  37. func TestProducingMessages(t *testing.T) {
  38. checkKafkaAvailability(t)
  39. client, err := NewClient("functional_test", []string{kafkaAddr}, nil)
  40. if err != nil {
  41. t.Fatal(err)
  42. }
  43. defer client.Close()
  44. consumerConfig := NewConsumerConfig()
  45. consumerConfig.OffsetMethod = OffsetMethodNewest
  46. consumer, err := NewConsumer(client, "single_partition", 0, "functional_test", consumerConfig)
  47. if err != nil {
  48. t.Fatal(err)
  49. }
  50. defer consumer.Close()
  51. producerConfig := NewProducerConfig()
  52. producerConfig.Partitioner = &ConstantPartitioner{Constant: 0}
  53. producer, err := NewProducer(client, producerConfig)
  54. if err != nil {
  55. t.Fatal(err)
  56. }
  57. for i := 1; i <= TestBatchSize; i++ {
  58. err = producer.SendMessage("single_partition", nil, StringEncoder(fmt.Sprintf("testing %d", i)))
  59. if err != nil {
  60. t.Fatal(err)
  61. }
  62. }
  63. producer.Close()
  64. events := consumer.Events()
  65. for i := 1; i <= TestBatchSize; i++ {
  66. select {
  67. case <-time.After(10 * time.Second):
  68. t.Fatal("Not received any more events in the last 10 seconds.")
  69. case event := <-events:
  70. if string(event.Value) != fmt.Sprintf("testing %d", i) {
  71. t.Fatalf("Unexpected message with index %d: %s", i, event.Value)
  72. }
  73. }
  74. }
  75. }