【问题标题】:Why aren't the Netty HTTP handlers sharable?为什么 Netty HTTP 处理程序不可共享?
【发布时间】:2020-02-20 05:31:55
【问题描述】:

Netty 实例化了一组请求处理程序类whenever a new connection is opened。这对于像 websocket 这样的连接将在 websocket 的生命周期内保持打开状态的东西似乎很好。

当使用 Netty 作为每秒可以接收数千个请求的 HTTP 服务器时,这似乎对垃圾收集产生了相当大的影响。每个请求都会实例化几个类(在我的例子中是 10 个处理程序类),然后在几毫秒后对它们进行垃圾收集。

在中等负载约 1000 个请求/秒的 HTTP 服务器中,这将是 一万 个类实例化和垃圾收集每秒

看来我们可以简单地 查看下面的答案创建可共享的处理程序,使用ChannelHandler.Sharable 消除这种巨大的GC 开销。他们只需要是线程安全的。

但是,我发现库中打包的所有非常基本的 HTTP 处理程序不可共享,例如 HttpServerCodecHttpObjectAggregator。此外,没有一个HTTP handler examples 是可共享的。 99% 的示例代码和教程似乎并不关心它。在 Norman Maurer 的 book(Netty 的作者)中,只有一处说明了使用共享处理程序的原因:

为什么要共享频道处理程序?

安装单个频道处理程序的常见原因 多个 ChannelPipelines 中的 ChannelHandler 是为了收集统计信息 跨多个渠道。

在任何地方都没有提到 GC 负载问题。


Netty 已经在常规生产环境中使用了将近十年。可以说是现有的用于高并发非阻塞 IO 的最常用的 Java 库。

换句话说,它的设计目的远远超过我每秒中等 1000 个请求。

我错过了什么让 GC 负载不成问题吗?

或者,我是否应该尝试实现自己的 Sharable 处理程序,具有类似的功能来解码、编码和编写 HTTP 请求和响应?

