[pravega-022] pravega源码分析--Controller子项目--关于netty[02]

2024-06-11 09:08

本文主要是介绍[pravega-022] pravega源码分析--Controller子项目--关于netty[02],希望对大家解决编程问题提供一定的参考价值,需要的开发者们随着小编来一起学习吧!

1.一个netty echo服务的例子。client向server发消息,server向client返回同样的消息。程序分两部分,server项目和client项目。这里包含了netty的核心要素。更多的细节请参考前文提到的参考资料即可。

2.server项目

2.1 目录结构

├── build.gradle
├── settings.gradle
└── src
    ├── main
        ├── java
           ├── EchoServerHandler.java
           └── Main.java
2.2 build.gradle文件内容

group 'com.brian.demo.netty'
version '1.0-SNAPSHOT'apply plugin: 'java'sourceCompatibility = 1.8repositories {mavenCentral()
}dependencies {compile group: 'io.netty', name: 'netty-all', version: '4.1.12.Final'testCompile group: 'junit', name: 'junit', version: '4.12'
}

2.3 Main.java文件内容

import io.netty.bootstrap.ServerBootstrap;
import io.netty.channel.ChannelFuture;
import io.netty.channel.ChannelInitializer;
import io.netty.channel.EventLoopGroup;
import io.netty.channel.nio.NioEventLoopGroup;
import io.netty.channel.socket.SocketChannel;
import io.netty.channel.socket.nio.NioServerSocketChannel;import java.net.InetSocketAddress;//主要参考资料:
// 《netty in action》
// https://github.com/waylau/essential-netty-in-action
public class Main {private final int port;public Main(int port) {this.port = port;}public static void main(String args[]) throws Exception {new Main(8008).start();}public void start() throws Exception {//NioEventLoopGroup是一个线程池,一个线程可以处理多个channel,一个channel只对应一个线程EventLoopGroup group = new NioEventLoopGroup();//echo服务handlerfinal EchoServerHandler echoServerHandler = new EchoServerHandler();//创建 引导服务器ServerBootstrap sbs = new ServerBootstrap();try {//设置线程池sbs.group(group)//指定NIO传输的channel.channel(NioServerSocketChannel.class)//本地端口.localAddress(new InetSocketAddress(port))//配置handler pipeline。每个channel有一个pipline,// 包含多个handler,事件从pipline逐个经过handler进行处理。.childHandler(new ChannelInitializer<SocketChannel>() {@Overridepublic void initChannel(SocketChannel ch) throws Exception {ch.pipeline().addLast(echoServerHandler);}});//绑定服务器,然后同步等待服务器关闭。ChannelFuture future = sbs.bind().sync();System.out.println(Main.class.getName() +" started and listening for connections on " + future.channel().localAddress());//关闭channel,同步等待future.channel().closeFuture().sync();} finally {//释放线程池group.shutdownGracefully().sync();}}
}

2.4 EchoServerHandler.java文件内容

import io.netty.buffer.ByteBuf;
import io.netty.buffer.Unpooled;
import io.netty.channel.ChannelFutureListener;
import io.netty.channel.ChannelHandler;
import io.netty.channel.ChannelHandlerContext;
import io.netty.channel.ChannelInboundHandlerAdapter;
import io.netty.util.CharsetUtil;//Sharable注解,声明这个Handler的实例可以被多个channel共享使用
@ChannelHandler.Sharable
//ChannelInboundHandler接口,处理入站事件
public class EchoServerHandler extends ChannelInboundHandlerAdapter {//每个信息入站都会调用channelRead@Overridepublic void channelRead(ChannelHandlerContext ctx, Object msg){//收到的消息是Object msg,强转撑ByteBufByteBuf inBuf = (ByteBuf)msg;//在控制台输出接收到的消息System.out.println("Server received:" + inBuf.toString(CharsetUtil.UTF_8));//再把消息不做修改重新写给客户端。注意,此时数据还没有flush,仍然在服务端。ctx.write(inBuf);}//通知处理器,当下的channelread()读取消息是本批处理的最后一条消息的时候,调用本函数@Overridepublic void channelReadComplete(ChannelHandlerContext ctx) throws Exception {//把所有数据冲刷flush到客户端,关闭通道。至此操作完成。Listern以future方式通知操作完成。ctx.writeAndFlush(Unpooled.EMPTY_BUFFER).addListener(ChannelFutureListener.CLOSE);}//读操作遇到异常@Overridepublic void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) throws Exception {cause.printStackTrace();ctx.close();}
}

3. Client项目

3.1 文件目录结构

.
├── build.gradle
├── settings.gradle
└── src
    ├── main
        ├── java
            ├── EchoClientHandler.java
            └── Main.java
    
3.2 build.gradle文件内容

group 'com.brian.demo.netty'
version '1.0-SNAPSHOT'apply plugin: 'java'sourceCompatibility = 1.8repositories {mavenCentral()
}dependencies {compile group: 'io.netty', name: 'netty-all', version: '4.1.12.Final'testCompile group: 'junit', name: 'junit', version: '4.12'
}

3.3 Main.java文件内容

import io.netty.bootstrap.Bootstrap;
import io.netty.channel.ChannelFuture;
import io.netty.channel.ChannelInitializer;
import io.netty.channel.EventLoopGroup;
import io.netty.channel.nio.NioEventLoopGroup;
import io.netty.channel.socket.SocketChannel;
import io.netty.channel.socket.nio.NioSocketChannel;import java.net.InetSocketAddress;public class Main {private final String host;private final int port;public Main(){this.host="127.0.0.1";this.port=8008;}public void start() throws Exception{EventLoopGroup group = new NioEventLoopGroup();try {Bootstrap b = new Bootstrap();b.group(group).channel(NioSocketChannel.class).remoteAddress(new InetSocketAddress(host,port)).handler(new ChannelInitializer<SocketChannel>() {@Overrideprotected void initChannel(SocketChannel ch) throws Exception {ch.pipeline().addLast(new EchoClientHandler());}});ChannelFuture future = b.connect().sync();future.channel().closeFuture().sync();}finally {group.shutdownGracefully();}}public static void main(String args[])throws Exception{new Main().start();}
}

       

3.4 EchoClientHandler.java文件内容

import io.netty.buffer.ByteBuf;
import io.netty.buffer.Unpooled;
import io.netty.channel.ChannelHandler;
import io.netty.channel.ChannelHandlerContext;
import io.netty.channel.SimpleChannelInboundHandler;
import io.netty.util.CharsetUtil;@ChannelHandler.Sharable
public class EchoClientHandler extends SimpleChannelInboundHandler<ByteBuf> {@Overrideprotected void channelRead0(ChannelHandlerContext ctx, ByteBuf msg) throws Exception {System.out.println("client received:" + msg.toString(CharsetUtil.UTF_8));}@Overridepublic void channelActive(ChannelHandlerContext ctx) throws Exception {ctx.writeAndFlush(Unpooled.copiedBuffer("hello,world", CharsetUtil.UTF_8));}@Overridepublic void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) throws Exception {cause.printStackTrace();ctx.close();}
}

