我需要实现一个公平的排队系统,以便根据当前排队的消息中该标头的所有值,基于某些消息标头的值以循环方式处理消息.
系统中的消息自然按某些属性分组,其中有数千个可能的值,当前排队的消息的值集随时间而变化.类比将是具有标题的消息,该标题是在消息创建时的毫秒部分时间.因此,标头将具有介于0和999之间的值,并且将在当前排队的所有消息中存在值的一些分布.
我需要能够按顺序使用消息,使得没有特定值优先于任何其他值.如果排队消息的头值是这样分发的
value | count
------|-------
A | 3
B | 3
C | 2
Run Code Online (Sandbox Code Playgroud)
那么消费订单就是A,B,C,A,B,C,A,B.
如果将具有其他值的消息添加到队列,则应将它们自动添加到循环序列中.
这意味着对当前排队的消息有一些了解,但不要求消费者掌握这些知识; 经纪人可能有以某种方式订购交货的机制.
可以接受公平排队开始的某个阈值.也就是说,如果阈值为10,那么顺序处理具有相同值的10个消息是可接受的,但处理的第11个消息应该是顺序的下一个值.如果唯一排队的消息具有该值,则Next可能是相同的值.
可能的值的数量可能排除了简单地为每个队列创建队列,并且迭代队列,尽管尚未经过测试.
我们正在使用HornetQ,但如果有替代方案可以提供这些语义,那么我很想知道.
消息是作业,标头值是用户ID.正在寻求的是,在某些限制内,任何特定用户的任何工作都不会过度拖延任何其他用户的工作; 生成100万个作业的用户不会导致其他用户的后续作业等待处理该百万个作业.
HornetQ队列中的消费者按创建顺序进行评估,因此向队列添加选择性消费者不会阻止任何全能消费者接收与过滤器匹配的消息.
JMS组似乎没有帮助,因为它将给定的组(用户?)绑定到给定的消费者.
一个潜在的解决方案是基于需求(例如:来自同一用户10级连续的消息)在主题创建选择性的消费者,与一些管理所有选择性消费者的生命周期,以确保捕获所有不处理相同的信息.虽然可能这似乎有一些繁重的同步要求.
小智 0
要考虑的第一个选择是拥有一个多线程消费应用程序。假设每个会话/消费者都有一个线程,则可以使用选择器设置同步或异步接收。每个选择器都与特定用户相关。
假设 JVM 在线程分派方面是合理公平的(我很乐意假设)并且应用程序代码中不存在任何死锁,我断言需求将得到满足。一个线程可能会被一百万个用户的作业卡住,其余的不会受到影响。
然而,如果需要单线程应用程序,那么 JMS 规范本身就没有任何帮助。当然,他们可能是供应商扩展,这可能会有所帮助。然而,另一个选项是应用程序查看每条消息并将其放入用户 ID 的特定队列中。最终的消费应用程序本身将在这些队列之间“循环”以获取工作。需要另一个应用程序,但您有一个非常确定的系统。
| 归档时间: |
|
| 查看次数: |
334 次 |
| 最近记录: |