rabbimq-java(com.rabbitmq.amqp-client)

系列: http://zxb1985.iteye.com/category/267524

rabbitmq学习2:Work Queues

在前面的已经提到了一对一的情况;现在一个生产者与多个消费者的情况(Work Queues)。

Work Queues的示意图如下:

对于上图的模型中对于c端的worker来说。RabbitMQ服务器可能一直发送多个消息给一个worker,而另一个可能几乎不做任何事情。这样就会导致一个worker很忙,而另一个却很空闲。这种情况可能都不想出现。如何解决这个问题呢。当然最理想的情况是均匀分配消息给每个worker。我们可能通过channel.basicQos(1)方法(prefetchCount=1)来设置同一时间每次发给一个消息给一个worker。示意图如下:

P端的程序如下:

Java代码

packagecom.abin.rabbitmq;

importcom.rabbitmq.client.Channel;

importcom.rabbitmq.client.Connection;

importcom.rabbitmq.client.ConnectionFactory;

importcom.rabbitmq.client.MessageProperties;

publicclassNewTask?{

privatestaticfinalString?TASK_QUEUE_NAME?="task_queue";

publicstaticvoidmain(String[]?argv)throwsException?{

ConnectionFactory?factory?=newConnectionFactory();

factory.setHost("localhost");

Connection?connection?=?factory.newConnection();

Channel?channel?=?connection.createChannel();

//声明此队列并且持久化

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

String?message?=?getMessage(argv);

channel.basicPublish("",?TASK_QUEUE_NAME,

MessageProperties.PERSISTENT_TEXT_PLAIN,?message.getBytes());//持久化消息

System.out.println("?[x]?Sent?'"+?message?+"'");

channel.close();

connection.close();

}

privatestaticString?getMessage(String[]?strings)?{

if(strings.length?<1)

return"Hello?World!";

returnjoinStrings(strings,"?");

}

privatestaticString?joinStrings(String[]?strings,?String?delimiter)?{

intlength?=?strings.length;

if(length?==0)

return"";

StringBuilder?words?=newStringBuilder(strings[0]);

for(inti?=1;?i?<?length;?i++)?{

words.append(delimiter).append(strings[i]);

}

returnwords.toString();

}

}

多次运行此程序并传入的参数分别为“First message”,“Secondmessage”,“Thirdmessage”,“Fourth message”,“Fifth message”

C端的程序如下:

Java代

packagecom.abin.rabbitmq;

importcom.rabbitmq.client.Channel;

importcom.rabbitmq.client.Connection;

importcom.rabbitmq.client.ConnectionFactory;

importcom.rabbitmq.client.QueueingConsumer;

publicclassWorker?{

privatestaticfinalString?TASK_QUEUE_NAME?="task_queue";

publicstaticvoidmain(String[]?argv)throwsException?{

ConnectionFactory?factory?=newConnectionFactory();

factory.setHost("localhost");

Connection?connection?=?factory.newConnection();

Channel?channel?=?connection.createChannel();

//声明此队列并且持久化

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

System.out.println("?[*]?Waiting?for?messages.?To?exit?press?CTRL+C");

channel.basicQos(1);//告诉RabbitMQ同一时间给一个消息给消费者

/*?We're?about?to?tell?the?server?to?deliver?us?the?messages?from?the?queue.

*?Since?it?will?push?us?messages?asynchronously,

*?we?provide?a?callback?in?the?form?of?an?object?that?will?buffer?the?messages

*?until?we're?ready?to?use?them.?That?is?what?QueueingConsumer?does.*/

QueueingConsumer?consumer?=newQueueingConsumer(channel);

/*

把名字为TASK_QUEUE_NAME的Channel的值回调给QueueingConsumer,即使一个worker在处理消息的过程中停止了,这个消息也不会失效

*/

channel.basicConsume(TASK_QUEUE_NAME,false,?consumer);

while(true)?{

QueueingConsumer.Delivery?delivery?=?consumer.nextDelivery();//得到消息传输信息

String?message?=newString(delivery.getBody());

System.out.println("?[x]?Received?'"+?message?+"'");

doWork(message);

System.out.println("?[x]?Done");

channel.basicAck(delivery.getEnvelope().getDeliveryTag(),false);//下一个消息

}

}

privatestaticvoiddoWork(String?task)throwsInterruptedException?{

for(charch?:?task.toCharArray())?{

if(ch?=='.')

Thread.sleep(1000);//这里是假装我们很忙

}

}

}

