前言
netty框架馬上就進入尾聲了,小編沒有特別深入的講解,第一是網路編程確實挺難的,第二用好netty其實是挺不容易的一件事情,尤其都是異步的情況下,今天小編繼續為大家帶來開發實戰,上次分享了redis客戶端和websocket彈幕功能的簡單實作,這次為大家帶來相對比較高檔的rpc框架底層網路通信,今天主要以dubbo為例,希望大家有所識訓,
RPC
定義
RPC為遠程服務呼叫,即客戶端遠程呼叫服務端的方法,然后服務端回傳回應或例外,常用的RPC解決方案有JAVA RMI,webService,Http Invoker,Dubbo,SpringCloud等等,
去中心化架構
傳統集中式架構:
下圖是小編理解的中心化架構

中心化架構其優點為:架構簡單,客戶端呼叫的時候可以跨語言,缺點的話所有的呼叫都會進過ngnix(這里就可以理解為中心,大家都得進過他,無論是呼叫還是回傳回應),ngnix一旦掛了之后就會是服務掛了,當然如果ngnix部署集群也會讓架構變的復雜,
去中心化架構
如下圖:

這里的話客戶端呼叫服務端不需要進過中心直連呼叫,
去中心化架構簡單描述后,繼續來看一下rpc框架組成,
框架組成

這個架構是不是似曾相識,這里如果看過小編的dubbo分享就覺得面熟了,
上面最小化實作rpc框架就是最下面的rpc協議就可以了,這樣兩個服務就可以通信即可,接著咱們來介紹一下rpc協議,
協議報文
這里小編直接用dubbo協議來說明報文:

上面請求頭主要有16個位元組,如果看過小編的前基本應用的使用后,看到這個就可以使用netty來進行訊息的拆包以及編解碼,那小編接下來繼續說明編解碼的程序,
概設程序
這邊在寫代碼之前,小編先分享一下設計的思路,以及一些呼叫的邏輯,
首先是編解碼:編解碼占網路傳輸中必須且固定的,先看下圖:

具體已經在上圖解釋清楚了,接下來看器呼叫程序即各個組件功能,

上圖是比較簡單的,接下來是至關重要的的從客戶端到服務整個流程的呼叫程序

