Go+Kafka实现延迟消息的实现示例
作者:jiaxwu 发布时间:2024-05-22 10:14:29
标签: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
}
简单实现例子: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


猜你喜欢
- Mybatis插入mysql报主键重复的问题首先思路是这样的,先去数据表里面去找有没有这个主键的数据(如果有会有返回值,如果没有则返回nul
- 简介一款跨平台/无依赖的自动化测试工具,目测只能控制鼠标/键盘/获取屏幕尺寸/弹出消息框/截屏。安装pip install pyautogu
- 直接po截图和代码下面是CheckFormDemo.html<!DOCTYPE html><html><hea
- 使用iframe嵌入网页,页面可自适应在项目中遇到要嵌入第三方网页的需求,因为没有同第三方页面交互的需求,只需展示即可,所以最终决定使用 i
- 一、FBVFBV(function base views) 就是在视图里使用函数处理请求。二、CBVCBV(class base views
- 用Python进行爬取网页文字的代码:#!/usr/bin/python# -*- coding: UTF-8 -*-import requ
- Windows•安装lxml最好的安装方式是通过wheel文件来安装,http://www.lfd.uci.edu/~gohlke/pyth
- 目前网络数据库的应用已经成为最为广泛的应用之一了,并且关于数据库的安全性,性能都是企业最为关心的事情。数据库渐渐成为企业的命脉,优化查询就解
- 函数介绍Socket对象方法:服务端:函数描述.bind()绑定地址关键字,AF_INET下以元组的形式表示地址。常用bind((host,
- 本文实例讲述了javascript设置和获取cookie的方法。分享给大家供大家参考,具体如下:1. 设置cookiefunction se
- 我们想要知道数目的总和,只要通过+就能实现,这是我们在做题上经常用到的符号。但是在python中不能直接使用,我们需要借助一些代码或者函数帮
- mysqldump工具备份备份整个数据库$> mysqldump -u root -h host -p dbname > bac
- 官方文档介绍链接:append方法介绍DataFrame.append(other, ignore_index=False, verify_
- 本文从多个角度来讲解如何在Access数据库上如何上传并且显示上所上传图片。在 * 站制做过程中,需要上传图片、显示图片,上传的图片要能够保
- Python包导入报错的问题首先,一般来说,写一个小demo可能一个文件就够了,但是要是做一个小项目,可能需要拆分成很多零散的文件,放在不同
- 最近常有厦门的客户通过网站上的联系方式加我QQ,询问网站改版的情况。几乎每日都要针对客户网站存在的问题做一番分析,然后客户以价格等其他因素结
- 在使用npm 的过程中,搜索网上的资料基本上可以看到类似如下的描述:“npm是国外的,使用起来比较慢,我们这里使用淘宝的cnpm镜像”。初体
- 如下所示:string =" { "status": "error", "mes
- Matplotlib可以无缝的处理LaTex字体,在图中加入数学公式from matplotlib.patches import Polyg
- 1.SocketServer模块编写的TCP服务器端代码Socketserver原理图服务端:import SocketServer &nb