Netty thread model vs long-running taks

Viewed 733

Update: OK, since it was my first post here, I was probably too verbose looking at the zero feedback I received. So, I decided to add a diagram to better visualize my problem. Maybe some one will take a look (I'm looking at you, Norman Maurer ;o) )

netty threading problem

Can one please confirm my findings and explain how I can prevent the long running task to block the other short requests? More generally, what is the recommended solution that scales without locking in case I have to deal with long running jobs on a netty server?


I'm learning netty 4/5 for the past 2 weeks and I read a lot about its threading model (including "Netty in Action") then I checked my understanding with a simple test application where I wanted to see the response of the netty server against multiple TCP simple concurrent clients, some of them requesting long-running jobs on the server. My results were not as expected.

My setup:

  • netty 4.1.45
  • jdk 1.8 on windows 10
  • 15 concurrent connections are opened to the server
  • 1 client requests a long running-job (simulated as a sleep on the Server side)

My server initialization is as follows:

NioEventLoopGroup bossGroup = new NioEventLoopGroup(1);
NioEventLoopGroup workGroup = new NioEventLoopGroup();

ServerBootstrap bootstrap = new ServerBootstrap();

try {
    bootstrap
        .group(bossGroup, workGroup)
        .channel(NioServerSocketChannel.class)
        .childHandler(new MyChannelInitializer(new DefaultEventExecutorGroup(5)));
            
    ChannelFuture future = bootstrap.bind(new InetSocketAddress(SERVER_PORT));
    future.addListener(new MyChannelFutureListener());

    future.channel().closeFuture().sync();
            
} catch (InterruptedException e) {
    logger.error("Error: ", e);
} finally {
    bossGroup.shutdownGracefully().sync();
}

Inside the channel's initialiser I have the following code:

public class MyChannelInitializer extends ChannelInitializer<Channel> {

    private EventExecutorGroup handlerGroup = null;

    public MyChannelInitializer() {
    }

    public MyChannelInitializer(EventExecutorGroup handlerGroup) {
        this.handlerGroup = handlerGroup;
    }

    @Override
    protected void initChannel(Channel ch) throws Exception {
        ch.pipeline().addLast(new LoggingHandler(LogLevel.INFO));
        ch.pipeline().addLast(handlerGroup, new MyTestInboudHandler());
    }
}

As you can see, I want my inbound handler to be executed in a separate eventExecutionGroup so I don't block the NIO event loops. The logic in my inbound handler is very simple:

public class MyTestInboudHandler extends SimpleChannelInboundHandler<ByteBuf> {

    private static final Logger logger = LogManager.getLogger((MyTestInboudHandler.class));

    @Override
    protected void channelRead0(ChannelHandlerContext ctx, ByteBuf msg) throws Exception {

        if (msg.isReadable()) {
            String content = msg.toString(CharsetUtil.UTF_8);
            int len = content.length();
            if (len > 120) {
                content = content.substring(0, 119);
            }

            String clientId = content.substring(0, content.indexOf(':'));
            String payload = content.substring(content.indexOf(':') + 1);

            logger.info(ctx.channel().hashCode() + ": received message of length: " + len + " -- " + content + "<<<");

            if (payload.equals("LONG_TASK_REQ")) {
                // simulate load run
                Thread.sleep(50_000L);
                ctx.writeAndFlush(Unpooled.copiedBuffer("DoneBIG", CharsetUtil.UTF_8)).addListener(ChannelFutureListener.CLOSE);
            } else {
                ctx.writeAndFlush(Unpooled.copiedBuffer("Done", CharsetUtil.UTF_8)).addListener(ChannelFutureListener.CLOSE);
            }

        } else {
            logger.info("ByteBuf not readable...");
        }
    }

    @Override
    public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) throws Exception {

        logger.error("{} - Exception in input handler: {} ", ctx.channel().hashCode(), cause.getMessage());
        ctx.close();

    }
}

