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

作者:jiaxwu 时间:2024-05-22 10:14:29 

前言

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

这篇文章主要是使用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
}
  • 简单实现例子:https://github.com/jiaxwu/dq/blob/main/kafka_delay_queue_producer.go

  • 包含延迟服务的生产者:https://github.com/jiaxwu/dq/tree/main/kafka_delay_queue_example

来源:https://juejin.cn/post/7057584094766432270

标签:Go,Kafka,延迟消息
0
投稿

猜你喜欢

  • 详解vscode实现远程linux服务器上Python开发

    2021-07-02 05:09:32
  • python实现指定字符串补全空格、前面填充0的方法

    2022-04-06 21:13:58
  • 用Dreamweaver MX实现网站批量更新

    2009-09-13 18:39:00
  • 基于Python数据可视化利器Matplotlib,绘图入门篇,Pyplot详解

    2022-11-20 07:59:16
  • 一文教会你在sqlserver中创建表

    2024-01-17 11:42:57
  • python flask框架实现重定向功能示例

    2022-01-16 07:14:51
  • 在ASP.NET 2.0中操作数据之三十一:使用DataList来一行显示多条记录

    2024-05-11 09:30:00
  • Python使用signal定时结束AsyncIOScheduler任务的问题

    2022-12-19 21:28:11
  • ORACLE 10g 安装教程[图文]

    2023-07-15 07:07:27
  • 浅谈MySQL安装starting the server失败的解决办法

    2024-01-25 06:37:22
  • Python编程使用matplotlib挑钻石seaborn画图入门教程

    2021-04-12 18:28:17
  • 零基础写python爬虫之抓取糗事百科代码分享

    2021-02-01 11:54:39
  • keras.utils.to_categorical和one hot格式解析

    2023-10-03 18:27:12
  • python微信跳一跳系列之色块轮廓定位棋盘

    2022-10-18 04:33:22
  • php中用socket模拟http中post或者get提交数据的示例代码

    2023-11-19 00:45:21
  • Oracle使用PL/SQL操作COM对象

    2010-07-21 12:56:00
  • python必学知识之文件操作(建议收藏)

    2021-10-01 16:28:25
  • oracle group by语句实例测试

    2024-01-26 16:11:40
  • python对gif图压缩的完美解决方案

    2021-06-19 03:09:00
  • Windows下python3.7安装教程

    2023-02-16 16:39:11
  • asp之家 网络编程 m.aspxhome.com