开启两个worker分别运行。运行结果如:

c1的结果:

Java代码

[*]?Waitingformessages.?To?exit?press?CTRL+C

[x]?Received'First?message'

[x]?Received'Third?message'

[x]?Received'Fifth?message'

c2的结果

Java代码

[*]?Waitingformessages.?To?exit?press?CTRL+C

[x]?Received'Second?message'

[x]?Received'Fourth?message'


rabbitmq学习3:Publish/Subscribe

在前面的Work Queue中的消息是均匀分配消息给消费者;如果我想把消息分发给所有的消费者呢?那应当怎么操作呢?这就是要下面提到的Publish/Subscribe(分布/订阅)。让我们开始Publish/Subscribe之旅吧!

Publish/Subscribe的工作示意图如下:

在上图中的X表示Exchange(交换区);Exchange的类型有:direct,topic,headers和fanout

Publish/Subscribe的Exchang的类型为fanout;声明Publish/Subscribe的Exchang代码如下:

Java代码

channel.exchangeDeclare("logs","fanout");

对于Work Queue中提到的发布消息的代码如下:

Java代码

channel.basicPublish("",?queueName,null,?message.getBytes());

但对于Publish/Subscribe中发布消息中的Queue的使用的是默认的;代码如下:

Java代码

channel.basicPublish("logs","",null,?message.getBytes());

Exchange和各Queue之间是如何通信的呢?主要是通过把Exchange和各Queue绑定(binding);示意代码如下:

Java代码

channel.queueBind(queueName,?exchangeName,"");

Publish/Subscribe加入绑定的工作示意图如下:

那我们就开始程序代码吧:P端的代码如下:

Java代码

packagecom.abin.rabbitmq;

importcom.rabbitmq.client.Channel;

importcom.rabbitmq.client.Connection;

importcom.rabbitmq.client.ConnectionFactory;

publicclassEmitLog?{

privatestaticfinalString?EXCHANGE_NAME?="logs";

publicstaticvoidmain(String[]?argv)throwsException?{

ConnectionFactory?factory?=newConnectionFactory();

factory.setHost("localhost");

Connection?connection?=?factory.newConnection();

Channel?channel?=?connection.createChannel();

channel.exchangeDeclare(EXCHANGE_NAME,"fanout");//声明Exchange

for(inti?=0;?i?<=2;?i++)?{

String?message?="hello?word!"+?i;

channel.basicPublish(EXCHANGE_NAME,"",null,?message.getBytes());

System.out.println("?[x]?Sent?'"+?message?+"'");

}

channel.close();

connection.close();

}

}

运行结果如下:

Java代码

[x]?Sent'hello?word!0'

[x]?Sent'hello?word!1'

[x]?Sent'hello?word!2'

C端的代码如下:

Java代码

packagecom.abin.rabbitmq;

importcom.rabbitmq.client.Channel;

importcom.rabbitmq.client.Connection;

importcom.rabbitmq.client.ConnectionFactory;

importcom.rabbitmq.client.QueueingConsumer;

publicclassReceiveLogsOne?{

privatestaticfinalString?EXCHANGE_NAME?="logs";

publicstaticvoidmain(String[]?argv)throwsException?{

ConnectionFactory?factory?=newConnectionFactory();

factory.setHost("localhost");

Connection?connection?=?factory.newConnection();

Channel?channel?=?connection.createChannel();

channel.exchangeDeclare(EXCHANGE_NAME,"fanout");

String?queueName?="log-fb1";

channel.queueDeclare(queueName,false,false,false,null);

channel.queueBind(queueName,?EXCHANGE_NAME,"");//把Queue、Exchange绑定

QueueingConsumer?consumer?=newQueueingConsumer(channel);

channel.basicConsume(queueName,true,?consumer);

while(true)?{

QueueingConsumer.Delivery?delivery?=?consumer.nextDelivery();

String?message?=newString(delivery.getBody());

System.out.println("?[x]?Received?'"+?message?+"'");

}

}

}

对于C端的代码我写了二个差不多的程序,只需要修改一下queueName。这样就形成了二个Queue;运行结果相同;