So if the client id I find in the request payload equals 13, it means there is a long job to be done (e.g. JDBC call to a database) which I simulate with a big 50 sec sleep. All other requests are handled instantly by simply replying "Done" and closing the connection.

The result I expected was to see that all clients finish their server's call instantly, except for the client 13 which would take 50 sec to complete. Below is the client implementation, also using netty:

public class MyNettyClient extends Thread {

    private static final Logger logger = LogManager.getLogger(MyNettyClient.class);

    private int id;

    private MyNettyClient() {
    }

    public MyNettyClient(int id) {
        this();
        this.id = id;
    }

    public void run() {

        logger.info("Netty client starting with id: " + this.id);

        EventLoopGroup group = new NioEventLoopGroup();

        try {
            Bootstrap bootstrap = new Bootstrap();

            bootstrap.group(group)
                .channel(NioSocketChannel.class)
                    .remoteAddress(new InetSocketAddress("127.0.0.1", 8778))
                .handler(new ChannelInitializer<SocketChannel>() {
                        @Override
                        protected void initChannel(SocketChannel ch) throws Exception {
                            ch.pipeline().addLast(new MyNettyClientHandler(id));
                        }
                    });

            ChannelFuture future = bootstrap.connect().sync();
            future.channel().closeFuture().sync();
        } catch (Exception e) {
            logger.error("Error from client " + id + " : ", e);
        } finally {
            logger.info("Exiting client: " + id);
            try {
                group.shutdownGracefully().sync();
            } catch (InterruptedException e) {
                logger.error("Error from client " + id, e);
            }
        }

    }

    public static void main(String[] args) throws InterruptedException, ExecutionException {

        ExecutorService executor = Executors.newCachedThreadPool();
        List<Future<?>> futureList = new ArrayList<Future<?>>();

        for (int i = 1; i < 15; ++i) {
            futureList.add(executor.submit(new MyNettyClient(i)));
        }

        for (Future<?> f : futureList) {
            f.get();
            if (f.isCancelled()) {
                logger.info("One future was cancelled: " + f.toString());
            }
        }

        logger.info("Calling shutdown on executor");
        executor.shutdown();

        logger.info("MAIN THREAD DONE");
    }
}

The logic for the handler writing the calls to the Netty server is below:

public class MyNettyClientHandler extends SimpleChannelInboundHandler<ByteBuf> {

    private static final Logger logger = LogManager.getLogger(MyNettyClientHandler.class);

    private int clientId;

    public MyNettyClientHandler(int id) {
        clientId = id;
    }

    @Override
    protected void channelRead0(ChannelHandlerContext ctx, ByteBuf msg) throws Exception {

        logger.info("@@@@@Message received: " + msg.toString(CharsetUtil.UTF_8));
    }

    @Override
    public void channelRead(ChannelHandlerContext ctx, Object msg) throws Exception {

        ByteBuf b = (ByteBuf) msg;
        String serverMsg = b.toString(CharsetUtil.UTF_8);
        logger.info("Read from server by client " + clientId + ": " + serverMsg);

    }

    @Override
    public void channelActive(ChannelHandlerContext ctx) throws Exception {

        String content = null;

        if (clientId % 13 == 0) {
            content = clientId + ":LONG_TASK_REQ";
        } else {
            content = clientId + ":SMALL_TASK_REQ";
        }
        ChannelFuture f = ctx.writeAndFlush(Unpooled.copiedBuffer(content, CharsetUtil.UTF_8)).sync();

        if (f.isDone()) {
             logger.info("Send for client {} completed successfully.", clientId);
        } else {
            logger.error("Future not done for client {}. Reason: {}", clientId, f.cause().getMessage());
        }

        logger.info("Finish sending from client " + clientId);
    }

    @Override
    public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) throws Exception {

        logger.error("Exception in client {}, type: {} ", clientId, cause.getMessage());

        ctx.close();
    }
}

