【问题标题】:Asynchronous update of promise in Netty NioNetty Nio 中 Promise 的异步更新
【发布时间】:2018-04-11 11:03:41
【问题描述】:

我有一个交换信息的服务器和客户端架构。我想从服务器返回连接通道的数量。我想使用 promise 将服务器的消息返回给客户端。我的代码是:

public static void callBack () throws Exception{

   String host = "localhost";
   int port = 8080;

   try {
       Bootstrap b = new Bootstrap();
       b.group(workerGroup);
       b.channel(NioSocketChannel.class);
       b.option(ChannelOption.SO_KEEPALIVE, true);
       b.handler(new ChannelInitializer<SocketChannel>() {
        @Override
           public void initChannel(SocketChannel ch) throws Exception {
            ch.pipeline().addLast(new RequestDataEncoder(), new ResponseDataDecoder(), new ClientHandler(promise));
           }
       });
       ChannelFuture f = b.connect(host, port).sync();
       //f.channel().closeFuture().sync();
   }
   finally {
    //workerGroup.shutdownGracefully();
   }
}

public static void main(String[] args) throws Exception {

  callBack();
  while (true) {

    Object msg = promise.get();
    System.out.println("The number if the connected clients is not two");
    int ret = Integer.parseInt(msg.toString());
    if (ret == 2){
        break;
    }
  }
  System.out.println("The number if the connected clients is two");
}

当我运行一个客户端时,它总是收到消息The number if the connected clients is not two,并且返回的数字总是一。当我运行第二个客户端时,它总是接收一个返回值,但是,第一个客户端仍然接收一个。对于第一个客户的情况,我找不到更新承诺的正确方法。

编辑: 客户端服务器:

public class ClientHandler extends ChannelInboundHandlerAdapter {
  public final Promise<Object> promise;
  public ClientHandler(Promise<Object> promise) {
      this.promise = promise;
  }

  @Override
  public void channelActive(ChannelHandlerContext ctx) throws Exception {
      RequestData msg = new RequestData();
      msg.setIntValue(123);
      msg.setStringValue("all work and no play makes jack a dull boy");
      ctx.writeAndFlush(msg);
  }

  @Override
  public void channelRead(ChannelHandlerContext ctx, Object msg) throws Exception {
      System.out.println(msg);
      promise.trySuccess(msg);
  }
} 

来自客户端处理程序的代码,用于存储从服务器接收到的消息到 Promise。

【问题讨论】:

  • 当你说promise时,你的意思是非阻塞吗?
  • 我的意思是这个Object msg = promise.get();,这个值有服务器的retuned消息。
  • 你可以关注我对不同问题的回答stackoverflow.com/questions/46852221/…
  • @konstantin 根据您的问题,您的意思是:promise.get() 返回连接到您服务器的客户端(通道)数量??
  • 在客户端处理程序中存储我从客户端读取的消息。 github.com/kristosh/netty-nio-Client-Server

标签: java sockets web netty nio


【解决方案1】:

在 Netty 框架中,PromiseFuture 是一次写入对象,这一原则使它们更易于在多线程环境中使用。

由于 Promise 没有做你想做的事,我们需要看看其他技术是否适合你的条件,你的条件基本上归结为:

  • 从多个线程读取
  • 仅从单个线程写入(因为在 Netty 通道内,读取方法只能由 1 个线程同时执行,除非通道被标记为可共享)

对于这些要求,最合适的匹配是 volatile 变量,因为它对于读取是线程安全的,并且可以由 1 个线程安全地更新,而无需担心写入顺序。

要更新您的代码以使用 volatile 变量,需要进行一些修改,因为我们无法轻松地将引用链接传递给函数内的变量,但我们必须传递一个更新后端变量的函数。

