<em>Mac</em>Book项目 2009年学校开始实施<em>Mac</em>Book项目,所有师生配备一本<em>Mac</em>Book,并同步更新了校园无线网络。学校每周进行电脑技术更新,每月发送技术支持资料,极大改变了教学及学习方式。因此2011
2021-06-01 09:32:01
延遲佇列是一個非常有用的工具,我們經常遇到需要使用延遲佇列的場景,比如延遲通知,訂單關閉等等。
這篇文章主要是使用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延遲訊息內容請搜尋it145.com以前的文章或繼續瀏覽下面的相關文章希望大家以後多多支援it145.com!
相關文章
<em>Mac</em>Book项目 2009年学校开始实施<em>Mac</em>Book项目,所有师生配备一本<em>Mac</em>Book,并同步更新了校园无线网络。学校每周进行电脑技术更新,每月发送技术支持资料,极大改变了教学及学习方式。因此2011
2021-06-01 09:32:01
综合看Anker超能充系列的性价比很高,并且与不仅和iPhone12/苹果<em>Mac</em>Book很配,而且适合多设备充电需求的日常使用或差旅场景,不管是安卓还是Switch同样也能用得上它,希望这次分享能给准备购入充电器的小伙伴们有所
2021-06-01 09:31:42
除了L4WUDU与吴亦凡已经多次共事,成为了明面上的厂牌成员,吴亦凡还曾带领20XXCLUB全队参加2020年的一场音乐节,这也是20XXCLUB首次全员合照,王嗣尧Turbo、陈彦希Regi、<em>Mac</em> Ova Seas、林渝植等人全部出场。然而让
2021-06-01 09:31:34
目前应用IPFS的机构:1 谷歌<em>浏览器</em>支持IPFS分布式协议 2 万维网 (历史档案博物馆)数据库 3 火狐<em>浏览器</em>支持 IPFS分布式协议 4 EOS 等数字货币数据存储 5 美国国会图书馆,历史资料永久保存在 IPFS 6 加
2021-06-01 09:31:24
开拓者的车机是兼容苹果和<em>安卓</em>,虽然我不怎么用,但确实兼顾了我家人的很多需求:副驾的门板还配有解锁开关,有的时候老婆开车,下车的时候偶尔会忘记解锁,我在副驾驶可以自己开门:第二排设计很好,不仅配置了一个很大的
2021-06-01 09:30:48
不仅是<em>安卓</em>手机,苹果手机的降价力度也是前所未有了,iPhone12也“跳水价”了,发布价是6799元,如今已经跌至5308元,降价幅度超过1400元,最新定价确认了。iPhone12是苹果首款5G手机,同时也是全球首款5nm芯片的智能机,它
2021-06-01 09:30:45