Ste*_*fan 10 java multithreading rabbitmq mongodb couchbase
我有一个Job Distributor人发布不同的消息Channels.
此外,我希望有两个(以及将来更多)Consumers从事不同任务并在不同机器上运行的人.(目前我只有一个,需要扩展它)
让我们来命名这些任务(仅举例):
FIBONACCI (生成斐波纳契数)RANDOMBOOKS (生成随机句子来写一本书)这些任务最长可达2-3小时,应分别平均分配给每个任务Consumer.
每个消费者都可以拥有x 并行线程来处理这些任务.所以我说:(这些数字只是示例,将被变量取代)
FIBONACCI和5个并行作业RANDOMBOOKSFIBONACCI和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拉方法最好.我对这些系统不熟悉,有没有人可以向我解释一下,如何实现这一目标?
这是我如何在couchbase上构建它的概念视图.
总而言之,每个工作者都会对孤立的作业进行查询(如果有的话),检查是否依次检查是否存在锁定文件,如果没有则创建一个并且它遵循上述常规锁定协议.如果没有孤立的作业,则它会查找过期的作业,并遵循锁定协议.如果没有过期作业,那么它只需要最旧的作业并遵循锁定协议.
当然,如果您的系统没有"过期"这样的事情,这也会有效,如果时效性无关紧要,那么您可以使用其他方法代替最旧的工作.
另一种方法可能是在1-N之间创建一个随机值,其中N是一个相当大的数字,比如工人数量的4倍,并且每个工作都用该值标记.每当一个工人去寻找工作时,它就可以掷骰子,看看是否有任何有这个号码的工作.如果没有,它会再次这样做,直到找到具有该号码的作业.这样,代替多个工人竞争少数"最老的"或最高优先级的工作,以及更多的锁定争用的可能性,它们将被分散......以阙中的时间为代价比FIFO情况更随机.
随机方法也可以应用于你必须适应负载值的情况(这样一台机器不会承担太多负载)而不是采用最老的候选者,只需采取随机候选形式可行的工作清单,并尝试这样做.
编辑添加:
在步骤12中,我说"可能输入一个随机数",我的意思是,如果工人知道优先级(例如:哪个人最需要完成工作),他们可以将一个代表这个的数字放入文件中.如果没有"需要"工作的概念,那么他们都可以掷骰子.他们用骰子的角色更新这个文件.然后他们两个都可以看着它,看看对方是怎么回事.如果他们输了,那么他们就会踢,而另一个工人知道它有它.这样,您可以解决哪个工作人员在没有大量复杂协议或协商的情况下完成工作.我假设这两个工作者都在这里击中相同的锁文件,它可以用两个锁文件和一个查找所有这些文件的查询来实现.如果经过一段时间后,没有工人推出更高的数字(并且新员工认为他的工作会知道其他人已经开始工作,所以他们会跳过它)你可以安全地接受工作,因为知道你是唯一的工人正在努力.
首先让我说我没有使用Java与RabbitMQ进行通信,因此我无法提供代码示例.这应该不是问题,因为那不是你要问的问题.这个问题更多的是关于您的应用程序的一般设计.
让我们分解一下,因为这里有很多问题.
一种方法是使用循环法,但这是相当粗糙的,并没有考虑到不同的任务可能需要不同的时间来完成.那么该怎么办.那么一种方法是设置prefetch为1.预取意味着消费者在本地缓存消息(注意:消息尚未消耗).通过将此值设置为1,不会发生预取.这意味着您的消费者只会知道并且只有内存中正在处理的消息.这使得只有在工作人员空闲时才能接收消息.
通过上述设置,可以从队列中读取消息,将其传递给您的某个线程,然后确认该消息.对所有可用线程执行此操作-1.您不想确认最后一条消息,因为这意味着您将打开以接收另一条消息,您将无法将该消息传递给您的某个工作人员.当其中一个线程结束时,那就是当你确认该消息时,这样你就会总是让你的线程处理某些事情.
这取决于你不想做什么,但总的来说,我会说你的制作人应该知道他们传递的是什么.这意味着您可以将其发送到某个交换机,或者更确切地说是将某个路由密钥发送到某个路由密钥,该路由密钥会将此消息传递给正确的队列,该队列将让消费者收听该消息,知道如何处理该消息.
我建议您阅读AMQP和RabbitMQ,这可能是一个很好的起点.
我的提案和您的设计中存在一个主要缺陷,那就是ACK我们在实际处理它之前的消息.这意味着当(不是)我们的应用程序崩溃时,我们无法重新创建ACKed消息.如果您事先知道要启动多少个线程,这可以解决.我不知道你是否可以动态更改预取计数,但不知怎的,我对此表示怀疑.
从RabbitMQ的经验来看,虽然有限,但您不应该害怕创建交换和队列,如果正确完成,这些可以极大地改善和简化您的应用程序设计.也许你不应该有一个启动一堆消费者线程的应用程序.相反,您可能希望使用某种类型的包装器,根据系统中的可用内存或类似内容启动使用者.如果你这样做,你可以确保没有消息丢失,如果你的应用程序崩溃,因为如果你这样做,你当然会在你完成它时确认消息.
如果有什么不清楚,或者我错过了你的观点,请告诉我,如果可以的话,我会尝试扩展我的答案或改进它.