Java和RabbitMQ - 排队和多线程 - 或Couchbase作为作业队列

Ste*_*fan 10 java multithreading rabbitmq mongodb couchbase

我有一个Job Distributor人发布不同的消息Channels.

此外,我希望有两个(以及将来更多)Consumers从事不同任务并在不同机器上运行的人.(目前我只有一个,需要扩展它)

让我们来命名这些任务(仅举例):

  • FIBONACCI (生成斐波纳契数)
  • RANDOMBOOKS (生成随机句子来写一本书)

这些任务最长可达2-3小时,应分别平均分配给每个任务Consumer.

每个消费者都可以拥有x 并行线程来处理这些任务.所以我说:(这些数字只是示例,将被变量取代)

  • 机器1可以使用3个并行作业FIBONACCI和5个并行作业RANDOMBOOKS
  • 机器2可以使用7个并行作业FIBONACCI和3个并行作业RANDOMBOOKS

我怎样才能实现这一目标?

我是否必须x为每个人启动Threads Channel来监听每个Consumer

我何时需要确认?

我目前只有一种方法Consumer是:x为每个任务启动线程 - 每个线程都是一个Defaultconsumer实现Runnable.在handleDelivery方法中,我打电话basicAck(deliveryTag,false)然后做工作.

进一步:我想将一些任务发送给特殊的消费者.如何结合上述公平分配实现这一目标?

这是我的代码 publishing

String QUEUE_NAME = "FIBONACCI";

Channel channel = this.clientManager.getRabbitMQConnection().createChannel();

channel.queueDeclare(QUEUE_NAME, true, false, false, null);

channel.basicPublish("", QUEUE_NAME,
                MessageProperties.BASIC,
                Control.getBytes(this.getArgument()));

channel.close();
Run Code Online (Sandbox Code Playgroud)

这是我的代码 Consumer

public final class Worker extends DefaultConsumer implements Runnable {
    @Override
    public void run() {

        try {
            this.getChannel().queueDeclare(this.jobType.toString(), true, false, false, null);
            this.getChannel().basicConsume(this.jobType.toString(), this);

            this.getChannel().basicQos(1);
        } catch (IOException e) {
            // catch something
        }
        while (true) {
            try {
                Thread.sleep(1000);
            } catch (InterruptedException e) {
                Control.getLogger().error("Exception!", e);
            }

        }
    }

    @Override
    public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] bytes) throws IOException {
        String routingKey = envelope.getRoutingKey();
        String contentType = properties.getContentType();
        this.getChannel().basicAck(deliveryTag, false); // Is this right?
        // Start new Thread for this task with my own ExecutorService

    }
}
Run Code Online (Sandbox Code Playgroud)

Worker在这种情况下,该类启动两次:一次为FIBUNACCI一次,一次为RANDOMBOOKS

UPDATE

正如答案所述,RabbitMQ不是最好的解决方案,但Couchbase或MongoDB拉方法最好.我对这些系统不熟悉,有没有人可以向我解释一下,如何实现这一目标?

nir*_*ana 7

这是我如何在couchbase上构建它的概念视图.

  1. 你有一些机器可以处理作业,还有一些机器(可能是相同的机器)可以创建工作.
  2. 您可以在couchbase中的存储桶中为每个作业创建一个文档(并将其类型设置为"作业",或者如果您将其与该存储桶中的其他数据混合在一起).
  3. 每个作业描述以及要完成的特定命令可以包括创建时间,到期时间(如果有特定时间)以及某种生成的工作值.该工作值将是任意单位.
  4. 每个工作的消费者都会知道它一次可以做多少个工作单位,有多少可用(因为其他工人可能正在工作).
  5. 因此,具有10个工作单元的容量的机器(其具有6个工作单元)将进行查询以寻找4个或更少工作单元的工作.
  6. 在couchbase中有一些视图,这些视图是逐步更新的map/reduce作业,我认为你只需要这里的地图阶段.您可以编写一个视图,通过该视图查询到期时间,输入系统的时间和工作单位数.通过这种方式,您可以获得"4个或更少工作单位的最迟期工作".
  7. 这种查询,随着容量的释放,将首先得到最多的过期工作,尽管你可以获得最大的过期工作,如果没有,那么最大的未逾期工作.(其中"逾期"是当前时间与工作到期日之间的差值.)
  8. Couchbase视图允许这样非常复杂的查询.虽然它们逐步更新,但它们并非完全实时.因此,您不会寻找一份工作,而是一份求职者名单(无论您希望如何订购).
  9. 因此,下一步是获取候选作业列表并检查第二个位置 - 可能是锁定文件的膜库(例如:RAM缓存,非持久性).锁文件将有多个阶段(这里你使用CRDT做一些分区解析逻辑或任何最适合你需要的方法.)
  10. 由于这个桶是基于ram的,它比视图更快,并且从总状态开始具有更少的延迟.如果没有锁定文件,则创建一个状态标志为"临时"的文件.
  11. 如果另一个工作人员获得相同的工作并看到锁定文件,那么它可以跳过该候选人并在列表中执行下一个工作.
  12. 如果不知何故,两名工人试图为同一个工作创建锁文件,就会发生冲突.在冲突的情况下,你可以踢.或者您可以拥有逻辑,其中每个工作人员对锁定文件进行更新(CRDT解析因此使这些幂等元素可以合并兄弟姐妹)可能放入随机数或某个优先级数字.
  13. 在指定的时间段(可能是几秒钟)之后,工作人员检查锁定文件,如果它不必参与任何种族分辨率更改,它会将锁定文件的状态从"临时"更改为"已采取" "
  14. 然后,它会以"已采取"状态或某些状态更新作业本身,以便在其他工作人员查找可用作业时不会显示在视图中.
  15. 最后,您需要添加另一个步骤,在执行查询之前获取上面描述的这些求职者,您会进行特殊查询以查找已执行的作业,但涉及的工作人员已经死亡.(例如:过期的工作).
  16. 了解工人何时死亡的一种方法是,放入membase存储桶的锁文件应该有一个到期时间,最终会导致它消失.可能这个时间可能很短,工人只需触摸它就可以更新到期时间(这在couchbase API中得到支持)
  17. 如果一个工人死了,最终它的锁定文件将消失,孤立的工作将被标记为"已被",但没有锁定文件,这是寻找工作的工人可以寻找的条件.

