Go+Kafka實(shí)現(xiàn)延遲消息的實(shí)現(xiàn)示例
前言
延遲隊(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ì)列里超過(guò)延遲時(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
到此這篇關(guān)于Go+Kafka實(shí)現(xiàn)延遲消息的實(shí)現(xiàn)示例的文章就介紹到這了,更多相關(guān)Go Kafka延遲消息內(nèi)容請(qǐng)搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!
相關(guān)文章
golang實(shí)現(xiàn)微信支付v3版本的方法
這篇文章主要介紹了golang實(shí)現(xiàn)微信支付v3版本的方法,本文給大家介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或工作具有一定的參考借鑒價(jià)值,需要的朋友可以參考下2021-03-03
Go語(yǔ)言中?Print?Printf和Println?的區(qū)別解析
這篇文章主要介紹了Go語(yǔ)言中?Print?Printf和Println?的區(qū)別,本文給大家介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或工作具有一定的參考借鑒價(jià)值,需要的朋友可以參考下2023-03-03
一文帶你了解Go語(yǔ)言中的I/O接口設(shè)計(jì)
I/O?操作在編程中扮演著至關(guān)重要的角色,它涉及程序與外部世界之間的數(shù)據(jù)交換,下面我們就來(lái)簡(jiǎn)單了解一下Go語(yǔ)言中的?I/O?接口設(shè)計(jì)吧2023-06-06
Go語(yǔ)言從單體服務(wù)到微服務(wù)設(shè)計(jì)方案詳解
這篇文章主要為大家介紹了Go語(yǔ)言從單體服務(wù)到微服務(wù)設(shè)計(jì)方案詳解,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步,早日升職加薪2023-03-03
Ubuntu安裝Go語(yǔ)言運(yùn)行環(huán)境
由于最近偏愛Ubuntu,在加上作為一門開源語(yǔ)言,在Linux上從源代碼開始搭建環(huán)境更讓人覺(jué)得有趣味性。讓我們直接先從Go語(yǔ)言的環(huán)境搭建開始2015-04-04
詳解Go語(yǔ)言如何使用xorm實(shí)現(xiàn)讀取mysql
xorm是go語(yǔ)言的常用orm之一,可以用來(lái)操作數(shù)據(jù)庫(kù)。本文就來(lái)和大家聊聊Go語(yǔ)言如何使用xorm實(shí)現(xiàn)讀取mysql功能,感興趣的小伙伴可以跟隨小編一起學(xué)習(xí)一下2022-11-11