The results I got show that when the long request is received by the Server and processed in 50 sec, it always blocks other requests made on different channels by different clients. It seems that channels are bound to a certain thread from the EventExecutionGroup at activation time and all the activity conducted on that channel afterwards will always be processed by the same thread, no matter how busy is that thread and how many idle threads are in the pool. This is something I really don't understand because the whole purpose of the EventExecutionGroup was to handle the long synchronous jobs without any interference to the other clients.

I know I could implement my own pool of working threads to handle and synchronize the long jobs but the question remains: why we can write ch.pipeline().addLast(handlerGroup, new MyTestInboudHandler()); if the result is so limiting.

Please advise.

1 Answers

An EventLoopGroup is just a collection of EventLoops. Each EventLoop is a single thread that is assigned to one or more Channels and belongs to that channel for the duration of its life. This relationship is 1:many where a single EventLoop can handle many Channels.

As channels are registered with the EventLoopGroup, it will assign them to either a new EventLoop (if capacity allows) or an existing EventLoop (via round robin I believe). This means that it is totally possible that a long running task will be running on a channel's EventLoop that just so happens to be handling Channels from other clients that are running short running tasks, leading to blocking that you experienced.

One thing you can do is offload these long running tasks to a separate ExecutorService from inside your MyTestInboundHandler. This way, a long running task will never block, rather just dispatch handling to another thread.

Keep note though, that when you offload a task to a separate ExecutorService, it will run in a separate thread, and to write to the channel, the write operation Channel.writeAndFlush(data) will actually enqueue the task to the Channel's EventLoop queue and run it on the next iteration, it will NOT run it immediately.

Ex:

@Override
protected void channelRead0(ChannelHandlerContext ctx, ByteBuf msg) throws Exception {
    myExecutor.submit(() -> {
        byte[] response = doSomeLongTask(msg);
        ctx.channel().writeAndFlush(response); // Submits to channel's EventLoop queue to be written
    }
}

Now to answer the main question: What's the difference between this and adding an EventExecutor to each handler?

You can almost achieve the same results by performing

ch.pipeline().addLast(handlerGroup, new MyTestInboudHandler());

The main difference being that [1] in the case that you add your own handlerGroup to each handler, each socket event will be run within that group(emphasis mine):

A user can specify an EventExecutor when he or she adds a handler to a ChannelPipeline.

  • If specified, the handler methods of the ChannelHandler are always invoked by the specified EventExecutor.
  • If unspecified, the handler methods are always invoked by the EventLoop that its associated Channel is registered to.

EventExecutor and EventLoop assigned to a handler or a channel are always single-threaded.

  • The handler methods will always be invoked by the same thread.
  • If multithreaded EventExecutor or EventLoop is specified, one of the threads will be chosen first and then the chosen thread will be used until deregistration.
  • If two handlers in the same pipeline are assigned with different >EventExecutors, they are invoked simultaneously. A user has to pay attention to thread safety if more than one handler access shared data even if the shared data is accessed only by the handlers in the same pipeline.

This means that for each handler event that occurs on a pipeline, and for a handler that is assigned its own executor, those entire methods (such as channelRead0()) are instead run in their own thread from the specified executor.

If you don't want to specify an executor for a handler, this can also be achieved by overriding each handler method and immediately offloading the call to a custom executor. NOTE: You will need to override each method, any that are not overridden will still be run in the netty IO EventLoop thread.

So in essence, they are both almost identical. By using an executor assigned to a handler, you are still constrained by a single thread as Netty will reuse the same thread for each handler. But the benefits are that each handler method will be automatically run in a separate thread than the channel's assigned EventLoop

By using your own executor and offloading tasks to it, the EventLoop is still the one calling each handler method in the same thread, but you have the ability to offload the actual work done to your own executor. This way the executor can also utilize more than just 1 thread to process each of these tasks concurrently.

Related