【问题标题】:Java Netty load testing issuesJava Netty 负载测试问题
【发布时间】:2012-01-24 10:40:48
【问题描述】:

我使用文本协议编写了​​接受连接和轰炸消息(约 100 字节)的服务器,并且我的实现能够使用第 3 方客户端发送大约 400K/秒的环回消息。我为这个任务选择了 Netty,SUSE 11 RealTime,JRockit RTS。 但是当我开始基于 Netty 开发自己的客户端时,我面临着吞吐量的急剧下降(从 400K 降至 1.3K msg/sec)。客户端的代码非常简单。请您提供建议或展示如何编写更有效的客户端的示例。实际上,我更关心延迟,但从吞吐量测试开始,我认为环回时有 1.5Kmsg/秒是不正常的。 附言客户端的目的只是接收来自服务器的消息,很少发送 heartbits。

Client.java

public class Client {

private static ClientBootstrap bootstrap;
private static Channel connector;
public static boolean start()
{
    ChannelFactory factory =
        new NioClientSocketChannelFactory(
                Executors.newCachedThreadPool(),
                Executors.newCachedThreadPool());
    ExecutionHandler executionHandler = new ExecutionHandler( new OrderedMemoryAwareThreadPoolExecutor(16, 1048576, 1048576));

    bootstrap = new ClientBootstrap(factory);

    bootstrap.setPipelineFactory( new ClientPipelineFactory() );

    bootstrap.setOption("tcpNoDelay", true);
    bootstrap.setOption("keepAlive", true);
    bootstrap.setOption("receiveBufferSize", 1048576);
    ChannelFuture future = bootstrap
            .connect(new InetSocketAddress("localhost", 9013));
    if (!future.awaitUninterruptibly().isSuccess()) {
        System.out.println("--- CLIENT - Failed to connect to server at " +
                           "localhost:9013.");
        bootstrap.releaseExternalResources();
        return false;
    }

    connector = future.getChannel();

    return connector.isConnected();
}
public static void main( String[] args )
{
    boolean started = start();
    if ( started )
        System.out.println( "Client connected to the server" );
}

}

ClientPipelineFactory.java

public class ClientPipelineFactory  implements ChannelPipelineFactory{

private final ExecutionHandler executionHandler;
public ClientPipelineFactory( ExecutionHandler executionHandle )
{
    this.executionHandler = executionHandle;
}
@Override
public ChannelPipeline getPipeline() throws Exception {
    ChannelPipeline pipeline = pipeline();
    pipeline.addLast("framer", new DelimiterBasedFrameDecoder(
              1024, Delimiters.lineDelimiter()));
    pipeline.addLast( "executor", executionHandler);
    pipeline.addLast("handler", new MessageHandler() );

    return pipeline;
}

}

MessageHandler.java
public class MessageHandler extends SimpleChannelHandler{

long max_msg = 10000;
long cur_msg = 0;
long startTime = System.nanoTime();
@Override
public void messageReceived(ChannelHandlerContext ctx, MessageEvent e) {
    cur_msg++;

    if ( cur_msg == max_msg )
    {
        System.out.println( "Throughput (msg/sec) : " + max_msg* NANOS_IN_SEC/(     System.nanoTime() - startTime )   );
        cur_msg = 0;
        startTime = System.nanoTime();
    }
}

@Override
public void exceptionCaught(ChannelHandlerContext ctx, ExceptionEvent e) {
    e.getCause().printStackTrace();
    e.getChannel().close();
}

}

更新。在服务器端,有一个定期线程写入接受的客户端通道。并且通道很快变得不可写。 更新 N2。在管道中添加了 OrderedMemoryAwareExecutor,但吞吐量仍然非常低(大约 4k msg/sec)

已修复。我将 executor 放在整个管道堆栈的前面,它成功了!

【问题讨论】:

  • 我会在消息中发送时间戳并获取每条消息的延迟。这可能会让您更详细地了解延迟是什么。如果您只在同一台主机上进行通信,并且延迟对您很重要,您可以考虑改用共享内存。
  • 其实服务器会是远程的。我实现了模拟器来测试客户端每秒可以处理多少消息。事实证明,朴素的 netty 客户端实现很慢。
  • 您的代码显示为 ItchClientPipelineFactory,但粘贴的代码用于 ClientPipelineFactory。这只是一个命名错误还是 ItchClientPipelineFactory 中包含非最佳代码(即没有使用正确的消息处理程序)并且您忘记了您仍在使用它的情况?
  • 你看到这篇文章了吗,试试它的建议:stackoverflow.com/questions/8444267/…
  • 感谢链接,但我已经看到了,但它没有帮助。似乎服务器写入客户端可以处理的速度更快(服务器上的通道最终变得不可写)。但我无法弄清楚客户有什么问题。为什么它不能处理更多的消息...

标签: java performance netty low-latency throughput


【解决方案1】:

如果服务器发送固定大小(~100 字节)的消息,您可以将 ReceiveBufferSizePredictor 设置为客户端引导程序,这将优化读取

bootstrap.setOption("receiveBufferSizePredictorFactory",
            new AdaptiveReceiveBufferSizePredictorFactory(MIN_PACKET_SIZE, INITIAL_PACKET_SIZE, MAX_PACKET_SIZE));

根据您发布的代码段:客户端的 nio 工作线程正在处理管道中的所有内容,因此它将忙于解码和执行消息处理程序。您必须添加一个执行处理程序。

您已经说过,频道从服务器端变得不可写,因此您可能需要在服务器引导程序中调整水印大小。您可以定期监控写入缓冲区大小(写入队列大小),并确保通道由于消息无法写入网络而变得不可写入。可以通过像下面这样的 util 类来完成。

package org.jboss.netty.channel.socket.nio;

import org.jboss.netty.channel.Channel;

public final class NioChannelUtil {
  public static long getWriteTaskQueueCount(Channel channel) {
    NioSocketChannel nioChannel = (NioSocketChannel) channel;
    return nioChannel.writeBufferSize.get();
  }
}

【讨论】:

  • 杰斯坦,谢谢你的回答。我在实际业务逻辑处理程序(我测量吞吐量)之前在管道中添加了 ExecutionHandler,但它没有成功。吞吐量仍然很低,服务器通道不可写。我试图找到类 NioSocketChannel,但失败了。
  • @EgorLakomkin,所以你的问题还没有解决?,NioSocketChannel 丢失或包org.jboss.netty.channel.socket.nio 丢失? Netty 版本是什么?
  • Netty 4.x 似乎不允许访问 writeBufferSize 对象,而且包似乎也发生了变化。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2020-12-31
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多