容小編解釋一下:
- 從客戶端到服務端的呼叫涉及到了四個執行緒,分別是客戶端以及服務端的業務執行緒和IO執行緒
- 發起掉用寫入訊息體的內容是上面的Transfer,而編碼request則為客戶端發起的請求,而Transfer中的Request包含了介面,方法以及引數
- 寫入到socket中的為bytebuf,其經過Bytebuf -> head -> unsafe -> nio socket (doWrite) -> java nio channel -> socket
- 寫入到socket則到達服務端,服務端的io執行緒通過多路復用選擇器select輪詢,之后呼叫read
- 讀取到內容后,這邊的讀取流程和上面相似 unsafe read -> pipeline fireChannelRead,拿到了ByteBuf,然后解碼request,根據上面的解碼工具類
- 從ByteBuf拿到Transfer,之后交給業務的handler,之后涉及到服務端業務執行緒的處理,業務處理后回傳了結果或者報錯資訊(這里同上其實是Transfer)
- 之后又交還給io執行緒,再次將Response進行編碼(服務端的回應),Bytebuf寫到socket
- 回到客戶端io執行緒后,再次有select進行輪詢,讀取到內容Bytebuf,解碼成Transfer,Transfer中的response進行反序列化拿到結果填充回執,
- 客戶端拿到回執,釋放等待,
注意事項
第一:加入客戶端A和B的請求,客戶端怎樣拿到服務端回來的A回應和B相應呢,這里就需要Transfer里面的id,即協議中的id,不過如果是協議中的id,那就需要做請求的時候放入到一個map中來保存,
第二:既然知道使用id來區分請求回應,那什么時候放入到map中, 怎么保證執行緒安全,那最好是執行緒安全的map,不過高并發的時候,對系統很不友好,所以放入到map的時候也在io執行緒中執行,
第三:如何在io執行緒放入map中,這里是用eventloop的submit,寫入訊息完成后監聽并寫入map
講完理論小編不是純粹的理論派,還是代碼實戰派
代碼實戰
編解碼工具以及傳輸類
Transfer類
public class Transfer {
public static final byte STATUS_ERROR = 0;
public static final byte STATUS_OK = 1;
public static final byte STATUS_ILLEGAL = 2;
public static final byte SERIALIZABLE_JAVA=1;
public static final byte SERIALIZABLE_HESSIAN2=2;
public static final byte SERIALIZABLE_JSON=3;
boolean request;
byte serializableId; // 1:java 2:hessian2 3:json
boolean twoWay;
boolean heartbeat;
long id;
byte status; // 1正常 0失敗 2請求非法
Object target;
public Transfer(long id) {
this.id = id;
}
}
編解碼工具類
public class RpcCodec extends ByteToMessageCodec {
private static final int HEADER_LENGTH = 16;
private static final short MAGIC = 0xdad;
private static final ByteBuf MAGIC_BUF = Unpooled.copyShort(MAGIC);
private static final byte FLAG_REQUEST = (byte) 0x80;//1000 0000
private static final byte FLAG_TWO_WAY = (byte) 0x40; //0100 0000
private static final byte FLAG_EVENT = (byte) 0x20; //0010 0000
private static final int SERIALIZATION_MASK = 0x1f; //0001 1111
// 編碼
@Override
protected void encode(ChannelHandlerContext ctx, Object msg, ByteBuf out) {
if (msg instanceof Transfer) {
doEncode((Transfer) msg, out);
} else {
throw new IllegalArgumentException();
}
}
//解碼
@Override
protected void decode(ChannelHandlerContext ctx, ByteBuf in, List out) {
Transfer transfer = doDecode(in);
if (transfer != null) {
out.add(transfer);
}
}
// 編碼
protected void doEncode(Transfer data, ByteBuf buf) {
byte[] header = new byte[HEADER_LENGTH];
Bytes.short2bytes(MAGIC, header);
header[2] = data.serializableId;
if (data.request) header[2] |= FLAG_REQUEST;
if (data.twoWay) header[2] |= FLAG_TWO_WAY;
if (data.heartbeat) header[2] |= FLAG_EVENT;
if (!data.request) header[3] = data.status;
Bytes.long2bytes(data.id, header, 4);// id 占8個位元組
int len = 0;
byte[] body = new byte[0];
if (!data.heartbeat) {
body = serialize(data.serializableId, data.target);
len = body.length;
}
Bytes.int2bytes(len, header, 12);
buf.writeBytes(header);
buf.writeBytes(body);
}
// 解碼
protected Transfer doDecode(ByteBuf in) {
int index = ByteBufUtil.indexOf(MAGIC_BUF, in);
//是否有魔數
if (index < 0) {
return null;
}
//訊息頭是否完整
if (!in.isReadable(index + HEADER_LENGTH)) {
return null;
}
byte[] header = new byte[HEADER_LENGTH];
// in.getBytes(index, header);
ByteBuf slice = in.slice();
slice.readBytes(header);
int length = Bytes.bytes2int(header, 12);
//訊息體是否完整
if (!in.isReadable(index + HEADER_LENGTH + length)) {
return null;//需要更多的位元組
}
Transfer transfer = new Transfer(Bytes.bytes2long(header, 4));
transfer.heartbeat = (header[2] & FLAG_EVENT) != 0;
transfer.request = (header[2] & FLAG_REQUEST) != 0;
transfer.twoWay = (header[2] & FLAG_TWO_WAY) != 0;
transfer.serializableId = (byte) (header[2] & SERIALIZATION_MASK);
transfer.status = header[3];
if (!transfer.heartbeat) {
byte[] content = new byte[length];
// in.getBytes(index + HEADER_LENGTH, bytes);
slice.readBytes(content);
transfer.target = deserialize(transfer.serializableId, content);
}
//跳過已經讀取的
in.skipBytes(index + HEADER_LENGTH + length);
return transfer;
}
// 序列化
private byte[] serialize(byte serializableId, Object target) {
if (serializableId == Transfer.SERIALIZABLE_JAVA) { //JAVA
ByteArrayOutputStream out;
try {
out = new ByteArrayOutputStream();
ObjectOutputStream stream = new ObjectOutputStream(out);
stream.writeObject(target);
} catch (IOException e) {
throw new RuntimeException(e);
}
return out.toByteArray();
} else {
throw new UnsupportedOperationException();
}
}
// 反序列化
private Object deserialize(byte serializableId, byte[] bytes) {
if (serializableId == Transfer.SERIALIZABLE_JAVA) { //JAVA
try {
ObjectInputStream stream =
new ObjectInputStream(new ByteArrayInputStream(bytes));
return stream.readObject();
} catch (IOException | ClassNotFoundException e) {
throw new RuntimeException(e);
}
} else {
throw new UnsupportedOperationException();
}
}
}
客戶端代碼
public class RpcClient {
static AtomicLong atomicLong = new AtomicLong(100);
private Channel channel;
private Map<Long, Promise<Response>> results = new HashMap<>();
public static long getNextId() {
return atomicLong.getAndIncrement();
}
public void init(String address, int port) throws InterruptedException {
Bootstrap bootstrap = new Bootstrap();
bootstrap.group(new NioEventLoopGroup(1))
.channel(NioSocketChannel.class);
bootstrap.handler(new ChannelInitializer<Channel>() {
@Override
protected void initChannel(Channel ch) {
ch.pipeline().addLast("codec", new RpcCodec());
ch.pipeline().addLast("resultSet", new ResultFill());// 結果集填充
}
});
ChannelFuture connect = bootstrap.connect(address, port);
channel = connect.sync().channel();
System.out.println("連接成功");
//
// 每隔 兩秒發送心跳
channel.eventLoop().scheduleWithFixedDelay(() -> {
Transfer transfer=new Transfer(getNextId());
transfer.heartbeat=true;
channel.writeAndFlush(transfer);
},2000,2000,TimeUnit.MILLISECONDS);
}
public Response invokerRemote(Class serverInterface,
String methodDesc,
Object[] args) throws InterruptedException, ExecutionException, TimeoutException {
Request request = new Request(serverInterface.getName(), methodDesc);
request.setArgs(args);
Transfer transfer = new Transfer(getNextId());
transfer.request=true;
transfer.serializableId=Transfer.SERIALIZABLE_JAVA;
transfer.target = request;
DefaultPromise<Response> resultPromise = new DefaultPromise(channel.eventLoop());
// 寫入成功后添加 結果
channel.writeAndFlush(transfer).addListener(future ->
{// IO執行緒
if (future.cause() != null) {// 寫入失敗
resultPromise.setFailure(future.cause()); //寫入失敗必須處理
} else { // 寫入成功
results.put(transfer.id, resultPromise);
}
}
);
return resultPromise.get(10000, TimeUnit.MILLISECONDS);
}
private class ResultFill extends SimpleChannelInboundHandler<Transfer> {
@Override
protected void channelRead0(ChannelHandlerContext ctx, Transfer msg) {
if (msg.heartbeat) {
System.out.println(String.format("服務端心跳回傳:%s",
ctx.channel().remoteAddress()));
} else {
Promise<Response> promise = results.remove(msg.id);
promise.setSuccess((Response) msg.target); // 填充結果
}
}
}
public <T> T getRemoteService(Class<T> serviceInterface) {
assert serviceInterface.isInterface();
Object o = Proxy.newProxyInstance(getClass().getClassLoader(), new Class[]{serviceInterface}, new InvocationHandler() {
@Override
public Object invoke(Object proxy, Method method, Object[] args) throws Exception {
if (Object.class.equals(method.getDeclaringClass())) {
return method.invoke(this, args);
}
String methodDescriptor = method.getName()+Type.getMethodDescriptor(method);
Response response = invokerRemote(serviceInterface, methodDescriptor, args);
if (response.getError() != null) {
throw new RuntimeException("遠程服務呼叫例外:", response.getError());
}
return response.getResult();
}
});
return (T) o;
}
}
服務端代碼:
public class RpcServer {
ExecutorService threadPool = Executors.newFixedThreadPool(500);
private Map<String, ServiceBean> register = new HashMap<>();
public void start(int port) throws InterruptedException {
ServerBootstrap bootstrap = new ServerBootstrap();
EventLoopGroup boss = new NioEventLoopGroup(1);
EventLoopGroup work = new NioEventLoopGroup(8);
bootstrap.group(boss, work).channel(NioServerSocketChannel.class).childHandler(new ChannelInitializer<Channel>() {
@Override
protected void initChannel(Channel ch) throws Exception {
ch.pipeline().addLast("codec", new RpcCodec());
ch.pipeline().addLast("dispatch", new Dispatch());
}
}).bind(port).sync();
System.out.println("服務啟動成功");
}
private class Dispatch extends SimpleChannelInboundHandler<Transfer> {
@Override
protected void channelRead0(ChannelHandlerContext ctx, Transfer transfer) {
if (transfer.heartbeat) { // 心跳處理
Transfer t = new Transfer(transfer.id);
t.heartbeat = true;
t.request = false;
ctx.writeAndFlush(t);// 回傳心跳
} else {
threadPool.submit(() -> {
Transfer to = doDispatchRequest(transfer);
ctx.writeAndFlush(to);// 非IO執行緒 異步提交到IO
});
}
}
// 業務請求處理
Transfer doDispatchRequest(Transfer from) {
Request request = (Request) from.target;
Transfer to = new Transfer(from.id);
to.request = false;
to.serializableId = from.serializableId;
Response response = new Response();
try {
String serverId = request.getClassName() + request.getMethodDesc();
ServiceBean serverBean = register.get(serverId);
if (serverBean == null) {
throw new IllegalArgumentException("找不到服務" + serverId);
}
Object result = serverBean.invoke(request.getArgs());
response.setResult(result);
to.status = Transfer.STATUS_OK;
} catch (Throwable e) {
e.printStackTrace();
response.setError(e);
to.status = Transfer.STATUS_ERROR;
}
to.target = response;
return to;
}
}
private static class ServiceBean {
Method method;
Object target;
public ServiceBean(Method method, Object target) {
this.method = method;
this.target = target;
}
public Object invoke(Object[] args) throws Exception {
return method.invoke(target, args);
}
}
public void registerServer(Class serviceInterface, Object serverBean) {
assert serviceInterface.isInterface();
for (Method method : serviceInterface.getMethods()) {
int modifiers = method.getModifiers();
if (Modifier.isStatic(modifiers) || Modifier.isNative(modifiers)) {
continue;
}
String methodDescriptor = Type.getMethodDescriptor(method);
String key = serviceInterface.getName() +method.getName()+ methodDescriptor;
register.put(key, new ServiceBean(method, serverBean));
}
}
}
測驗類
public class RpcTest {
@Test
public void startServerTest() throws InterruptedException, IOException {
RpcServer server = new RpcServer();
server.registerServer(UserService.class, new UserServiceImpl());
server.start(8080);
System.in.read();
}
public static void main(String[] args) throws InterruptedException, IOException {
RpcClient client = new RpcClient();
client.init("127.0.0.1", 8080);
UserService service = client.getRemoteService(UserService.class);
BufferedReader in = new BufferedReader(new InputStreamReader(System.in));
while (true) {
String s = in.readLine();
System.out.println(service.getUser(1));
}
}
// 多執行緒并發呼叫
@Test
public void syncTest() throws InterruptedException, IOException {
RpcClient client = new RpcClient();
client.init("127.0.0.1", 8080);
UserService service = client.getRemoteService(UserService.class);
ExecutorService executor = Executors.newFixedThreadPool(100);
for (int i = 0; i < 100; i++) {
int id = i;
executor.execute(() -> {
User user = service.getUser(id);
System.out.println(user);
assert user.getId().equals(id);
});
}
System.in.read();
}
}
對于UserSevice大家可以自己建一個介面和隨便實作一下即可,
總結
利用netty實作rpc還是挺有難度的,小編是進過將近一周才陸續寫完這篇博客,希望小編已經充分講清楚了,假如大家將netty的理論與實戰結合完畢,那相信和小編一樣有長足進步,加油!
轉載請註明出處,本文鏈接:https://www.uj5u.com/qukuanlian/295725.html
標籤:區塊鏈
