main.go 3.8 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148
  1. package main
  2. import (
  3. "context"
  4. "flag"
  5. "log"
  6. "os"
  7. "os/signal"
  8. "strings"
  9. "sync"
  10. "syscall"
  11. "github.com/Shopify/sarama"
  12. )
  13. // Sarma configuration options
  14. var (
  15. brokers = ""
  16. version = ""
  17. group = ""
  18. topics = ""
  19. oldest = true
  20. verbose = false
  21. )
  22. func init() {
  23. flag.StringVar(&brokers, "brokers", "", "Kafka bootstrap brokers to connect to, as a comma separated list")
  24. flag.StringVar(&group, "group", "", "Kafka consumer group definition")
  25. flag.StringVar(&version, "version", "2.1.1", "Kafka cluster version")
  26. flag.StringVar(&topics, "topics", "", "Kafka topics to be consumed, as a comma seperated list")
  27. flag.BoolVar(&oldest, "oldest", true, "Kafka consumer consume initial ofset from oldest")
  28. flag.BoolVar(&verbose, "verbose", false, "Sarama logging")
  29. flag.Parse()
  30. if len(brokers) == 0 {
  31. panic("no Kafka bootstrap brokers defined, please set the -brokers flag")
  32. }
  33. if len(topics) == 0 {
  34. panic("no topics given to be consumed, please set the -topics flag")
  35. }
  36. if len(group) == 0 {
  37. panic("no Kafka consumer group defined, please set the -group flag")
  38. }
  39. }
  40. func main() {
  41. log.Println("Starting a new Sarama consumer")
  42. if verbose {
  43. sarama.Logger = log.New(os.Stdout, "[sarama] ", log.LstdFlags)
  44. }
  45. version, err := sarama.ParseKafkaVersion(version)
  46. if err != nil {
  47. log.Panicf("Error parsing Kafka version: %v", err)
  48. }
  49. /**
  50. * Construct a new Sarama configuration.
  51. * The Kafka cluster version has to be defined before the consumer/producer is initialized.
  52. */
  53. config := sarama.NewConfig()
  54. config.Version = version
  55. if oldest {
  56. config.Consumer.Offsets.Initial = sarama.OffsetOldest
  57. }
  58. /**
  59. * Setup a new Sarama consumer group
  60. */
  61. consumer := Consumer{
  62. ready: make(chan bool, 0),
  63. }
  64. ctx, cancel := context.WithCancel(context.Background())
  65. client, err := sarama.NewConsumerGroup(strings.Split(brokers, ","), group, config)
  66. if err != nil {
  67. log.Panicf("Error creating consumer group client: %v", err)
  68. }
  69. wg := &sync.WaitGroup{}
  70. go func() {
  71. wg.Add(1)
  72. defer wg.Done()
  73. for {
  74. if err := client.Consume(ctx, strings.Split(topics, ","), &consumer); err != nil {
  75. log.Panicf("Error from consumer: %v", err)
  76. }
  77. // check if context was cancelled, signaling that the consumer should stop
  78. if ctx.Err() != nil {
  79. return
  80. }
  81. consumer.ready = make(chan bool, 0)
  82. }
  83. }()
  84. <-consumer.ready // Await till the consumer has been set up
  85. log.Println("Sarama consumer up and running!...")
  86. sigterm := make(chan os.Signal, 1)
  87. signal.Notify(sigterm, syscall.SIGINT, syscall.SIGTERM)
  88. select {
  89. case <-ctx.Done():
  90. log.Println("terminating: context cancelled")
  91. case <-sigterm:
  92. log.Println("terminating: via signal")
  93. }
  94. cancel()
  95. wg.Wait()
  96. if err = client.Close(); err != nil {
  97. log.Panicf("Error closing client: %v", err)
  98. }
  99. }
  100. // Consumer represents a Sarama consumer group consumer
  101. type Consumer struct {
  102. ready chan bool
  103. }
  104. // Setup is run at the beginning of a new session, before ConsumeClaim
  105. func (consumer *Consumer) Setup(sarama.ConsumerGroupSession) error {
  106. // Mark the consumer as ready
  107. close(consumer.ready)
  108. return nil
  109. }
  110. // Cleanup is run at the end of a session, once all ConsumeClaim goroutines have exited
  111. func (consumer *Consumer) Cleanup(sarama.ConsumerGroupSession) error {
  112. return nil
  113. }
  114. // ConsumeClaim must start a consumer loop of ConsumerGroupClaim's Messages().
  115. func (consumer *Consumer) ConsumeClaim(session sarama.ConsumerGroupSession, claim sarama.ConsumerGroupClaim) error {
  116. // NOTE:
  117. // Do not move the code below to a goroutine.
  118. // The `ConsumeClaim` itself is called within a goroutine, see:
  119. // https://github.com/Shopify/sarama/blob/master/consumer_group.go#L27-L29
  120. for message := range claim.Messages() {
  121. log.Printf("Message claimed: value = %s, timestamp = %v, topic = %s", string(message.Value), message.Timestamp, message.Topic)
  122. session.MarkMessage(message, "")
  123. }
  124. return nil
  125. }