RabbitMQ消息隊列中多路復用Channel信道詳解
什么叫消息隊列
消息(Message)是指在應用間傳送的數(shù)據。
消息可以非常簡單,比如只包含文本字符串,也可以更復雜,可能包含嵌入對象。
消息隊列(Message Queue)是一種應用間的通信方式,消息發(fā)送后可以立即返回,由消息系統(tǒng)來確保消息的可靠傳遞。
消息發(fā)布者只管把消息發(fā)布到 MQ 中而不用管誰來取,消息使用者只管從 MQ 中取消息而不管是誰發(fā)布的。
這樣發(fā)布者和使用者都不用知道對方的存在。
最簡單的理解是開了兩個定時器的線程,分別是上傳參數(shù)和下載數(shù)據兩個不同的線程,他們之間通過數(shù)據庫(也就是這里的queue)進行異步聯(lián)通,但這里并不是用的定時器而是用的觀察者模式
RabbitMQ 基本概念
上面只是最簡單抽象的描述,具體到 RabbitMQ 則有更詳細的概念需要解釋。
上面介紹過 RabbitMQ 是 AMQP 協(xié)議的一個開源實現(xiàn),所以其內部實際上也是 AMQP 中的基本概念:

RabbitMQ 內部結構
- Message
- 消息,消息是不具名的,它由消息頭和消息體組成。消息體是不透明的,而消息頭則由一系列的可選屬性組成,這些屬性包括routing-key(路由鍵)、priority(相對于其他消息的優(yōu)先權)、delivery-mode(指出該消息可能需要持久性存儲)等。
- Publisher
- 消息的生產者,也是一個向交換器發(fā)布消息的客戶端應用程序。
- Exchange
- 交換器,用來接收生產者發(fā)送的消息并將這些消息路由給服務器中的隊列。
- Binding
- 綁定,用于消息隊列和交換器之間的關聯(lián)。一個綁定就是基于路由鍵將交換器和消息隊列連接起來的路由規(guī)則,所以可以將交換器理解成一個由綁定構成的路由表。
- Queue
- 消息隊列,用來保存消息直到發(fā)送給消費者。它是消息的容器,也是消息的終點。一個消息可投入一個或多個隊列。消息一直在隊列里面,等待消費者連接到這個隊列將其取走。
- Connection
- 網絡連接,比如一個TCP連接。
- Channel
- 信道,多路復用連接中的一條獨立的雙向數(shù)據流通道。信道是建立在真實的TCP連接內地虛擬連接,AMQP 命令都是通過信道發(fā)出去的,不管是發(fā)布消息、訂閱隊列還是接收消息,這些動作都是通過信道完成。因為對于操作系統(tǒng)來說建立和銷毀 TCP 都是非常昂貴的開銷,所以引入了信道的概念,以復用一條 TCP 連接。
- Consumer
- 消息的消費者,表示一個從消息隊列中取得消息的客戶端應用程序。
- Virtual Host
- 虛擬主機,表示一批交換器、消息隊列和相關對象。虛擬主機是共享相同的身份認證和加密環(huán)境的獨立服務器域。每個 vhost 本質上就是一個 mini 版的 RabbitMQ 服務器,擁有自己的隊列、交換器、綁定和權限機制。vhost 是 AMQP 概念的基礎,必須在連接時指定,RabbitMQ 默認的 vhost 是 / 。
- Broker
- 表示消息隊列服務器實體。
對于理解多路復用,需要講下NIO:
其實,多路復用是一種思想,多路是指多個客戶端連接線路(TCP、Channel),復用是指使用一個線程重復使用,總的來說,就是單線程能同時處理多個請求。要實現(xiàn)這一點,就得改造BIO的連接模式了,BIO是客戶端直接連接服務端,NIO采用的是多路復用器 (Selector),相當于客戶端的連接不會直接連接服務端,而是連接到多路復用器。


這樣做的好處,就是把服務端和客戶端隔離了,如果不直接連接,服務器端就不會阻塞,多路復用器會將收到的消息做為事件請求發(fā)送給服務端,但服務端在處理事件的時候對于其他客戶端來說還是阻塞的,這些事件有不同類型。

