Netty TCP服务端打印客户端连接数量

1. 自建hadler 继承  ChannelDuplexHandler 

   1.1完整代码

package com.lgdz.netty.server;

import com.codahale.metrics.ConsoleReporter; import com.codahale.metrics.Gauge; import com.codahale.metrics.MetricRegistry; import com.codahale.metrics.jmx.JmxReporter; import io.netty.channel.ChannelDuplexHandler; import io.netty.channel.ChannelHandler; import io.netty.channel.ChannelHandlerContext; import io.netty.channel.ChannelPromise; import lombok.extern.slf4j.Slf4j;

import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicLong;

@Slf4j @ChannelHandler.Sharable public class MetricsHandler extends ChannelDuplexHandler {

1**private** AtomicLong **totalConnectionNumber** \= **new** AtomicLong(); 2 3{ 4 5 MetricRegistry metricRegistry = **new** MetricRegistry(); 6 7 metricRegistry.register(**"totalConnectionNumber"**, **new** Gauge<Long>() { 8 @Override

public Long getValue() { return totalConnectionNumber.longValue(); } });

1 ConsoleReporter consoleReporter = ConsoleReporter._forRegistry_(metricRegistry).build(); 2 consoleReporter.start(10, TimeUnit._SECONDS_); **//10秒检测一次TCP客户端连接数** 3 4 JmxReporter jmxReporter = JmxReporter._forRegistry_(metricRegistry).build(); 5 jmxReporter.start(); 6 7} 8 9@Override

public void channelActive(ChannelHandlerContext ctx) throws Exception { totalConnectionNumber.incrementAndGet(); super.channelActive(ctx); }

@Override

public void channelInactive(ChannelHandlerContext ctx) throws Exception { totalConnectionNumber.decrementAndGet(); super.channelInactive(ctx); }

@Override

public void channelRead(ChannelHandlerContext ctx, Object msg) throws Exception { _//_在这里可以处理硬件发送过来的数据 _// log.debug("数据对象长度:" + ((byte[]) msg).length); //java.lang.ClassCastException: java.lang.String cannot be cast to [B _ } @Override public void write(ChannelHandlerContext ctx, Object msg, ChannelPromise promise) throws Exception { super.write(ctx, msg, promise); } }

2.把 metricsHandler 添加到  pipeline   。 

   2.1完整代码

package com.lgdz.netty.server;

import io.netty.bootstrap.ServerBootstrap; import io.netty.channel.Channel; import io.netty.channel.ChannelInitializer; import io.netty.channel.ChannelOption; 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 io.netty.handler.logging.LogLevel; import io.netty.handler.logging.LoggingHandler; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.beans.BeansException; import org.springframework.beans.factory.annotation.Value; import org.springframework.context.ApplicationContext; import org.springframework.context.ApplicationContextAware; import org.springframework.stereotype.Component;

import javax.annotation.PostConstruct; import javax.annotation.PreDestroy; import javax.annotation.Resource; import java.util.HashMap; import java.util.Map;

_/** _ * @Author: _代 _ * _@Description:netty__服务器配置 _ * @Date: _Created in 10:20 2020/12/29 _ */ @Component _//_实现__ApplicationContextAware__以获得__ApplicationContext__中的所有__bean public class NettyServer implements ApplicationContextAware {

1**private static final** Logger _logger_ \= LoggerFactory._getLogger_(NettyServer.**class**); 2**private** Channel **channel**; 3**private** EventLoopGroup **bossGroup**; 4**private** EventLoopGroup **workerGroup**; 5@Resource

private HelloServerInHandler helloServerInHandler;

1**private** Map<String, Object> **exportServiceMap** \= **new** HashMap<String, Object>(); 2 3@Value(**"${dai.server.host}"**) 4String **host**; 5 6@Value(**"${rpcServer.ioThreadNum:5}"**) 7**int** **ioThreadNum**; 8_//__内核为此套接口排队的最大连接个数,对于给定的监听套接口,内核要维护两个队列,未链接队列和已连接队列大小总和最大值_

@Value("${rpcServer.backlog:1024}") int backlog;

1@Value(**"${dai.server.port}"**) 2**int** **port**; 3 4_/\*\*

_ * _启动 _ * @throws _InterruptedException _ */

@PostConstruct public void start() { logger.info("begin to start rpc server"); // 主从 Reactor _多线程模式 _ bossGroup = new NioEventLoopGroup(); workerGroup = new NioEventLoopGroup(ioThreadNum);

1 MetricsHandler metricsHandler = **new** MetricsHandler(); 2 3 ServerBootstrap serverBootstrap = **new** ServerBootstrap(); 4 serverBootstrap.group(**bossGroup**, **workerGroup**) 5 6 .channel(NioServerSocketChannel.**class**) 7 .option(ChannelOption._SO\_BACKLOG_, **backlog**) 8 _//__注意是__childOption

_ .childOption(ChannelOption.SO_KEEPALIVE, true) .childOption(ChannelOption.TCP_NODELAY, true) .childHandler(new ChannelInitializer<SocketChannel>() { @Override protected void initChannel(SocketChannel socketChannel) throws Exception { socketChannel.pipeline() _//真实数据最大字节数为__Integer.MAX_VALUE,解码时自动去掉前面四个字节 _ _//io.netty.handler.codec.DecoderException: java.lang.IndexOutOfBoundsException: readerIndex(900) + length(176) exceeds writerIndex(1024): UnpooledUnsafeDirectByteBuf(ridx: 900, widx: 1024, cap: 1024) _ .addLast("logging", new LoggingHandler(LogLevel.INFO)) _/* .addLast(new MyCustomMessageDecoder()) _ .addLast(new MyEncode())*/ .addLast(helloServerInHandler)

                     _/\*       .addLast("decoder",new MyDecode())

_ .addLast("encoder",new MyEncode())*/ /* .addLast(new MyCustomMessageDecoder()) .addLast(new MyEncode())*/ .addLast("metricHandler", metricsHandler); //检测客户端连接数

1 } 2 }); 3 4 **try** { 5 **channel** \= serverBootstrap.bind(**host**,**port**).sync().channel(); 6 } **catch** (InterruptedException e) { 7 **channel**.close(); 8 **return**; 9 } 10 _logger_.info(**"========================================================================================"**); 11 _logger_.info(**"NettyRPC server listening on port "** \+ **port** \+ **" and ready for connections..."**); 12 _logger_.info(**"========================================================================================"**); 13} 14 15@PreDestroy

public void stop() { logger.info("destroy server resources"); if (null == channel) { logger.error("server channel is null"); } bossGroup.shutdownGracefully(); workerGroup.shutdownGracefully(); channel.closeFuture().syncUninterruptibly(); bossGroup = null; workerGroup = null; channel = null; }

_/\*\*

_ * _利用此方法获取__spring ioc__接管的所有__bean _ * @param _ctx _ * @throws _BeansException _ */ public void setApplicationContext(ApplicationContext ctx) throws BeansException { Map<String, Object> serviceMap = ctx.getBeansWithAnnotation(ServiceExporter.class); // 获取所有带有 ServiceExporter 注解的 Spring Bean logger.info("取到所有的RPC:{}", serviceMap); if (serviceMap != null && serviceMap.size() > 0) { for (Object serviceBean : serviceMap.values()) { String interfaceName = serviceBean.getClass().getAnnotation(ServiceExporter.class) .targetInterface() .getName(); logger.info("register service mapping:{}",interfaceName); exportServiceMap.put(interfaceName, serviceBean); } }else{ System.out.println("kong======================================="); } } }

3.效果

3.1

点赞
收藏

评论区

加载中...

相关推荐

MySQL:[Err] 1292 - Incorrect datetime value: ‘0000-00-00 00:00:00‘ for column ‘CREATE_TIME‘ at row 1

文章目录问题用navicat导入数据时,报错:原因这是因为当前的MySQL不支持datetime为0的情况。解决修改sql\mode:sql\mode:SQLMode定义了MySQL应支持的SQL语法、数据校验等,这样可以更容易地在不同的环境中使用MySQL。全局s

Oracle 分组与拼接字符串同时使用

SELECTT.,ROWNUMIDFROM(SELECTT.EMPLID,T.NAME,T.BU,T.REALDEPART,T.FORMATDATE,SUM(T.S0)S0,MAX(UPDATETIME)CREATETIME,LISTAGG(TOCHAR(

MySQL部分从库上面因为大量的临时表tmp_table造成慢查询

背景描述Time:20190124T00:08:14.70572408:00User@Host:@Id:Schema:sentrymetaLast_errno:0Killed:0Query_time:0.315758Lock_

皕杰报表之UUID

​在我们用皕杰报表工具设计填报报表时,如何在新增行里自动增加id呢?能新增整数排序id吗?目前可以在新增行里自动增加id,但只能用uuid函数增加UUID编码,不能新增整数排序id。uuid函数说明:获取一个UUID,可以在填报表中用来创建数据ID语法:uuid()或uuid(sep)参数说明:sep布尔值,生成的uuid中是否包含分隔符'',缺省为

手写Java HashMap源码

HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程22

java将前端的json数组字符串转换为列表

记录下在前端通过ajax提交了一个json数组的字符串,在后端如何转换为列表。前端数据转化与请求varcontracts{id:'1',name:'yanggb合同1'},{id:'2',name:'yanggb合同2'},{id:'3',name:'yang