NSQ源码分析(四)——inFlightPqueue和PriorityQueue优先级队列

2023-12-16 16:32

本文主要是介绍NSQ源码分析(四)——inFlightPqueue和PriorityQueue优先级队列,希望对大家解决编程问题提供一定的参考价值,需要的开发者们随着小编来一起学习吧!

在Channel结构体中用到了两种优先级队列pqueue.PriorityQueue和inFlightPqueue。

deferredMessages map[MessageID]*pqueue.Item
deferredPQ       pqueue.PriorityQueue
deferredMutex    sync.MutexinFlightMessages map[MessageID]*Message
inFlightPQ       inFlightPqueue
inFlightMutex    sync.Mutex

       其中deferredMessages和inFlightMessages 使用map储存了MessageID和Message的对应关系,用于根据MessageID获取对应的Message。而deferredPQ和inFlightPQ是两种优先级队列。deferredPQ队列中储存了延时消息和消息投递失败需要等待指定时间后重新投递的消息,inFlightPQ队列中储存了正在投递但还没确认投递成功的消息。

 

一、PriorityQueue优先级队列

       源码位置在 nsq/internal/pqueue/pqueue.go文件中

      

PriorityQueue实现了Golang源码包heap中的接口,是最小堆。在了解PriorityQueue之前需要先对堆的概念有了解,可以参考:

https://blog.csdn.net/skh2015java/article/details/83183681

PriorityQueue队列中储存的是Item指针,Item结构体及字段说明如下:

