版權說明:本文檔由用戶提供并上傳,收益歸屬內(nèi)容提供方,若內(nèi)容存在侵權,請進行舉報或認領
文檔簡介
第Go+Kafka實現(xiàn)延遲消息的實現(xiàn)示例目錄前言原理簡單的實現(xiàn)生產(chǎn)者延遲服務消費者改進點通用的延遲服務生產(chǎn)者負責延遲服務總結
前言
延遲隊列是一個非常有用的工具,我們經(jīng)常遇到需要使用延遲隊列的場景,比如延遲通知,訂單關閉等等。
這篇文章主要是使用Go+Kafka實現(xiàn)延遲消息。
使用了sarama客戶端。
原理
Kafka實現(xiàn)延遲消息分為下面三步:
生產(chǎn)者把消息發(fā)送到延遲隊列延遲服務把延遲隊列里超過延遲時間的消息寫入真實隊列消費者消費真實隊列里的消息
簡單的實現(xiàn)
生產(chǎn)者
生產(chǎn)者只是把消息發(fā)送到延遲隊列
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)
}
延遲服務
延遲服務會訂閱延遲隊列的消息,并把超時消息發(fā)送到真實隊列
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ā)送到真實隊列
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
}
消費者
消費者只是訂閱真實隊列并消費消息
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
}
改進點
通用的延遲服務
可以把延遲服務封裝成一個通用的服務,這樣生產(chǎn)者可以直接把消息發(fā)送給延遲服務,讓延遲服務去處理剩下的邏輯。
延遲服務可以提供多個延時等級,比如5s、10s、30s、1m、5m、10m、1h、2h等,類似于RocketMQ。
生產(chǎn)者負責延遲服務
也可以讓生產(chǎn)者負責延遲服務,讓生產(chǎn)者自己把延遲隊列里面的消息發(fā)送到真實隊列。
下面是一個簡單的實現(xiàn):
//KafkaDelayQueueProducer延遲隊列生產(chǎn)者,包含了生產(chǎn)者和延遲服務
typeKafkaDelayQueueProducerstruct{
producersarama.SyncProducer//生產(chǎn)者
delayTopicstring//延遲服務主題
//NewKafkaDelayQueueProducer創(chuàng)建延遲隊列生產(chǎn)者
//producer生產(chǎn)者
//delayServiceConsumerGroup延遲服務消費者
//delayTime延遲時間
//delayTopic延遲服務主題
//realTopic真實隊列主題
funcNewKafkaDelayQueueProducer(producersarama.SyncProducer,delayServiceConsumerGroupsarama.ConsumerGroup,
delayTimetime.Duration,delayTopic,realTopicstring)*KafkaDelayQueueProducer{
//啟動延遲服務
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延遲服務消費者
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ā)送到真實隊列
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)系上傳者。文件的所有權益歸上傳用戶所有。
- 3. 本站RAR壓縮包中若帶圖紙,網(wǎng)頁內(nèi)容里面會有圖紙預覽,若沒有圖紙預覽就沒有圖紙。
- 4. 未經(jīng)權益所有人同意不得將文件中的內(nèi)容挪作商業(yè)或盈利用途。
- 5. 人人文庫網(wǎng)僅提供信息存儲空間,僅對用戶上傳內(nèi)容的表現(xiàn)方式做保護處理,對用戶上傳分享的文檔內(nèi)容本身不做任何修改或編輯,并不能對任何下載內(nèi)容負責。
- 6. 下載文件中如有侵權或不適當內(nèi)容,請與我們聯(lián)系,我們立即糾正。
- 7. 本站不保證下載資源的準確性、安全性和完整性, 同時也不承擔用戶因使用這些下載資源對自己和他人造成任何形式的傷害或損失。
最新文檔
- 警務室調解制度
- 用電基礎知識培訓
- 2025高一政治期末模擬卷01(考試版)【測試范圍:必修1全冊+必修2全冊】(新高考用)含答案
- 醫(yī)院愛崗敬業(yè)培訓課件
- 國考公安考試試題及答案
- 2026年上半年浙江杭州市婦產(chǎn)科醫(yī)院(杭州市婦幼保健院)高層次、緊缺專業(yè)人才招聘15人(總)備考考試試題附答案解析
- 2026某事業(yè)單位招聘保潔崗位1人備考考試題庫附答案解析
- JIS D 9101-2012 自行車術語標準 Cycles - Terminology
- 2026福建福州市平潭綜合實驗區(qū)黨工委黨校(區(qū)行政學院、區(qū)社會主義學院)招聘編外工作人員1人備考考試題庫附答案解析
- 2026福建龍巖鑫達彩印有限公司龍巖鑫利來酒店分公司(第一批)招聘3人參考考試試題附答案解析
- 2025屆高考小說專題復習-小說敘事特征+課件
- 部編版二年級下冊寫字表字帖(附描紅)
- 干部履歷表(中共中央組織部2015年制)
- GB/T 5657-2013離心泵技術條件(Ⅲ類)
- GB/T 3518-2008鱗片石墨
- GB/T 17622-2008帶電作業(yè)用絕緣手套
- GB/T 1041-2008塑料壓縮性能的測定
- 400份食物頻率調查問卷F表
- 滑坡地質災害治理施工
- 實驗動物從業(yè)人員上崗證考試題庫(含近年真題、典型題)
- 可口可樂-供應鏈管理
評論
0/150
提交評論