網(wǎng)站首頁 編程語言 正文
前言
延遲隊(duì)列是一個(gè)非常有用的工具,我們經(jīng)常遇到需要使用延遲隊(duì)列的場(chǎng)景,比如延遲通知,訂單關(guān)閉等等。
這篇文章主要是使用Go+Kafka實(shí)現(xiàn)延遲消息。
使用了sarama客戶端。
原理
Kafka實(shí)現(xiàn)延遲消息分為下面三步:
- 生產(chǎn)者把消息發(fā)送到
延遲隊(duì)列
- 延遲服務(wù)把
延遲隊(duì)列
里超過延遲時(shí)間的消息寫入真實(shí)隊(duì)列
- 消費(fèi)者消費(fèi)
真實(shí)隊(duì)列
里的消息
簡(jiǎn)單的實(shí)現(xiàn)
生產(chǎn)者
生產(chǎn)者只是把消息發(fā)送到延遲隊(duì)列
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) }
延遲服務(wù)
延遲服務(wù)會(huì)訂閱延遲隊(duì)列
的消息,并把超時(shí)消息
發(fā)送到真實(shí)隊(duì)列
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() { // 如果消息已經(jīng)超時(shí),把消息發(fā)送到真實(shí)隊(duì)列 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 }
消費(fèi)者
消費(fèi)者只是訂閱真實(shí)隊(duì)列
并消費(fèi)消息
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 }
改進(jìn)點(diǎn)
通用的延遲服務(wù)
可以把延遲服務(wù)封裝成一個(gè)通用的服務(wù),這樣生產(chǎn)者可以直接把消息發(fā)送給延遲服務(wù),讓延遲服務(wù)去處理剩下的邏輯。
延遲服務(wù)可以提供多個(gè)延時(shí)等級(jí),比如5s、10s、30s、1m、5m、10m、1h、2h等,類似于RocketMQ。
生產(chǎn)者負(fù)責(zé)延遲服務(wù)
也可以讓生產(chǎn)者負(fù)責(zé)延遲服務(wù),讓生產(chǎn)者自己把延遲隊(duì)列里面的消息發(fā)送到真實(shí)隊(duì)列。
下面是一個(gè)簡(jiǎn)單的實(shí)現(xiàn):
// KafkaDelayQueueProducer 延遲隊(duì)列生產(chǎn)者,包含了生產(chǎn)者和延遲服務(wù) type KafkaDelayQueueProducer struct { producer sarama.SyncProducer // 生產(chǎn)者 delayTopic string // 延遲服務(wù)主題 } // NewKafkaDelayQueueProducer 創(chuàng)建延遲隊(duì)列生產(chǎn)者 // producer 生產(chǎn)者 // delayServiceConsumerGroup 延遲服務(wù)消費(fèi)者 // delayTime 延遲時(shí)間 // delayTopic 延遲服務(wù)主題 // realTopic 真實(shí)隊(duì)列主題 func NewKafkaDelayQueueProducer(producer sarama.SyncProducer, delayServiceConsumerGroup sarama.ConsumerGroup, delayTime time.Duration, delayTopic, realTopic string) *KafkaDelayQueueProducer { // 啟動(dòng)延遲服務(wù) 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 發(fā)送消息 func (q *KafkaDelayQueueProducer) SendMessage(msg *sarama.ProducerMessage) (partition int32, offset int64, err error) { msg.Topic = q.delayTopic return q.producer.SendMessage(msg) } // DelayServiceConsumer 延遲服務(wù)消費(fèi)者 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() { // 如果消息已經(jīng)超時(shí),把消息發(fā)送到真實(shí)隊(duì)列 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 }
總結(jié)
使用中間隊(duì)列
+輪詢
可以很容易的在Kafka實(shí)現(xiàn)延遲消息,如果需要一個(gè)通用的延遲隊(duì)列也可以實(shí)現(xiàn)一個(gè)通用的延遲服務(wù),也可以讓消費(fèi)者負(fù)責(zé)延遲服務(wù)的功能。
完整代碼:
- 簡(jiǎn)單實(shí)現(xiàn)例子:https://github.com/jiaxwu/dq/blob/main/kafka_delay_queue_producer.go
- 包含延遲服務(wù)的生產(chǎn)者:https://github.com/jiaxwu/dq/tree/main/kafka_delay_queue_example
原文鏈接:https://juejin.cn/post/7057584094766432270
相關(guān)推薦
- 2022-11-06 Python操作XML文件的使用指南_python
- 2022-06-21 Oracle新增和刪除用戶_oracle
- 2022-05-09 Python數(shù)據(jù)結(jié)構(gòu)與算法之鏈表,無序鏈表詳解_python
- 2023-07-30 el-tab在切換時(shí)echarts圖表顯示不正常
- 2022-11-10 android時(shí)間選擇控件之TimePickerView使用方法詳解_Android
- 2022-08-19 Mac上Chrome瀏覽器快捷鍵匯總
- 2022-08-31 pandas基礎(chǔ)?Series與Dataframe與numpy對(duì)二進(jìn)制文件輸入輸出_python
- 2022-05-13 三分鐘搞懂react-hooks及實(shí)例代碼_React
- 最近更新
-
- window11 系統(tǒng)安裝 yarn
- 超詳細(xì)win安裝深度學(xué)習(xí)環(huán)境2025年最新版(
- Linux 中運(yùn)行的top命令 怎么退出?
- MySQL 中decimal 的用法? 存儲(chǔ)小
- get 、set 、toString 方法的使
- @Resource和 @Autowired注解
- Java基礎(chǔ)操作-- 運(yùn)算符,流程控制 Flo
- 1. Int 和Integer 的區(qū)別,Jav
- spring @retryable不生效的一種
- Spring Security之認(rèn)證信息的處理
- Spring Security之認(rèn)證過濾器
- Spring Security概述快速入門
- Spring Security之配置體系
- 【SpringBoot】SpringCache
- Spring Security之基于方法配置權(quán)
- redisson分布式鎖中waittime的設(shè)
- maven:解決release錯(cuò)誤:Artif
- restTemplate使用總結(jié)
- Spring Security之安全異常處理
- MybatisPlus優(yōu)雅實(shí)現(xiàn)加密?
- Spring ioc容器與Bean的生命周期。
- 【探索SpringCloud】服務(wù)發(fā)現(xiàn)-Nac
- Spring Security之基于HttpR
- Redis 底層數(shù)據(jù)結(jié)構(gòu)-簡(jiǎn)單動(dòng)態(tài)字符串(SD
- arthas操作spring被代理目標(biāo)對(duì)象命令
- Spring中的單例模式應(yīng)用詳解
- 聊聊消息隊(duì)列,發(fā)送消息的4種方式
- bootspring第三方資源配置管理
- GIT同步修改后的遠(yuǎn)程分支