Go+Kafka实现延迟消息的实现示例

 更新时间:2022年07月25日 08:29:15   作者:jiaxwu  
本文主要介绍了Go+Kafka实现延迟消息的实现示例,文中通过示例代码介绍的非常详细,对大家的学习或者工作具有一定的参考学习价值,需要的朋友们下面随着小编来一起学习学习吧

前言

延迟队列是一个非常有用的工具,我们经常遇到需要使用延迟队列的场景,比如延迟通知,订单关闭等等。

这篇文章主要是使用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延迟消息内容请搜索脚本之家以前的文章或继续浏览下面的相关文章希望大家以后多多支持脚本之家!

相关文章

  • Go语言中一些不常见的命令参数详解

    Go语言中一些不常见的命令参数详解

    这篇文章主要给大家介绍了关于Go语言中一些不常见的命令参数的相关资料,文中通过示例代码介绍的非常详细,对大家的学习或者工作具有一定的参考学习价值,需要的朋友们下面随着小编来一起学习学习吧。
    2017-12-12
  • Go语言实现彩色输出示例详解

    Go语言实现彩色输出示例详解

    这篇文章主要为大家介绍了Go语言实现彩色输出示例详解,有需要的朋友可以借鉴参考下,希望能够有所帮助,祝大家多多进步,早日升职加薪
    2022-09-09
  • Golang并发发送HTTP请求的各种方法

    Golang并发发送HTTP请求的各种方法

    在 Golang 领域,并发发送 HTTP 请求是优化 Web 应用程序的一项重要技能,本文探讨了实现此目的的各种方法,从基本的 goroutine 到涉及通道和sync.WaitGroup 的高级技术,需要的朋友可以参考下
    2024-02-02
  • go使用errors.Wrapf()代替log.Error()方法示例

    go使用errors.Wrapf()代替log.Error()方法示例

    这篇文章主要为大家介绍了go使用errors.Wrapf()代替log.Error()的方法示例详解,有需要的朋友可以借鉴参考下,希望能够有所帮助,祝大家多多进步,早日升职加薪
    2023-08-08
  • Goland使用delve进行远程调试的详细教程

    Goland使用delve进行远程调试的详细教程

    网上给出的使用delve进行远程调试,都需要先在本地交叉编译或者在远程主机上编译出可运行的程序,然后再用delve在远程启动程序,本教程会将上面的步骤简化为只需要两步,1,在远程运行程序2,在本地启动调试,需要的朋友可以参考下
    2024-08-08
  • Go 中 time.After 可能导致的内存泄露问题解析

    Go 中 time.After 可能导致的内存泄露问题解析

    这篇文章主要介绍了Go 中 time.After 可能导致的内存泄露,本文给大家介绍的非常详细,对大家的学习或工作具有一定的参考借鉴价值,需要的朋友可以参考下
    2023-05-05
  • Golang限流库与漏桶和令牌桶的使用介绍

    Golang限流库与漏桶和令牌桶的使用介绍

    这篇文章主要介绍了golang限流库以及漏桶与令牌桶的实现原理,文中通过示例代码介绍的非常详细,对大家的学习或者工作具有一定的参考学习价值,需要的朋友们下面随着小编来一起学习吧
    2023-03-03
  • 解决golang时间字符串转time.Time的坑

    解决golang时间字符串转time.Time的坑

    这篇文章主要介绍了解决golang时间字符串转time.Time的坑,具有很好的参考价值,希望对大家有所帮助。一起跟随小编过来看看吧
    2021-04-04
  • Go语言数据类型详细介绍

    Go语言数据类型详细介绍

    这篇文章主要介绍了Go语言数据类型详细介绍,Go语言数据类型包含基础类型和复合类型两大类,下文关于这两类型的相关介绍,需要的小伙伴可以参考一下
    2022-03-03
  • 1行Go代码实现反向代理的示例

    1行Go代码实现反向代理的示例

    这篇文章主要介绍了1行Go代码实现反向代理的示例,小编觉得挺不错的,现在分享给大家,也给大家做个参考。一起跟随小编过来看看吧
    2018-08-08

最新评论