Go+Kafka實(shí)現(xiàn)延遲消息的實(shí)現(xiàn)示例_第1頁
Go+Kafka實(shí)現(xiàn)延遲消息的實(shí)現(xiàn)示例_第2頁
Go+Kafka實(shí)現(xiàn)延遲消息的實(shí)現(xiàn)示例_第3頁
Go+Kafka實(shí)現(xiàn)延遲消息的實(shí)現(xiàn)示例_第4頁
Go+Kafka實(shí)現(xiàn)延遲消息的實(shí)現(xiàn)示例_第5頁
已閱讀5頁,還剩2頁未讀, 繼續(xù)免費(fèi)閱讀

下載本文檔

版權(quán)說明:本文檔由用戶提供并上傳,收益歸屬內(nèi)容提供方,若內(nèi)容存在侵權(quán),請進(jìn)行舉報或認(rèn)領(lǐng)

文檔簡介

第Go+Kafka實(shí)現(xiàn)延遲消息的實(shí)現(xiàn)示例目錄前言原理簡單的實(shí)現(xiàn)生產(chǎn)者延遲服務(wù)消費(fèi)者改進(jìn)點(diǎn)通用的延遲服務(wù)生產(chǎn)者負(fù)責(zé)延遲服務(wù)總結(jié)

前言

延遲隊(duì)列是一個非常有用的工具,我們經(jīng)常遇到需要使用延遲隊(duì)列的場景,比如延遲通知,訂單關(guān)閉等等。

這篇文章主要是使用Go+Kafka實(shí)現(xiàn)延遲消息。

使用了sarama客戶端。

原理

Kafka實(shí)現(xiàn)延遲消息分為下面三步:

生產(chǎn)者把消息發(fā)送到延遲隊(duì)列延遲服務(wù)把延遲隊(duì)列里超過延遲時間的消息寫入真實(shí)隊(duì)列消費(fèi)者消費(fèi)真實(shí)隊(duì)列里的消息

簡單的實(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ù)會訂閱延遲隊(duì)列的消息,并把超時消息發(fā)送到真實(shí)隊(duì)列

iferr=consumerGroup.Consume(context.Background(),

[]string{kafka_delay_queue_test.DelayTopic},consumer);err!=nil{

break

}