type Item struct {Value interface{} //储存的消息内容Priority int64 //优先级的时间点Index int //在切片中的索引值}

   这三个方法用于排序,按Priority字段排序

func (pq PriorityQueue) Len() int {return len(pq)}func (pq PriorityQueue) Less(i, j int) bool {return pq[i].Priority < pq[j].Priority}func (pq PriorityQueue) Swap(i, j int) {pq[i], pq[j] = pq[j], pq[i]pq[i].Index = ipq[j].Index = j}

  常用方法:

  

    New函数用于初始化队列并指定cap

   
func New(capacity int) PriorityQueue {return make(PriorityQueue, 0, capacity)}

  Push函数用于向队列中添加元素

  Pop函数用于取出队列中优先级最高的元素(即Priority值最小,也就是根部元素)

  PeekAndShift(max int64) 用于判断根部的元素是否超过max,如果超过则返回nil,如果没有则返回并移除根部元素(根部元素是最小值)

    在实际的项目中Priority字段存的是时间戳,比如说5分钟之后投递本条消息,则Priority字段存的就是5分钟之后的时间戳。而PeekAndShift(max int64)中max值就是当前时间,如果队列中的根部元素大于当前时间戳max的值,则说明队列中没有可以投递的消息,故返回nil。如果小于等于,则根部元素存放的消息可以投递,就是就返回并移除该根部元素。

func (pq *PriorityQueue) PeekAndShift(max int64) (*Item, int64) {if pq.Len() == 0 {return nil, 0}item := (*pq)[0] //获取根部元素if item.Priority > max {return nil, item.Priority - max}heap.Remove(pq, 0) //移除根部元素,并重新排列堆,使根部元素为最小值return item, 0}

 

  二、inFlightPqueue优先级队列

     

源码位置:nsq/nsqd/in_flight_pqueue.go文件中

  inFlightPqueue队列中存放的元素是Message的指针,也是最小堆,和PriorityQueue队列类似。inFlightPqueue队列按Message中的pri字段进行排序的(pri也是时间戳,是投递消息的超时时间)

  常用方法:

     Swap函数用于交换两个元素的位置,索 引值也随之改变

     Push函数向队列中添加元素

     Pop函数移除队列中的根部元素

     Remove(i int)函数移除队列中指点位置的元素

     PeekAndShit函数用于判断根部的元素是否超过max,如果超过则返回nil,如果没有则返回并移除根部元素(根部元素是最小值)

本节主要对两种队列进行了解和学习,下篇的Channel会进一步探讨这两种队列中消息的传递。

这篇关于NSQ源码分析(四)——inFlightPqueue和PriorityQueue优先级队列的文章就介绍到这儿,希望我们推荐的文章对编程师们有所帮助!



http://www.chinasem.cn/article/501136

相关文章

SpringKafka错误处理(重试机制与死信队列)

《SpringKafka错误处理(重试机制与死信队列)》SpringKafka提供了全面的错误处理机制,通过灵活的重试策略和死信队列处理,下面就来介绍一下,具有一定的参考价值,感兴趣的可以了解一下... 目录引言一、Spring Kafka错误处理基础二、配置重试机制三、死信队列实现四、特定异常的处理策略五

Python 迭代器和生成器概念及场景分析

《Python迭代器和生成器概念及场景分析》yield是Python中实现惰性计算和协程的核心工具,结合send()、throw()、close()等方法,能够构建高效、灵活的数据流和控制流模型,这... 目录迭代器的介绍自定义迭代器省略的迭代器生产器的介绍yield的普通用法yield的高级用法yidle

C++ Sort函数使用场景分析

《C++Sort函数使用场景分析》sort函数是algorithm库下的一个函数,sort函数是不稳定的,即大小相同的元素在排序后相对顺序可能发生改变,如果某些场景需要保持相同元素间的相对顺序,可使... 目录C++ Sort函数详解一、sort函数调用的两种方式二、sort函数使用场景三、sort函数排序

Java调用C++动态库超详细步骤讲解(附源码)

《Java调用C++动态库超详细步骤讲解(附源码)》C语言因其高效和接近硬件的特性,时常会被用在性能要求较高或者需要直接操作硬件的场合,:本文主要介绍Java调用C++动态库的相关资料,文中通过代... 目录一、直接调用C++库第一步:动态库生成(vs2017+qt5.12.10)第二步:Java调用C++

kotlin中const 和val的区别及使用场景分析

《kotlin中const和val的区别及使用场景分析》在Kotlin中,const和val都是用来声明常量的,但它们的使用场景和功能有所不同,下面给大家介绍kotlin中const和val的区别,... 目录kotlin中const 和val的区别1. val:2. const:二 代码示例1 Java

Go标准库常见错误分析和解决办法

《Go标准库常见错误分析和解决办法》Go语言的标准库为开发者提供了丰富且高效的工具,涵盖了从网络编程到文件操作等各个方面,然而,标准库虽好,使用不当却可能适得其反,正所谓工欲善其事,必先利其器,本文将... 目录1. 使用了错误的time.Duration2. time.After导致的内存泄漏3. jsO

Python实现无痛修改第三方库源码的方法详解

《Python实现无痛修改第三方库源码的方法详解》很多时候,我们下载的第三方库是不会有需求不满足的情况,但也有极少的情况,第三方库没有兼顾到需求,本文将介绍几个修改源码的操作,大家可以根据需求进行选择... 目录需求不符合模拟示例 1. 修改源文件2. 继承修改3. 猴子补丁4. 追踪局部变量需求不符合很

Spring事务中@Transactional注解不生效的原因分析与解决

《Spring事务中@Transactional注解不生效的原因分析与解决》在Spring框架中,@Transactional注解是管理数据库事务的核心方式,本文将深入分析事务自调用的底层原理,解释为... 目录1. 引言2. 事务自调用问题重现2.1 示例代码2.2 问题现象3. 为什么事务自调用会失效3

找不到Anaconda prompt终端的原因分析及解决方案

《找不到Anacondaprompt终端的原因分析及解决方案》因为anaconda还没有初始化,在安装anaconda的过程中,有一行是否要添加anaconda到菜单目录中,由于没有勾选,导致没有菜... 目录问题原因问http://www.chinasem.cn题解决安装了 Anaconda 却找不到 An

Spring定时任务只执行一次的原因分析与解决方案

《Spring定时任务只执行一次的原因分析与解决方案》在使用Spring的@Scheduled定时任务时,你是否遇到过任务只执行一次,后续不再触发的情况?这种情况可能由多种原因导致,如未启用调度、线程... 目录1. 问题背景2. Spring定时任务的基本用法3. 为什么定时任务只执行一次?3.1 未启用