4. 如果遇到gradle不能导入nettty包的情况,关闭项目,删除~/.gradle目录所有文件,然后重新在idea做import项目即可。
   
   

这篇关于[pravega-022] pravega源码分析--Controller子项目--关于netty[02]的文章就介绍到这儿,希望我们推荐的文章对编程师们有所帮助!



http://www.chinasem.cn/article/1050735

相关文章

MyBatis Plus 中 update_time 字段自动填充失效的原因分析及解决方案(最新整理)

《MyBatisPlus中update_time字段自动填充失效的原因分析及解决方案(最新整理)》在使用MyBatisPlus时,通常我们会在数据库表中设置create_time和update... 目录前言一、问题现象二、原因分析三、总结:常见原因与解决方法对照表四、推荐写法前言在使用 MyBATis

Python主动抛出异常的各种用法和场景分析

《Python主动抛出异常的各种用法和场景分析》在Python中,我们不仅可以捕获和处理异常,还可以主动抛出异常,也就是以类的方式自定义错误的类型和提示信息,这在编程中非常有用,下面我将详细解释主动抛... 目录一、为什么要主动抛出异常?二、基本语法:raise关键字基本示例三、raise的多种用法1. 抛

github打不开的问题分析及解决

《github打不开的问题分析及解决》:本文主要介绍github打不开的问题分析及解决,具有很好的参考价值,希望对大家有所帮助,如有错误或未考虑完全的地方,望不吝赐教... 目录一、找到github.com域名解析的ip地址二、找到github.global.ssl.fastly.net网址解析的ip地址三

Mysql的主从同步/复制的原理分析

《Mysql的主从同步/复制的原理分析》:本文主要介绍Mysql的主从同步/复制的原理分析,具有很好的参考价值,希望对大家有所帮助,如有错误或未考虑完全的地方,望不吝赐教... 目录为什么要主从同步?mysql主从同步架构有哪些?Mysql主从复制的原理/整体流程级联复制架构为什么好?Mysql主从复制注意

java -jar命令运行 jar包时运行外部依赖jar包的场景分析

《java-jar命令运行jar包时运行外部依赖jar包的场景分析》:本文主要介绍java-jar命令运行jar包时运行外部依赖jar包的场景分析,本文给大家介绍的非常详细,对大家的学习或工作... 目录Java -jar命令运行 jar包时如何运行外部依赖jar包场景:解决:方法一、启动参数添加: -Xb

Apache 高级配置实战之从连接保持到日志分析的完整指南

《Apache高级配置实战之从连接保持到日志分析的完整指南》本文带你从连接保持优化开始,一路走到访问控制和日志管理,最后用AWStats来分析网站数据,对Apache配置日志分析相关知识感兴趣的朋友... 目录Apache 高级配置实战:从连接保持到日志分析的完整指南前言 一、Apache 连接保持 - 性

Linux中的more 和 less区别对比分析

《Linux中的more和less区别对比分析》在Linux/Unix系统中,more和less都是用于分页查看文本文件的命令,但less是more的增强版,功能更强大,:本文主要介绍Linu... 目录1. 基础功能对比2. 常用操作对比less 的操作3. 实际使用示例4. 为什么推荐 less?5.

spring-gateway filters添加自定义过滤器实现流程分析(可插拔)

《spring-gatewayfilters添加自定义过滤器实现流程分析(可插拔)》:本文主要介绍spring-gatewayfilters添加自定义过滤器实现流程分析(可插拔),本文通过实例图... 目录需求背景需求拆解设计流程及作用域逻辑处理代码逻辑需求背景公司要求,通过公司网络代理访问的请求需要做请

Java集成Onlyoffice的示例代码及场景分析

《Java集成Onlyoffice的示例代码及场景分析》:本文主要介绍Java集成Onlyoffice的示例代码及场景分析,本文给大家介绍的非常详细,对大家的学习或工作具有一定的参考借鉴价值,需要... 需求场景:实现文档的在线编辑,团队协作总结:两个接口 + 前端页面 + 配置项接口1:一个接口,将o

IDEA下"File is read-only"可能原因分析及"找不到或无法加载主类"的问题

《IDEA下Fileisread-only可能原因分析及找不到或无法加载主类的问题》:本文主要介绍IDEA下Fileisread-only可能原因分析及找不到或无法加载主类的问题,具有很好的参... 目录1.File is read-only”可能原因2.“找不到或无法加载主类”问题的解决总结1.File