typeConsumerstruct{

producersarama.SyncProducer

delaytime.Duration

funcNewConsumer(producersarama.SyncProducer,delaytime.Duration)*Consumer{

returnConsumer{

producer:producer,

delay:delay,

func(c*Consumer)ConsumeClaim(sessionsarama.ConsumerGroupSession,claimsarama.ConsumerGroupClaim)error{

formessage:=rangeclaim.Messages(){

//如果消息已經(jīng)超時,把消息發(fā)送到真實(shí)隊(duì)列

now:=time.Now()

ifnow.Sub(message.Timestamp)=c.delay{

_,_,err:=ducer.SendMessage(sarama.ProducerMessage{

Topic:kafka_delay_queue_test.RealTopic,

Key:sarama.ByteEncoder(message.Key),

Value:sarama.ByteEncoder(message.Value),

iferr==nil{

session.MarkMessage(message,"")

continue

//否則休眠一秒

time.Sleep(time.Second)

returnnil

returnnil

}

消費(fèi)者

消費(fèi)者只是訂閱真實(shí)隊(duì)列并消費(fèi)消息

iferr=consumerGroup.Consume(context.Background(),

[]string{kafka_delay_queue_test.RealTopic},consumer);err!=nil{

break

}

typeConsumerstruct{}

funcNewConsumer()*Consumer{

returnConsumer{}

func(c*Consumer)ConsumeClaim(sessionsarama.ConsumerGroupSession,claimsarama.ConsumerGroupClaim)error{

formessage:=rangeclaim.Messages(){

fmt.Println("收到消息:",message.Value,message.Timestamp)

session.MarkMessage(message,"")

returnnil

}

改進(jìn)點(diǎn)

通用的延遲服務(wù)

可以把延遲服務(wù)封裝成一個通用的服務(wù),這樣生產(chǎn)者可以直接把消息發(fā)送給延遲服務(wù),讓延遲服務(wù)去處理剩下的邏輯。

延遲服務(wù)可以提供多個延時等級,比如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ì)列。

下面是一個簡單的實(shí)現(xiàn):

//KafkaDelayQueueProducer延遲隊(duì)列生產(chǎn)者,包含了生產(chǎn)者和延遲服務(wù)

typeKafkaDelayQueueProducerstruct{

producersarama.SyncProducer//生產(chǎn)者

delayTopicstring//延遲服務(wù)主題

//NewKafkaDelayQueueProducer創(chuàng)建延遲隊(duì)列生產(chǎn)者

//producer生產(chǎn)者

//delayServiceConsumerGroup延遲服務(wù)消費(fèi)者

//delayTime延遲時間

//delayTopic延遲服務(wù)主題

//realTopic真實(shí)隊(duì)列主題

funcNewKafkaDelayQueueProducer(producersarama.SyncProducer,delayServiceConsumerGroupsarama.ConsumerGroup,

delayTimetime.Duration,delayTopic,realTopicstring)*KafkaDelayQueueProducer{

//啟動延遲服務(wù)

consumer:=NewDelayServiceConsumer(producer,delayTime,realTopic)

gofunc(){

for{

iferr:=delayServiceConsumerGroup.Consume(context.Background(),

[]string{delayTopic},consumer);err!=nil{

break

returnKafkaDelayQueueProducer{

producer:producer,

delayTopic:delayTopic,

//SendMessage發(fā)送消息

func(q*KafkaDelayQueueProducer)SendMessage(msg*sarama.ProducerMessage)(partitionint32,offsetint64,errerror){

msg.Topic=q.delayTopic

returnducer.SendMessage(msg)

//DelayServiceConsumer延遲服務(wù)消費(fèi)者

typeDelayServiceConsumerstruct{

producersarama.SyncProducer

delaytime.Duration

realTopicstring

funcNewDelayServiceConsumer(producersarama.SyncProducer,delaytime.Duration,

realTopicstring)*DelayServiceConsumer{

returnDelayServiceConsumer{

producer:producer,

delay:delay,

realTopic:realTopic,

func(c*DelayServiceConsumer)ConsumeClaim(sessionsarama.ConsumerGroupSession,

claimsarama.ConsumerGroupClaim)error{

formessage:=rangeclaim.Messages(){

//如果消息已經(jīng)超時,把消息發(fā)送到真實(shí)隊(duì)列

now:=time.Now()

ifnow.Sub(message.Timestamp)=c.delay{

_,_,err:=ducer.SendMessage(sarama.ProducerMessage{

Topic:c.realTopic,

Key:sarama.ByteEncoder(message.Key),

Value:sarama.ByteEncoder(message.Value),

iferr==nil{

session.MarkMessage(message,"")

continue

//否則休眠一秒

time.Sleep(time.Second)

returnnil

returnnil

func(c*DelayServiceConsumer)Setup(sarama.ConsumerGroupSession)error{

returnnil

func(c*DelayServiceConsumer)Cleanup(sarama.ConsumerGroupSession)error{

溫馨提示

  • 1. 本站所有資源如無特殊說明,都需要本地電腦安裝OFFICE2007和PDF閱讀器。圖紙軟件為CAD,CAXA,PROE,UG,SolidWorks等.壓縮文件請下載最新的WinRAR軟件解壓。
  • 2. 本站的文檔不包含任何第三方提供的附件圖紙等,如果需要附件,請聯(lián)系上傳者。文件的所有權(quán)益歸上傳用戶所有。
  • 3. 本站RAR壓縮包中若帶圖紙,網(wǎng)頁內(nèi)容里面會有圖紙預(yù)覽,若沒有圖紙預(yù)覽就沒有圖紙。
  • 4. 未經(jīng)權(quán)益所有人同意不得將文件中的內(nèi)容挪作商業(yè)或盈利用途。
  • 5. 人人文庫網(wǎng)僅提供信息存儲空間,僅對用戶上傳內(nèi)容的表現(xiàn)方式做保護(hù)處理,對用戶上傳分享的文檔內(nèi)容本身不做任何修改或編輯,并不能對任何下載內(nèi)容負(fù)責(zé)。
  • 6. 下載文件中如有侵權(quán)或不適當(dāng)內(nèi)容,請與我們聯(lián)系,我們立即糾正。
  • 7. 本站不保證下載資源的準(zhǔn)確性、安全性和完整性, 同時也不承擔(dān)用戶因使用這些下載資源對自己和他人造成任何形式的傷害或損失。

評論

0/150

提交評論