总而言之,每个工作者都会对孤立的作业进行查询(如果有的话),检查是否依次检查是否存在锁定文件,如果没有则创建一个并且它遵循上述常规锁定协议.如果没有孤立的作业,则它会查找过期的作业,并遵循锁定协议.如果没有过期作业,那么它只需要最旧的作业并遵循锁定协议.

当然,如果您的系统没有"过期"这样的事情,这也会有效,如果时效性无关紧要,那么您可以使用其他方法代替最旧的工作.

另一种方法可能是在1-N之间创建一个随机值,其中N是一个相当大的数字,比如工人数量的4倍,并且每个工作都用该值标记.每当一个工人去寻找工作时,它就可以掷骰子,看看是否有任何有这个号码的工作.如果没有,它会再次这样做,直到找到具有该号码的作业.这样,代替多个工人竞争少数"最老的"或最高优先级的工作,以及更多的锁定争用的可能性,它们将被分散......以阙中的时间为代价比FIFO情况更随机.

随机方法也可以应用于你必须适应负载值的情况(这样一台机器不会承担太多负载)而不是采用最老的候选者,只需采取随机候选形式可行的工作清单,并尝试这样做.

编辑添加:

在步骤12中,我说"可能输入一个随机数",我的意思是,如果工人知道优先级(例如:哪个人最需要完成工作),他们可以将一个代表这个的数字放入文件中.如果没有"需要"工作的概念,那么他们都可以掷骰子.他们用骰子的角色更新这个文件.然后他们两个都可以看着它,看看对方是怎么回事.如果他们输了,那么他们就会踢,而另一个工人知道它有它.这样,您可以解决哪个工作人员在没有大量复杂协议或协商的情况下完成工作.我假设这两个工作者都在这里击中相同的锁文件,它可以用两个锁文件和一个查找所有这些文件的查询来实现.如果经过一段时间后,没有工人推出更高的数字(并且新员工认为他的工作会知道其他人已经开始工作,所以他们会跳过它)你可以安全地接受工作,因为知道你是唯一的工人正在努力.


Dan*_*roa 5

首先让我说我没有使用Java与RabbitMQ进行通信,因此我无法提供代码示例.这应该不是问题,因为那不是你要问的问题.这个问题更多的是关于您的应用程序的一般设计.

让我们分解一下,因为这里有很多问题.

将任务划分给不同的消费者

一种方法是使用循环法,但这是相当粗糙的,并没有考虑到不同的任务可能需要不同的时间来完成.那么该怎么办.那么一种方法是设置prefetch1.预取意味着消费者在本地缓存消息(注意:消息尚未消耗).通过将此值设置为1,不会发生预取.这意味着您的消费者只会知道并且只有内存中正在处理的消息.这使得只有在工作人员空闲时才能接收消息.

什么时候承认

通过上述设置,可以从队列中读取消息,将其传递给您的某个线程,然后确认该消息.对所有可用线程执行此操作-1.您不想确认最后一条消息,因为这意味着您将打开以接收另一条消息,您将无法将该消息传递给您的某个工作人员.当其中一个线程结束时,那就是当你确认该消息时,这样你就会总是让你的线程处理某些事情.

传递特殊信息

这取决于你不想做什么,但总的来说,我会说你的制作人应该知道他们传递的是什么.这意味着您可以将其发送到某个交换机,或者更确切地说是将某个路由密钥发送到某个路由密钥,该路由密钥会将此消息传递给正确的队列,该队列将让消费者收听该消息,知道如何处理该消息.

我建议您阅读AMQP和RabbitMQ,这可能是一个很好的起点.

注意事项

我的提案和您的设计中存在一个主要缺陷,那就是ACK我们在实际处理它之前的消息.这意味着当(不是)我们的应用程序崩溃时,我们无法重新创建ACKed消息.如果您事先知道要启动多少个线程,这可以解决.我不知道你是否可以动态更改预取计数,但不知怎的,我对此表示怀疑.

一些想法

从RabbitMQ的经验来看,虽然有限,但您不应该害怕创建交换和队列,如果正确完成,这些可以极大地改善和简化您的应用程序设计.也许你不应该有一个启动一堆消费者线程的应用程序.相反,您可能希望使用某种类型的包装器,根据系统中的可用内存或类似内容启动使用者.如果你这样做,你可以确保没有消息丢失,如果你的应用程序崩溃,因为如果你这样做,你当然会在你完成它时确认消息.

推荐阅读

如果有什么不清楚,或者我错过了你的观点,请告诉我,如果可以的话,我会尝试扩展我的答案或改进它.