首頁 > 軟體

Go+Kafka實現延遲訊息的實現範例

2022-07-25 10:01:39

前言

延遲佇列是一個非常有用的工具,我們經常遇到需要使用延遲佇列的場景,比如延遲通知,訂單關閉等等。

這篇文章主要是使用Go+Kafka實現延遲訊息。

使用了sarama使用者端。

原理

Kafka實現延遲訊息分為下面三步:

  • 生產者把訊息傳送到延遲佇列
  • 延遲服務把延遲佇列裡超過延遲時間的訊息寫入真實佇列
  • 消費者消費真實佇列裡的訊息

簡單的實現

生產者

生產者只是把訊息傳送到延遲佇列

msg := &sarama.ProducerMessage{
   Topic: kafka_delay_queue_test.DelayTopic,
   Value: sarama.ByteEncoder("test" + strconv.Itoa(i)),
}
if _, _, err := producer.SendMessage(msg); err != nil {
   log.Println(err)
}

延遲服務

延遲服務會訂閱延遲佇列的訊息,並把超時訊息傳送到真實佇列

if err = consumerGroup.Consume(context.Background(),
   []string{kafka_delay_queue_test.DelayTopic}, consumer); err != nil {
   break
}
type Consumer struct {
   producer sarama.SyncProducer
   delay    time.Duration
}

func NewConsumer(producer sarama.SyncProducer, delay time.Duration) *Consumer {
   return &Consumer{
      producer: producer,
      delay:    delay,
   }
}

func (c *Consumer) ConsumeClaim(session sarama.ConsumerGroupSession, claim sarama.ConsumerGroupClaim) error {
   for message := range claim.Messages() {
      // 如果訊息已經超時,把訊息傳送到真實佇列
      now := time.Now()
      if now.Sub(message.Timestamp) >= c.delay {
         _, _, err := c.producer.SendMessage(&sarama.ProducerMessage{
            Topic: kafka_delay_queue_test.RealTopic,
            Key:   sarama.ByteEncoder(message.Key),
            Value: sarama.ByteEncoder(message.Value),
         })
         if err == nil {
            session.MarkMessage(message, "")
         }
         continue
      }
      // 否則休眠一秒
      time.Sleep(time.Second)
      return nil
   }
   return nil
}

消費者

消費者只是訂閱真實佇列並消費訊息

if err = consumerGroup.Consume(context.Background(), 
   []string{kafka_delay_queue_test.RealTopic}, consumer); err != nil {
   break
}
type Consumer struct{}

func NewConsumer() *Consumer {
   return &Consumer{}
}

func (c *Consumer) ConsumeClaim(session sarama.ConsumerGroupSession, claim sarama.ConsumerGroupClaim) error {
   for message := range claim.Messages() {
      fmt.Println("收到訊息:", message.Value, message.Timestamp)
      session.MarkMessage(message, "")
   }
   return nil
}

改進點

通用的延遲服務

可以把延遲服務封裝成一個通用的服務,這樣生產者可以直接把訊息傳送給延遲服務,讓延遲服務去處理剩下的邏輯。

延遲服務可以提供多個延時等級,比如5s、10s、30s、1m、5m、10m、1h、2h等,類似於RocketMQ。

生產者負責延遲服務

也可以讓生產者負責延遲服務,讓生產者自己把延遲佇列裡面的訊息傳送到真實佇列。

下面是一個簡單的實現:

// KafkaDelayQueueProducer 延遲佇列生產者,包含了生產者和延遲服務
type KafkaDelayQueueProducer struct {
   producer   sarama.SyncProducer // 生產者
   delayTopic string              // 延遲服務主題
}

// NewKafkaDelayQueueProducer 建立延遲佇列生產者
// producer 生產者
// delayServiceConsumerGroup 延遲服務消費者
// delayTime 延遲時間
// delayTopic 延遲服務主題
// realTopic 真實佇列主題
func NewKafkaDelayQueueProducer(producer sarama.SyncProducer, delayServiceConsumerGroup sarama.ConsumerGroup,
   delayTime time.Duration, delayTopic, realTopic string) *KafkaDelayQueueProducer {
   // 啟動延遲服務
   consumer := NewDelayServiceConsumer(producer, delayTime, realTopic)
   go func() {
      for {
         if err := delayServiceConsumerGroup.Consume(context.Background(),
            []string{delayTopic}, consumer); err != nil {
            break
         }
      }
   }()
   return &KafkaDelayQueueProducer{
      producer:   producer,
      delayTopic: delayTopic,
   }
}

// SendMessage 傳送訊息
func (q *KafkaDelayQueueProducer) SendMessage(msg *sarama.ProducerMessage) (partition int32, offset int64, err error) {
   msg.Topic = q.delayTopic
   return q.producer.SendMessage(msg)
}

// DelayServiceConsumer 延遲服務消費者
type DelayServiceConsumer struct {
   producer  sarama.SyncProducer
   delay     time.Duration
   realTopic string
}

func NewDelayServiceConsumer(producer sarama.SyncProducer, delay time.Duration,
   realTopic string) *DelayServiceConsumer {
   return &DelayServiceConsumer{
      producer:  producer,
      delay:     delay,
      realTopic: realTopic,
   }
}

func (c *DelayServiceConsumer) ConsumeClaim(session sarama.ConsumerGroupSession,
   claim sarama.ConsumerGroupClaim) error {
   for message := range claim.Messages() {
      // 如果訊息已經超時,把訊息傳送到真實佇列
      now := time.Now()
      if now.Sub(message.Timestamp) >= c.delay {
         _, _, err := c.producer.SendMessage(&sarama.ProducerMessage{
            Topic: c.realTopic,
            Key:   sarama.ByteEncoder(message.Key),
            Value: sarama.ByteEncoder(message.Value),
         })
         if err == nil {
            session.MarkMessage(message, "")
         }
         continue
      }
      // 否則休眠一秒
      time.Sleep(time.Second)
      return nil
   }
   return nil
}

func (c *DelayServiceConsumer) Setup(sarama.ConsumerGroupSession) error {
   return nil
}

func (c *DelayServiceConsumer) Cleanup(sarama.ConsumerGroupSession) error {
   return nil
}

總結

使用中間佇列+輪詢可以很容易的在Kafka實現延遲訊息,如果需要一個通用的延遲佇列也可以實現一個通用的延遲服務,也可以讓消費者負責延遲服務的功能。

完整程式碼:

到此這篇關於Go+Kafka實現延遲訊息的實現範例的文章就介紹到這了,更多相關Go Kafka延遲訊息內容請搜尋it145.com以前的文章或繼續瀏覽下面的相關文章希望大家以後多多支援it145.com!


IT145.com E-mail:sddin#qq.com