运行结果可能如下:

Java代

[x]?Received'hello?word!0'

[x]?Received'hello?word!1'

[x]?Received'hello?word!2'

最后编辑于
?著作权归作者所有,转载或内容合作请联系作者
  • 序言:七十年代末,一起剥皮案震惊了整个滨河市,随后出现的几起案子,更是在滨河造成了极大的恐慌,老刑警刘岩,带你破解...
    沈念sama阅读 214,128评论 6 493
  • 序言:滨河连续发生了三起死亡事件,死亡现场离奇诡异,居然都是意外死亡,警方通过查阅死者的电脑和手机,发现死者居然都...
    沈念sama阅读 91,316评论 3 388
  • 文/潘晓璐 我一进店门,熙熙楼的掌柜王于贵愁眉苦脸地迎上来,“玉大人,你说我怎么就摊上这事?!?“怎么了?”我有些...
    开封第一讲书人阅读 159,737评论 0 349
  • 文/不坏的土叔 我叫张陵,是天一观的道长。 经常有香客问我,道长,这世上最难降的妖魔是什么? 我笑而不...
    开封第一讲书人阅读 57,283评论 1 287
  • 正文 为了忘掉前任,我火速办了婚礼,结果婚礼上,老公的妹妹穿的比我还像新娘。我一直安慰自己,他们只是感情好,可当我...
    茶点故事阅读 66,384评论 6 386
  • 文/花漫 我一把揭开白布。 她就那样静静地躺着,像睡着了一般。 火红的嫁衣衬着肌肤如雪。 梳的纹丝不乱的头发上,一...
    开封第一讲书人阅读 50,458评论 1 292
  • 那天,我揣着相机与录音,去河边找鬼。 笑死,一个胖子当着我的面吹牛,可吹牛的内容都是我干的。 我是一名探鬼主播,决...
    沈念sama阅读 39,467评论 3 412
  • 文/苍兰香墨 我猛地睁开眼,长吁一口气:“原来是场噩梦啊……” “哼!你这毒妇竟也来了?” 一声冷哼从身侧响起,我...
    开封第一讲书人阅读 38,251评论 0 269
  • 序言:老挝万荣一对情侣失踪,失踪者是张志新(化名)和其女友刘颖,没想到半个月后,有当地人在树林里发现了一具尸体,经...
    沈念sama阅读 44,688评论 1 306
  • 正文 独居荒郊野岭守林人离奇死亡,尸身上长有42处带血的脓包…… 初始之章·张勋 以下内容为张勋视角 年9月15日...
    茶点故事阅读 36,980评论 2 328
  • 正文 我和宋清朗相恋三年,在试婚纱的时候发现自己被绿了。 大学时的朋友给我发了我未婚夫和他白月光在一起吃饭的照片。...
    茶点故事阅读 39,155评论 1 342
  • 序言:一个原本活蹦乱跳的男人离奇死亡,死状恐怖,灵堂内的尸体忽然破棺而出,到底是诈尸还是另有隐情,我是刑警宁泽,带...
    沈念sama阅读 34,818评论 4 337
  • 正文 年R本政府宣布,位于F岛的核电站,受9级特大地震影响,放射性物质发生泄漏。R本人自食恶果不足惜,却给世界环境...
    茶点故事阅读 40,492评论 3 322
  • 文/蒙蒙 一、第九天 我趴在偏房一处隐蔽的房顶上张望。 院中可真热闹,春花似锦、人声如沸。这庄子的主人今日做“春日...
    开封第一讲书人阅读 31,142评论 0 21
  • 文/苍兰香墨 我抬头看了看天上的太阳。三九已至,却和暖如春,着一层夹袄步出监牢的瞬间,已是汗流浃背。 一阵脚步声响...
    开封第一讲书人阅读 32,382评论 1 267
  • 我被黑心中介骗来泰国打工, 没想到刚下飞机就差点儿被人妖公主榨干…… 1. 我叫王不留,地道东北人。 一个月前我还...
    沈念sama阅读 47,020评论 2 365
  • 正文 我出身青楼,却偏偏与公主长得像,于是被迫代替她去往敌国和亲。 传闻我的和亲对象是个残疾皇子,可洞房花烛夜当晚...
    茶点故事阅读 44,044评论 2 352

推荐阅读更多精彩内容