【问题讨论】:

    标签: java garbage-collection netty instantiation java-11


    【解决方案1】:

    虽然我们始终致力于在 netty 中产生尽可能少的 GC,但在某些情况下这实际上是不可能的。例如,http 编解码器等保持每个连接的状态,因此无法共享这些状态(即使它们是线程安全的)。

    解决这个问题的唯一方法是将它们池化,但我认为还有其他对象更有可能导致 GC 问题,对于这些我们尝试在可能的情况下进行池化。

    【讨论】:

    • 感谢您的回答,感谢 Netty!我正在尝试按照您的建议实施频道池,但所有 implementations 似乎都严格适用于 Netty clients,您知道您可以指点我的任何服务器实现吗?
    【解决方案2】:

    TL;DR:

    如果您使用默认 HTTP 处理程序达到使 GC 成为问题所需的容量,那么无论如何是时候使用代理服务器进行扩展了。


    在诺曼的回答之后,我最终尝试了一个非常简单的可共享 HTTP 编解码器/聚合器 POC,看看这是否值得追求。

    我的可共享解码器与 RFC 7230 相差甚远,但它满足了我当前项目的要求。

    然后我使用httperfvisualvm 来了解GC 负载差异的概念。对于我的努力,我的 GC 率只下降了 10%。换句话说,它真的没有太大的区别。

    唯一真正值得赞赏的效果是,与使用打包的非共享 HTTP 编解码器 + 聚合器相比,我在运行 1000 个请求/秒时的错误减少了 5%。这只发生在我以 1000 个请求/秒的速度持续超过 10 秒时。

    最后我不会去追求它。将其制作成完全符合 HTTP 的解码器所需的时间量根本不值得花时间。

    出于参考目的,这里是我尝试过的组合可共享解码器/聚合器:

    import java.util.concurrent.ConcurrentHashMap;
    
    import io.netty.buffer.ByteBuf;
    import io.netty.channel.ChannelHandler.Sharable;
    import io.netty.channel.ChannelHandlerContext;
    import io.netty.channel.ChannelId;
    import io.netty.channel.ChannelInboundHandlerAdapter;
    
    @Sharable
    public class SharableHttpDecoder extends ChannelInboundHandlerAdapter {
    
        private static final ConcurrentHashMap<ChannelId, SharableHttpRequest> MAP = 
                new ConcurrentHashMap<ChannelId, SharableHttpRequest>();
        
        @Override
        public void channelRead(ChannelHandlerContext ctx, Object msg) 
            throws Exception 
        {        
            if (msg instanceof ByteBuf) 
            {
                ByteBuf buf = (ByteBuf) msg;
                ChannelId channelId = ctx.channel().id();
                SharableHttpRequest request = MAP.get(channelId);
                                        
                if (request == null)
                {
                    request = new SharableHttpRequest(buf);
                    buf.release();
                    if (request.isComplete()) 
                    {
                        ctx.fireChannelRead(request);
                    }
                    else
                    {
                        MAP.put(channelId, request);
                    }
                }
                else
                {
                    request.append(buf);
                    buf.release();
                    if (request.isComplete()) 
                    {
                        ctx.fireChannelRead(request);
                    }
                }
            }
            else
            {
                // TODO send 501
                System.out.println("WTF is this? " + msg.getClass().getName());
                ctx.fireChannelRead(msg);
            }
        }
        
        @Override
        public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) 
            throws Exception 
        {
            System.out.println("Unable to handle request on channel: " + 
                ctx.channel().id().asLongText());
            cause.printStackTrace(System.err);
            
            // TODO send 500
            ctx.fireExceptionCaught(cause);
            ctx.close();
        }
        
    }
    

    解码器创建的用于在管道上处理的结果对象:

    import java.util.Arrays;
    import java.util.HashMap;
    import io.netty.buffer.ByteBuf;
    
    public class SharableHttpRequest
    {
        
        private static final byte SPACE = 32;
        private static final byte COLON = 58;
        private static final byte CARRAIGE_RETURN = 13;
        
        private HashMap<Header,String> myHeaders;
        private Method myMethod;
        private String myPath;
        private byte[] myBody;
        private int myIndex = 0;
        
        public SharableHttpRequest(ByteBuf buf)
        {
            try
            {
                myHeaders = new HashMap<Header,String>();
                final StringBuilder builder = new StringBuilder(8);
                parseRequestLine(buf, builder);
                while (parseNextHeader(buf, builder));
                parseBody(buf);
            }
            catch (Exception e)
            {
                e.printStackTrace(System.err);
            }
        }
        
        public String getHeader(Header name)
        {
            return myHeaders.get(name);
        }
        
        public Method getMethod()
        {
            return myMethod;
        }
        
        public String getPath()
        {
            return myPath;
        }
        
        public byte[] getBody()
        {
            return myBody;
        }
        
        public boolean isComplete()
        {
            return myIndex >= myBody.length;
        }
        
        public void append(ByteBuf buf)
        {
            int length = buf.readableBytes();
            buf.getBytes(buf.readerIndex(), myBody, myIndex, length);
            myIndex += length;
        }
    
        private void parseRequestLine(ByteBuf buf, StringBuilder builder)
        {
            int idx = buf.readerIndex();
            int end = buf.writerIndex();
            for (; idx < end; ++idx)
            {
                byte next = buf.getByte(idx);
                
                // break on CR
                if (next == CARRAIGE_RETURN)
                {
                    break;
                }
                
                // we need the method
                else if (myMethod == null)
                {
                    if (next == SPACE)
                    {
                        myMethod = Method.fromBuilder(builder);
                        builder.delete(0, builder.length());
                        builder.ensureCapacity(100);
                    }
                    else
                    {
                        builder.append((char) next);
                    }
                }
                
                // we need the path
                else if (myPath == null)
                {
                    if (next == SPACE)
                    {
                        myPath = builder.toString();
                        builder.delete(0, builder.length());
                    }
                    else
                    {
                        builder.append((char) next);
                    }
                }
                
                // don't need the version right now
            }
            idx += 2; // skip line endings
            buf.readerIndex(idx);
        }
        
        private boolean parseNextHeader(ByteBuf buf, StringBuilder builder)
        {
            Header header = null;
            int idx = buf.readerIndex();
            int end = buf.writerIndex();
            for (; idx < end; ++idx)
            {
                byte next = buf.getByte(idx);
                
                // break on CR
                if (next == CARRAIGE_RETURN)
                {
                    if (header != Header.UNHANDLED)
                    {
                        myHeaders.put(header,builder.toString());
                        builder.delete(0, builder.length());
                    }
                    break;
                }
                
                else if (header == null)
                {
                    // we have the full header name
                    if (next == COLON)
                    {
                        header = Header.fromBuilder(builder);
                        builder.delete(0, builder.length());
                    }
    
                    // get header name as lower case for mapping purposes
                    else
                    {
                        builder.append(next > 64 && next < 91 ? 
                            (char) ( next | 32 ) : (char) next);
                    }
                }
                
                // we don't care about some headers
                else if (header == Header.UNHANDLED)
                {
                    continue;
                }
                
                // skip initial spaces
                else if (builder.length() == 0 && next == SPACE)
                {
                    continue;
                }
                
                // get the header value
                else
                {
                    builder.append((char) next);
                }
            }
            
            idx += 2; // skip line endings
            buf.readerIndex(idx);
            
            if (buf.getByte(idx) == CARRAIGE_RETURN)
            {
                idx += 2; // skip line endings
                buf.readerIndex(idx);
                return false;
            }
            else
            {
                return true;
            }
        }
        
        private void parseBody(ByteBuf buf)
        {
            int length = buf.readableBytes();
            if (length == 0)
            {
                myBody = new byte[0];
                myIndex = 1;
            }
            else
            {
                System.out.println("Content-Length: " + myHeaders.get(Header.CONTENT_LENGTH));
                if (myHeaders.get(Header.CONTENT_LENGTH) != null)
                {
                    int totalLength = Integer.valueOf(myHeaders.get(Header.CONTENT_LENGTH));
                    myBody = new byte[totalLength];
                    buf.getBytes(buf.readerIndex(), myBody, myIndex, length);
                    myIndex += length;
                }
                
                // TODO handle chunked
            }
        }
        
        
        
        
        public enum Method
        {
            GET(new char[]{71, 69, 84}), 
            POST(new char[]{80, 79, 83, 84}),
            UNHANDLED(new char[]{}); // could be expanded if needed
            
            private char[] chars;
    
            Method(char[] chars) 
            {
                this.chars = chars;
            }
            
            public static Method fromBuilder(StringBuilder builder) 
            {
                for (Method method : Method.values()) 
                {
                    if (method.chars.length == builder.length()) 
                    {
                        boolean match = true;
                        for (int i = 0; i < builder.length(); i++) 
                        {
                            if (method.chars[i] != builder.charAt(i)) 
                            {
                                match = false;
                                break;
                            }
                        }
                        
                        if (match)
                        {
                            return method;
                        }
                    }
                }
                return null;
            }
        }
        
        public enum Header
        {
            HOST(new char[]{104, 111, 115, 116}), 
            CONNECTION(new char[]{99, 111, 110, 110, 101, 99, 116, 105, 111, 110}),
            IF_MODIFIED_SINCE(new char[]{
                105, 102, 45, 109, 111, 100, 105, 102, 105, 101, 100, 45, 115, 
                105, 110, 99, 101}),
            COOKIE(new char[]{99, 111, 111, 107, 105, 101}),
            CONTENT_LENGTH(new char[]{
                99, 111, 110, 116, 101, 110, 116, 45, 108, 101, 110, 103, 116, 104}),
            UNHANDLED(new char[]{}); // could be expanded if needed
            
            private char[] chars;
    
            Header(char[] chars) 
            {
                this.chars = chars;
            }
            
            public static Header fromBuilder(StringBuilder builder) 
            {
                for (Header header : Header.values()) 
                {
                    if (header.chars.length == builder.length()) 
                    {                    
                        boolean match = true;
                        for (int i = 0; i < builder.length(); i++) 
                        {
                            if (header.chars[i] != builder.charAt(i)) 
                            {
                                match = false;
                                break;
                            }
                        }
                        
                        if (match)
                        {
                            return header;
                        }
                    }
                }
                return UNHANDLED;
            }
        }
    }
    

    一个简单的测试处理程序:

    import io.netty.buffer.ByteBuf;
    import io.netty.channel.ChannelFuture;
    import io.netty.channel.ChannelFutureListener;
    import io.netty.channel.ChannelHandler.Sharable;
    import io.netty.channel.ChannelHandlerContext;
    import io.netty.channel.SimpleChannelInboundHandler;
    import io.netty.util.CharsetUtil;
    
    @Sharable
    public class SharableHttpHandler extends SimpleChannelInboundHandler<SharableHttpRequest>
    {    
        @Override
        protected void channelRead0(ChannelHandlerContext ctx, SharableHttpRequest msg) 
            throws Exception
        {
            String message = "HTTP/1.1 200 OK\r\n" +
                    "Content-type: text/html\r\n" + 
                    "Content-length: 42\r\n\r\n" + 
                    "<html><body>Hello sharedworld</body><html>";
            
            ByteBuf buffer = ctx.alloc().buffer(message.length());
            buffer.writeCharSequence(message, CharsetUtil.UTF_8);
            ChannelFuture flushPromise = ctx.channel().writeAndFlush(buffer);
            flushPromise.addListener(ChannelFutureListener.CLOSE);
            if (!flushPromise.isSuccess()) 
            {
                flushPromise.cause().printStackTrace(System.err);
            }
        }    
    }
    

    使用这些可共享处理程序的完整管道:

    import tests.SharableHttpDecoder;
    import tests.SharableHttpHandler;
    
    import io.netty.channel.ChannelInitializer;
    import io.netty.channel.ChannelPipeline;
    import io.netty.channel.socket.SocketChannel;
    
    public class ServerPipeline extends ChannelInitializer<SocketChannel>
    {
        private final SharableHttpDecoder decoder = new SharableHttpDecoder();
        private final SharableHttpHandler handler = new SharableHttpHandler();
    
        @Override
        public void initChannel(SocketChannel channel)
        {
            ChannelPipeline pipeline = channel.pipeline();
            pipeline.addLast(decoder);
            pipeline.addLast(handler);
            
        }
    }
    

    以上内容针对这个(更常见的)非共享管道进行了测试:

    import static io.netty.handler.codec.http.HttpResponseStatus.OK;
    import static io.netty.handler.codec.http.HttpVersion.HTTP_1_1;
    
    import io.netty.buffer.ByteBuf;
    import io.netty.channel.ChannelFuture;
    import io.netty.channel.ChannelFutureListener;
    import io.netty.channel.ChannelHandlerContext;
    import io.netty.channel.ChannelInitializer;
    import io.netty.channel.ChannelPipeline;
    import io.netty.channel.SimpleChannelInboundHandler;
    import io.netty.channel.socket.SocketChannel;
    import io.netty.handler.codec.http.DefaultFullHttpResponse;
    import io.netty.handler.codec.http.FullHttpRequest;
    import io.netty.handler.codec.http.FullHttpResponse;
    import io.netty.handler.codec.http.HttpHeaderNames;
    import io.netty.handler.codec.http.HttpHeaderValues;
    import io.netty.handler.codec.http.HttpObjectAggregator;
    import io.netty.handler.codec.http.HttpServerCodec;
    import io.netty.handler.codec.http.HttpUtil;
    import io.netty.util.CharsetUtil;
    
    public class ServerPipeline extends ChannelInitializer<SocketChannel>
    {
    
        @Override
        public void initChannel(SocketChannel channel)
        {
            ChannelPipeline pipeline = channel.pipeline();
            pipeline.addLast(new HttpServerCodec());
            pipeline.addLast(new HttpObjectAggregator(65536));
            pipeline.addLast(new UnsharedHttpHandler());
            
        }
        
        class UnsharedHttpHandler extends SimpleChannelInboundHandler<FullHttpRequest>
        {
    
            @Override
            public void channelRead0(ChannelHandlerContext ctx, FullHttpRequest request) 
                throws Exception
            {
                String message = "<html><body>Hello sharedworld</body><html>";
                ByteBuf buffer = ctx.alloc().buffer(message.length());
                buffer.writeCharSequence(message.toString(), CharsetUtil.UTF_8);
    
                FullHttpResponse response = new DefaultFullHttpResponse(HTTP_1_1, OK, buffer);
                response.headers().set(HttpHeaderNames.CONTENT_TYPE, "text/html; charset=UTF-8");
                HttpUtil.setContentLength(response, response.content().readableBytes());
                response.headers().set(HttpHeaderNames.CONNECTION, HttpHeaderValues.CLOSE);
                ChannelFuture flushPromise = ctx.writeAndFlush(response);
                flushPromise.addListener(ChannelFutureListener.CLOSE);
                            
            }
        }
    }
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2019-09-16
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多