【问题标题】:How can I send multiple http requests asynchronous with Netty?如何使用 Netty 异步发送多个 http 请求?
【发布时间】:2016-09-04 04:21:07
【问题描述】:

我正在尝试向一台服务器异步发送大量 http 帖子请求。我的目标是将每个响应与其原始请求进行比较。

为此,我正在关注 Netty Snoop example

但是,这个例子(和其他 http 例子)没有介绍如何异步发送多个请求,也没有介绍如何将它们随后链接到相应的请求。

所有类似的问题(如this onethis onethis one,实现SimpleChannelUpstreamHandler类,该类来自netty 3,4.0不再存在(documentation netty 4.0

有人知道如何在 netty 4.0 中解决这个问题吗?

编辑:

我的问题是虽然我向频道写了很多消息,但我收到的回复很慢(1 个响应/秒,而希望收到几千个/秒)。为了澄清这一点,让我发布到目前为止我得到的东西。我确信我发送请求的服务器也可以处理大量流量。

到目前为止我得到了什么:

import java.net.URI
import java.nio.charset.StandardCharsets
import java.io.File

import io.netty.bootstrap.Bootstrap
import io.netty.buffer.{Unpooled, ByteBuf}
import io.netty.channel.{ChannelHandlerContext, SimpleChannelInboundHandler, ChannelInitializer}
import io.netty.channel.socket.SocketChannel
import io.netty.channel.socket.nio.NioSocketChannel
import io.netty.handler.codec.http._
import io.netty.handler.timeout.IdleStateHandler
import io.netty.util.{ReferenceCountUtil, CharsetUtil}
import io.netty.channel.nio.NioEventLoopGroup

import scala.io.Source

object ClientTest {

  val URL = System.getProperty("url", MY_URL)     
  val configuration = new Configuration

  def main(args: Array[String]) {
    println("Starting client")
    start()
  }

  def start(): Unit = {

    val group = new NioEventLoopGroup()

    try {

      val uri: URI = new URI(URL)
      val host: String= {val h = uri.getHost(); if (h != null) h else "127.0.0.1"}
      val port: Int = {val p = uri.getPort; if (p != -1) p else 80}

      val b = new Bootstrap()

      b.group(group)
      .channel(classOf[NioSocketChannel])
      .handler(new HttpClientInitializer())

      val ch = b.connect(host, port).sync().channel()

      val logFolder: File = new File(configuration.LOG_FOLDER)
      val fileToProcess: Array[File] = logFolder.listFiles()

      for (file <- fileToProcess){
        val name: String = file.getName()
        val source = Source.fromFile(configuration.LOG_FOLDER + "/" + name)

        val lineIterator: Iterator[String] = source.getLines()

        while (lineIterator.hasNext) {
            val line = lineIterator.next()
            val jsonString = parseLine(line)
            val request = createRequest(jsonString, uri, host)
            ch.writeAndFlush(request)
        }
        println("closing")
        ch.closeFuture().sync()
      }
    } finally {
      group.shutdownGracefully()
    }
  }

  private def parseLine(line: String) = {
    //do some parsing to get the json string I want
  }

  def createRequest(jsonString: String, uri: URI, host: String): FullHttpRequest = {
    val bytebuf: ByteBuf = Unpooled.copiedBuffer(jsonString, StandardCharsets.UTF_8)

    val request: FullHttpRequest = new DefaultFullHttpRequest(
      HttpVersion.HTTP_1_1, HttpMethod.POST, uri.getRawPath())
    request.headers().set(HttpHeaders.Names.HOST, host)
    request.headers().set(HttpHeaders.Names.CONNECTION, HttpHeaders.Values.KEEP_ALIVE)
    request.headers().set(HttpHeaders.Names.ACCEPT_ENCODING, HttpHeaders.Values.GZIP)
    request.headers().add(HttpHeaders.Names.CONTENT_TYPE, "application/json")

    request.headers().set(HttpHeaders.Names.CONTENT_LENGTH, bytebuf.readableBytes())
    request.content().clear().writeBytes(bytebuf)

    request
  }
}

class HttpClientInitializer() extends ChannelInitializer[SocketChannel] {

  override def initChannel(ch: SocketChannel) = {
  val pipeline = ch.pipeline()

  pipeline.addLast(new HttpClientCodec())

  //aggregates all http messages into one if content is chunked
  pipeline.addLast(new HttpObjectAggregator(1048576))

  pipeline.addLast(new IdleStateHandler(0, 0, 600))

  pipeline.addLast(new HttpClientHandler())
  }
}

class HttpClientHandler extends SimpleChannelInboundHandler[HttpObject] {

  override def channelRead0(ctx: ChannelHandlerContext, msg: HttpObject) {
    try {
      msg match {
        case res: FullHttpResponse =>
          println("response is: " + res.content().toString(CharsetUtil.US_ASCII))
          ReferenceCountUtil.retain(msg)
      }
    } finally {
      ReferenceCountUtil.release(msg)
    }
  }

  override def exceptionCaught(ctx: ChannelHandlerContext, e: Throwable) = {
    println("HttpHandler caught exception", e)
    ctx.close()
  }
}

【问题讨论】:

  • 写入通道不是异步的吗?作为 write 的结果,你会得到 Future,这取决于你如何处理它
  • 我也在学习 Netty 4.0。这是我对设计的理解。我要记住的第一件事是,在 Netty 4 中,您确信所有注册的处理程序都在单线程中执行,因此不需要同步,除非您使用共享处理程序。因此,您提交的所有请求都将通过该通道按顺序发送,并且将以相同的顺序接收响应。因此,在双工处理程序中为所有请求管理数据结构(如队列),您始终可以轮询相应请求以获取最新收到的响应。
  • 感谢您的回复!我的问题是虽然我向频道写了很多消息,但我收到的回复很慢(1 个回复/秒,而希望收到几千个/秒)。为了澄清这一点,让我发布我到目前为止所获得的信息。
  • 您能否使用更多线程来扩展您的事件循环组,并检查响应流的性能是否有所提高?

标签: java scala asynchronous netty nio


【解决方案1】:

ChannelFuture cf = channel.writeAndFlush(createRequest());

也不知道如何将它们随后链接到相应的请求。

Can netty assign multiple IO threads to the same Channel?

一旦分配给通道的工作线程在通道的生命周期内不会改变。所以我们没有从线程中受益。这是因为您保持连接处于活动状态,通道也保持活动状态。

要解决此问题,您可以考虑使用一个频道池(例如 30 个)。然后使用通道池来放置您的请求。

      int concurrent = 30;

  // Start the client.
  ChannelFuture[] channels = new ChannelFuture[concurrent];
  for (int i = 0; i < channels.length; i++) {
    channels[i] = b.connect(host, port).sync();
  }

  for (int i = 0; i < 1000; i++) {
      ChannelFuture requestHandle = process(channels[(i+1)%concurrent]); 
      // do something with the request handle       
  }

  for (int i = 0; i < channels.length; i++) {
    channels[i].channel().closeFuture().sync();
  }

HTH

【讨论】:

  • 我认为即使使用一个通道,他也应该从服务器端获得更多每秒 1 条消息。我的假设是作者堆积了频道,请求没有机会处理响应。
  • 并发 - 处理 1000 个请求的时间(毫秒)(因端点而异) 1 - 55219 10 - 48749 30 - 13364 100 - 29106 直到线程数增加的时间会提高性能,但稍后的上下文切换会影响性能。
  • user1582639 你提到的确实是问题之一。限制请求的数量可以大大提高性能。 @AmodPandey:使用频道池也可以,并进一步提高性能。但是,我仍然不明白如何将初始请求链接到相应的响应。 writeAndFlush 确实返回了一个 ChannelFurture,但是如果我添加一个侦听器,它将在请求发送后完成,这使我无法以某种方式将其连接到响应。
  • Mart - 凭借其设计,在 ClientTest 中我们只能确定请求是否成功放置到端点上。 HttpClientHandler 将处理来自服务器的任何响应(它不会知道相应的请求)。但是如果你真的想要请求和响应之间的相关性,那么有办法做到这一点。 1. 在通道的 ChannelOutboundHandler 中设置属性,您可以在 HttpClientHandler 中检索该属性。 2. 拥有返回请求中发送的请求 id 的服务器。另一方面,您可能会很好地使用 Apache Sync Http Client
  • 也看看github.com/netty/netty/blob/4.1/example/src/main/java/io/netty/…。如果您想查看上述建议的示例代码,请告诉我。
猜你喜欢
  • 1970-01-01
  • 2012-04-07
  • 2011-01-08
  • 2015-06-09
  • 2018-12-16
  • 2016-11-15
  • 2018-04-10
  • 2012-12-20
  • 2015-08-12
相关资源
最近更新 更多