private static volatile int connectedClients = 0;
public static void callBack () throws Exception{
    //....
           ch.pipeline().addLast(new RequestDataEncoder(), new ResponseDataDecoder(),
                                 new ClientHandler(i -> {connectedClients = i;});
    //....
}

public static void main(String[] args) throws Exception {

  callBack();
  while (true) {
    System.out.println("The number if the connected clients is not two");
    int ret = connectedClients;
    if (ret == 2){
        break;
    }
  }
  System.out.println("The number if the connected clients is two");
}

public class ClientHandler extends ChannelInboundHandlerAdapter {
  public final IntConsumer update;
  public ClientHandler(IntConsumer update) {
      this.update = update;
  }

  @Override
  public void channelActive(ChannelHandlerContext ctx) throws Exception {
      RequestData msg = new RequestData();
      msg.setIntValue(123);
      msg.setStringValue("all work and no play makes jack a dull boy");
      ctx.writeAndFlush(msg);
  }

  @Override
  public void channelRead(ChannelHandlerContext ctx, Object msg) throws Exception {
      System.out.println(msg);
      update.accept(Integer.parseInt(msg));
  }
} 

虽然上面的方法应该可行,但我们很快发现主类中的 while 循环占用了很大一部分 CPU 时间,这可能会影响本地客户端系统的其他部分,幸运的是,如果我们将其他部分添加到系统中,即同步。通过将connectedClients 的初始读取留在同步块之外,我们仍然可以在“真”情况下从快速读取中获益,而在“假”情况下,我们可以保护重要的 CPU 周期可用于系统的其他部分。

为了解决这个问题,我们在阅读时采用以下步骤:

  1. connectedClients 的值存储在单独的变量中
  2. 将此变量与目标值进行比较
  3. 如果是真的,那么尽早跳出循环
  4. 如果为 false,则进入同步块内
  5. 开始一个 while true 循环
  6. 再次读出变量,因为值现在可能会改变
  7. 检查条件,如果现在条件正确则中断
  8. 如果没有,请等待值发生变化

写作时还有以下几点:

  1. 同步
  2. 更新值
  3. 唤醒等待该值的所有其他线程

这可以在如下代码中实现:

private static volatile int connectedClients = 0;
private static final Object lock = new Object();
public static void callBack () throws Exception{
    //....
           ch.pipeline().addLast(new RequestDataEncoder(), new ResponseDataDecoder(),
                                 new ClientHandler(i -> {
               synchronized (lock) {
                   connectedClients = i;
                   lock.notifyAll();
               }
           });
    //....
}

public static void main(String[] args) throws Exception {

  callBack();
  int connected = connectedClients;
  if (connected != 2) {
      System.out.println("The number if the connected clients is not two before locking");
      synchronized (lock) {
          while (true) {
              connected = connectedClients;
              if (connected == 2)
                  break;
              System.out.println("The number if the connected clients is not two");
              lock.wait();
          }
      }
  }
  System.out.println("The number if the connected clients is two: " + connected );
}

服务器端更改

但是,并非所有问题都与客户端有关。

因为您发布了指向您的 github 存储库的链接,所以当有新人加入时,您永远不会从服务器向旧客户端发送请求。因为这没有完成,所以永远不会通知客户有关更改,请确保也这样做。

【讨论】:

  • 当我尝试将 connectedClients 添加到 initChannels 时(new ClientHandler(i -> {connectedClients = i;}) 正在接收不兼容的类型,它在接收 lambda 参数时需要 int。
  • @konstantin 你有没有注意到我调整了ClientHandler,我把参数从Promise 改成了IntConsumer。如果您有一个类型为 IntConsumer 的参数,您可以使用一个接受 1 个 int 作为输入的 lamba 块,这是我用来将结果传播到其他类的方法。
  • 好的,我正确地编写了代码。我运行了两个客户端,但是在它们两个中我都收到了以下消息:2017 年 11 月 2 日上午 11:43:09 io.netty.channel.DefaultChannelPipeline onUnhandledInboundException 警告:一个 exceptionCaught() 事件被触发,它到达了尾部的管道。这通常意味着管道中的最后一个处理程序没有处理异常。 java.lang.ClassCastException:chat.ResponseData 无法转换为 java.lang.Integer
  • 更新的代码可以在这里找到github.com/kristosh/netty-nio-Client-Server
  • 这可能是我从您的原始代码中遗漏的内容,在CLientHandler 中,而不是update.accept(Integer.parseInt(msg));,执行update.accept(((RequestData)msg).getIntValue());,这会直接从您的请求数据对象中提取大小
猜你喜欢
  • 1970-01-01
  • 2018-04-01
  • 1970-01-01
  • 2010-10-25
  • 1970-01-01
  • 1970-01-01
  • 2011-02-23
  • 2017-08-24
  • 2013-12-29
相关资源
最近更新 更多