[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

相关文章

Springboot中分析SQL性能的两种方式详解

《Springboot中分析SQL性能的两种方式详解》文章介绍了SQL性能分析的两种方式:MyBatis-Plus性能分析插件和p6spy框架,MyBatis-Plus插件配置简单,适用于开发和测试环... 目录SQL性能分析的两种方式:功能介绍实现方式:实现步骤:SQL性能分析的两种方式:功能介绍记录

最长公共子序列问题的深度分析与Java实现方式

《最长公共子序列问题的深度分析与Java实现方式》本文详细介绍了最长公共子序列(LCS)问题,包括其概念、暴力解法、动态规划解法,并提供了Java代码实现,暴力解法虽然简单,但在大数据处理中效率较低,... 目录最长公共子序列问题概述问题理解与示例分析暴力解法思路与示例代码动态规划解法DP 表的构建与意义动

C#使用DeepSeek API实现自然语言处理,文本分类和情感分析

《C#使用DeepSeekAPI实现自然语言处理,文本分类和情感分析》在C#中使用DeepSeekAPI可以实现多种功能,例如自然语言处理、文本分类、情感分析等,本文主要为大家介绍了具体实现步骤,... 目录准备工作文本生成文本分类问答系统代码生成翻译功能文本摘要文本校对图像描述生成总结在C#中使用Deep

Go中sync.Once源码的深度讲解

《Go中sync.Once源码的深度讲解》sync.Once是Go语言标准库中的一个同步原语,用于确保某个操作只执行一次,本文将从源码出发为大家详细介绍一下sync.Once的具体使用,x希望对大家有... 目录概念简单示例源码解读总结概念sync.Once是Go语言标准库中的一个同步原语,用于确保某个操

Redis主从/哨兵机制原理分析

《Redis主从/哨兵机制原理分析》本文介绍了Redis的主从复制和哨兵机制,主从复制实现了数据的热备份和负载均衡,而哨兵机制可以监控Redis集群,实现自动故障转移,哨兵机制通过监控、下线、选举和故... 目录一、主从复制1.1 什么是主从复制1.2 主从复制的作用1.3 主从复制原理1.3.1 全量复制

Redis主从复制的原理分析

《Redis主从复制的原理分析》Redis主从复制通过将数据镜像到多个从节点,实现高可用性和扩展性,主从复制包括初次全量同步和增量同步两个阶段,为优化复制性能,可以采用AOF持久化、调整复制超时时间、... 目录Redis主从复制的原理主从复制概述配置主从复制数据同步过程复制一致性与延迟故障转移机制监控与维

Redis连接失败:客户端IP不在白名单中的问题分析与解决方案

《Redis连接失败:客户端IP不在白名单中的问题分析与解决方案》在现代分布式系统中,Redis作为一种高性能的内存数据库,被广泛应用于缓存、消息队列、会话存储等场景,然而,在实际使用过程中,我们可能... 目录一、问题背景二、错误分析1. 错误信息解读2. 根本原因三、解决方案1. 将客户端IP添加到Re

Java汇编源码如何查看环境搭建

《Java汇编源码如何查看环境搭建》:本文主要介绍如何在IntelliJIDEA开发环境中搭建字节码和汇编环境,以便更好地进行代码调优和JVM学习,首先,介绍了如何配置IntelliJIDEA以方... 目录一、简介二、在IDEA开发环境中搭建汇编环境2.1 在IDEA中搭建字节码查看环境2.1.1 搭建步

Redis主从复制实现原理分析

《Redis主从复制实现原理分析》Redis主从复制通过Sync和CommandPropagate阶段实现数据同步,2.8版本后引入Psync指令,根据复制偏移量进行全量或部分同步,优化了数据传输效率... 目录Redis主DodMIK从复制实现原理实现原理Psync: 2.8版本后总结Redis主从复制实

锐捷和腾达哪个好? 两个品牌路由器对比分析

《锐捷和腾达哪个好?两个品牌路由器对比分析》在选择路由器时,Tenda和锐捷都是备受关注的品牌,各自有独特的产品特点和市场定位,选择哪个品牌的路由器更合适,实际上取决于你的具体需求和使用场景,我们从... 在选购路由器时,锐捷和腾达都是市场上备受关注的品牌,但它们的定位和特点却有所不同。锐捷更偏向企业级和专