Channel
了解了多路復用這個設計后,再講下Channel部分,Channel是一個雙向讀寫通道,是異步傳輸?shù)模跀?shù)據塊結構傳輸,BIO使用的是Stream流,基于字節(jié)傳輸。

關于性能方面,我也沒做過測試,但我知道一口咬定NIO比BIO性能要高效的言論,是錯誤的,NIO主要解決的不是性能的問題 Channel有很多實現(xiàn),因為本文介紹的是Socket,客戶端使用SocketChannel,服務端使用ServerSocketChannel
代碼實現(xiàn)
服務端:
import java.io.IOException;
import java.net.InetSocketAddress;
import java.nio.ByteBuffer;
import java.nio.channels.SelectionKey;
import java.nio.channels.Selector;
import java.nio.channels.ServerSocketChannel;
import java.nio.channels.SocketChannel;
import java.util.Iterator;
import java.util.Set;
public class NioServer {
public static void main(String[] args) throws IOException, InterruptedException {
// 打開Channel服務端并綁定監(jiān)聽一個端口
ServerSocketChannel ssc = ServerSocketChannel.open();
ssc.socket().bind(new InetSocketAddress(8459));
ssc.configureBlocking(false);
// 打開多路復用器 并注冊到 ServerSocketChannel 并監(jiān)聽連接事件
Selector selector = Selector.open();
ssc.register(selector, SelectionKey.OP_ACCEPT);
System.out.println("服務器已啟動...");
while (true) {
// 如果沒有事件發(fā)生 則select() 處于阻塞狀態(tài)
selector.select();
// 發(fā)生事件
Set<SelectionKey> selectionKeys = selector.selectedKeys();
Iterator<SelectionKey> iterator = selectionKeys.iterator();
// 處理事件
while (iterator.hasNext()) {
// 拿到具體事件
SelectionKey selectionKey = iterator.next();
// 判斷事件的類型
if (selectionKey.isAcceptable()) {
System.out.println("客戶端連接事件");
SocketChannel channel = ssc.accept();
channel.configureBlocking(false);
channel.register(selector, SelectionKey.OP_READ);
if (channel.finishConnect()) {
System.out.println("完成連接");
}
} else if (selectionKey.isReadable()) {
SocketChannel sc = (SocketChannel) selectionKey.channel();
ByteBuffer buffer = ByteBuffer.allocate(1024);
int read = sc.read(buffer);
System.out.println("收到的消息:" + new String(buffer.array(), 0, read));
// 響應客戶端 這里可有可無
buffer.clear();
buffer.put("已收到消息".getBytes());
// 將緩沖區(qū)各標志復位,因為向里面put了數(shù)據標志被改變要想從中讀取數(shù)據發(fā)向服務器,就要復位
buffer.flip();
sc.write(buffer);
// 設置監(jiān)聽事件的集合 這里把寫入事件加入
selectionKey.interestOps(selectionKey.interestOps() | SelectionKey.OP_WRITE);
System.out.println("服務器向客戶端發(fā)送已確認消息");
} else if (selectionKey.isWritable()) {
System.out.println("觸發(fā)往客戶端寫入數(shù)據事件");
// 發(fā)送完了就取消監(jiān)聽寫事件,否則會無限循環(huán)觸發(fā)寫事件
selectionKey.interestOps(selectionKey.interestOps() & ~SelectionKey.OP_WRITE);
}
// 事件完成后,將其移除
iterator.remove();
}
}
}客戶端
import java.io.IOException;
import java.net.InetSocketAddress;
import java.nio.ByteBuffer;
import java.nio.channels.SelectionKey;
import java.nio.channels.Selector;
import java.nio.channels.SocketChannel;
import java.util.Iterator;
import java.util.Scanner;
import java.util.Set;
public class NioClient {
public static void main(String[] args) throws IOException {
//打開選擇器
Selector selector = Selector.open();
//打開通道
SocketChannel socketChannel = SocketChannel.open();
//配置非阻塞模型
socketChannel.configureBlocking(false);
//連接遠程主機
socketChannel.connect(new InetSocketAddress("127.0.0.1", 8459));
//注冊事件
socketChannel.register(selector, SelectionKey.OP_CONNECT | SelectionKey.OP_READ );
//循環(huán)處理
new Thread(() -> {
while (true) {
try {
selector.select();
Set<SelectionKey> keys = selector.selectedKeys();
Iterator<SelectionKey> iter = keys.iterator();
while (iter.hasNext()) {
SelectionKey key = iter.next();
if (key.isConnectable()) {
//連接建立或者連接建立不成功
SocketChannel channel = (SocketChannel) key.channel();
//完成連接的建立
if (channel.finishConnect()) {
System.out.println("完成連接");
}
} else if (key.isReadable()) {
SocketChannel channel = (SocketChannel) key.channel();
ByteBuffer buffer = ByteBuffer.allocate(1024);
buffer.clear();
channel.read(buffer);
System.out.println("客戶端收到消息:" + new String(buffer.array()));
key.interestOps(key.interestOps() | SelectionKey.OP_WRITE);
} else if (key.isWritable()) {
System.out.println("客戶端向服務端寫入數(shù)據");
// 設置監(jiān)聽事件的集合 這里把寫入事件加入
key.interestOps(key.interestOps() & ~SelectionKey.OP_WRITE);
}
iter.remove();
}
} catch (IOException e) {
e.printStackTrace();
break;
}
}
}).start();
Scanner scanner = new Scanner(System.in);
while (true) {
System.out.println("請輸入要發(fā)送的字符串");
String str = scanner.nextLine();
ByteBuffer buffer = ByteBuffer.allocate(1024);
buffer.put(str.getBytes());
buffer.flip();
socketChannel.write(buffer);
}
}
}講完多路復用,接下來就是TCP實現(xiàn)的多路復用,整體和NIO很相似,都是把請求注冊到selector中,根據SelectionKey進行阻塞式的通信,因為有了注冊和SelectionKey使得單線程能同時處理多個請求成為可能
IO多路復用(multiplexing)屬于同步IO網絡模型
是以Reactor模式實現(xiàn)
常見的IO多路復用應用有:select、poll、epoll
本篇文章采用Java的NIO框架來實現(xiàn)單線程的IO多路復用
Reactor模式的組成角色
1. Reactor:負責派發(fā)IO事件給對應的角色處理。為了監(jiān)聽IO事件,select必須實現(xiàn)在Reactor中。
2. Acceptor:負責接受client的連線,然后給client綁定一個Handler并注冊IO事件到Reactor上監(jiān)聽。
3. Handler:負責處理與client交互的事件或行為。通常因為Handler要處理與所對應client交互的多個事件或行為,為了簡化設計,會以狀態(tài)模式來實現(xiàn)Handler。

代碼實現(xiàn)
[TCPReactor.java]
// Reactor線程
package server;
import java.io.IOException;
import java.net.InetSocketAddress;
import java.nio.channels.SelectionKey;
import java.nio.channels.Selector;
import java.nio.channels.ServerSocketChannel;
import java.util.Iterator;
import java.util.Set;
public class TCPReactor implements Runnable {
private final ServerSocketChannel ssc;
private final Selector selector;
public TCPReactor(int port) throws IOException {
selector = Selector.open();
ssc = ServerSocketChannel.open();
InetSocketAddress addr = new InetSocketAddress(port);
ssc.socket().bind(addr); // 在ServerSocketChannel綁定監(jiān)聽端口
ssc.configureBlocking(false); // 設置ServerSocketChannel為非阻塞
SelectionKey sk = ssc.register(selector, SelectionKey.OP_ACCEPT); // ServerSocketChannel向selector註冊一個OP_ACCEPT事件,然後返回該通道的key
sk.attach(new Acceptor(selector, ssc)); // 給定key一個附加的Acceptor對象
}
@Override
public void run() {
while (!Thread.interrupted()) { // 在線程被中斷前持續(xù)運行
System.out.println("Waiting for new event on port: " + ssc.socket().getLocalPort() + "...");
try {
if (selector.select() == 0) // 若沒有事件就緒則不往下執(zhí)行
continue;
} catch (IOException e) {
// TODO Auto-generated catch block
e.printStackTrace();
}
Set<SelectionKey> selectedKeys = selector.selectedKeys(); // 取得所有已就緒事件的key集合
Iterator<SelectionKey> it = selectedKeys.iterator();
while (it.hasNext()) {
dispatch((SelectionKey) (it.next())); // 根據事件的key進行調度
it.remove();
}
}
}
/*
* name: dispatch(SelectionKey key)
* description: 調度方法,根據事件綁定的對象開新線程
*/
private void dispatch(SelectionKey key) {
Runnable r = (Runnable) (key.attachment()); // 根據事件之key綁定的對象開新線程
if (r != null)
r.run();
}
} [Acceptor.java]
// 接受連線請求線程
package server;
import java.io.IOException;
import java.nio.channels.SelectionKey;
import java.nio.channels.Selector;
import java.nio.channels.ServerSocketChannel;
import java.nio.channels.SocketChannel;
public class Acceptor implements Runnable {
private final ServerSocketChannel ssc;
private final Selector selector;
public Acceptor(Selector selector, ServerSocketChannel ssc) {
this.ssc=ssc;
this.selector=selector;
}
@Override
public void run() {
try {
SocketChannel sc= ssc.accept(); // 接受client連線請求
System.out.println(sc.socket().getRemoteSocketAddress().toString() + " is connected.");
if(sc!=null) {
sc.configureBlocking(false); // 設置為非阻塞
SelectionKey sk = sc.register(selector, SelectionKey.OP_READ); // SocketChannel向selector註冊一個OP_READ事件,然後返回該通道的key
selector.wakeup(); // 使一個阻塞住的selector操作立即返回
sk.attach(new TCPHandler(sk, sc)); // 給定key一個附加的TCPHandler對象
}
} catch (IOException e) {
// TODO Auto-generated catch block
e.printStackTrace();
}
}
} 我們先來簡單點的,Handler不以??狀態(tài)模式實現(xiàn),只以比較直覺的方式實現(xiàn)。
[TCPHandler.java]
// Handler線程
package server;
import java.io.IOException;
import java.nio.ByteBuffer;
import java.nio.channels.SelectionKey;
import java.nio.channels.SocketChannel;
import java.util.concurrent.LinkedBlockingQueue;
import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;
public class TCPHandler implements Runnable {
private final SelectionKey sk;
private final SocketChannel sc;
int state;
public TCPHandler(SelectionKey sk, SocketChannel sc) {
this.sk = sk;
this.sc = sc;
state = 0; // 初始狀態(tài)設定為READING
}
@Override
public void run() {
try {
if (state == 0)
read(); // 讀取網絡數(shù)據
else
send(); // 發(fā)送網絡數(shù)據
} catch (IOException e) {
System.out.println("[Warning!] A client has been closed.");
closeChannel();
}
}
private void closeChannel() {
try {
sk.cancel();
sc.close();
} catch (IOException e1) {
e1.printStackTrace();
}
}
private synchronized void read() throws IOException {
// non-blocking下不可用Readers,因為Readers不支援non-blocking
byte[] arr = new byte[1024];
ByteBuffer buf = ByteBuffer.wrap(arr);
int numBytes = sc.read(buf); // 讀取字符串
if(numBytes == -1)
{
System.out.println("[Warning!] A client has been closed.");
closeChannel();
return;
}
String str = new String(arr); // 將讀取到的byte內容轉為字符串型態(tài)
if ((str != null) && !str.equals(" ")) {
process(str); // 邏輯處理
System.out.println(sc.socket().getRemoteSocketAddress().toString()
+ " > " + str);
state = 1; // 改變狀態(tài)
sk.interestOps(SelectionKey.OP_WRITE); // 通過key改變通道註冊的事件
sk.selector().wakeup(); // 使一個阻塞住的selector操作立即返回
}
}
private void send() throws IOException {
// get message from message queue
String str = "Your message has sent to "
+ sc.socket().getLocalSocketAddress().toString() + "\r\n";
ByteBuffer buf = ByteBuffer.wrap(str.getBytes()); // wrap自動把buf的position設為0,所以不需要再flip()
while (buf.hasRemaining()) {
sc.write(buf); // 回傳給client回應字符串,發(fā)送buf的position位置 到limit位置為止之間的內容
}
state = 0; // 改變狀態(tài)
sk.interestOps(SelectionKey.OP_READ); // 通過key改變通道註冊的事件
sk.selector().wakeup(); // 使一個阻塞住的selector操作立即返回
}
void process(String str) {
// do process(decode, logically process, encode)..
// ..
}
} 最后是主程序代碼
[Main.java]
package server;
import java.io.IOException;
public class Main {
public static void main(String[] args) {
// TODO Auto-generated method stub
try {
TCPReactor reactor = new TCPReactor(1333);
reactor.run();
} catch (IOException e) {
// TODO Auto-generated catch block
e.printStackTrace();
}
}
} 下面附上客戶端代碼:
[Client.java]
package main.pkg;
import java.io.BufferedReader;
import java.io.IOException;
import java.io.InputStreamReader;
import java.io.PrintWriter;
import java.net.Socket;
import java.net.UnknownHostException;
public class Client {
/**
* @param args
*/
public static void main(String[] args) {
// TODO Auto-generated method stub
String hostname=args[0];
int port = Integer.parseInt(args[1]);
//String hostname="127.0.0.1";
//int port=1333;
System.out.println("Connecting to "+ hostname +":"+port);
try {
Socket client = new Socket(hostname, port); // 連接至目的地
System.out.println("Connected to "+ hostname);
PrintWriter out = new PrintWriter(client.getOutputStream());
BufferedReader in = new BufferedReader(new InputStreamReader(client.getInputStream()));
BufferedReader stdIn = new BufferedReader(new InputStreamReader(System.in));
String input;
while((input=stdIn.readLine()) != null) { // 讀取輸入
out.println(input); // 發(fā)送輸入的字符串
out.flush(); // 強制將緩衝區(qū)內的數(shù)據輸出
if(input.equals("exit"))
{
break;
}
System.out.println("server: "+in.readLine());
}
client.close();
System.out.println("client stop.");
} catch (UnknownHostException e) {
// TODO Auto-generated catch block
System.err.println("Don't know about host: " + hostname);
} catch (IOException e) {
// TODO Auto-generated catch block
System.err.println("Couldn't get I/O for the socket connection");
}
}
}到此這篇關于RabbitMQ消息隊列中多路復用Channel信道詳解的文章就介紹到這了,更多相關RabbitMQ多路復用Channel內容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關文章希望大家以后多多支持腳本之家!
相關文章
Java實現(xiàn)日志文件監(jiān)聽并讀取相關數(shù)據的方法實踐
本文主要介紹了Java實現(xiàn)日志文件監(jiān)聽并讀取相關數(shù)據的方法實踐,文中通過示例代碼介紹的非常詳細,具有一定的參考價值,感興趣的小伙伴們可以參考一下2022-05-05
Java?json轉換實體類(JavaBean)及實體類(JavaBean)轉換json代碼示例
這篇文章主要介紹了兩種常見的JSON與Java實體類相互轉換的方法,分別是使用庫Jackson、Gson、Fastjson和在線工具,無論是將JSON轉換為Java實體類還是將Java實體類轉換為JSON,這些方法都能顯著簡化開發(fā)過程,需要的朋友可以參考下2024-12-12
Servlet實現(xiàn)統(tǒng)計頁面訪問次數(shù)功能
這篇文章主要介紹了Servlet實現(xiàn)統(tǒng)計頁面訪問次數(shù)功能,文中示例代碼介紹的非常詳細,具有一定的參考價值,感興趣的小伙伴們可以參考一下2021-04-04
java 中使用maven shade plugin 打可執(zhí)行Jar包
這篇文章主要介紹了java 中使用maven shade plugin 打可執(zhí)行Jar包的相關資料,需要的朋友可以參考下2017-05-05
Mybatis-Plus中getOne方法獲取最新一條數(shù)據的示例代碼
這篇文章主要介紹了Mybatis-Plus中getOne方法獲取最新一條數(shù)據,本文通過示例代碼給大家介紹的非常詳細,對大家的學習或工作具有一定的參考借鑒價值,需要的朋友可以參考下2023-05-05

