【问题标题】:Number of threads in Akka keep increasing. What could be wrong?Akka 中的线程数不断增加。有什么问题?
【发布时间】:2016-04-22 09:03:31
【问题描述】:

为什么线程数一直在增加?

看这张图片的右下角。

整体流程是这样的:

Akka HTTP Server API 
  -> on http request, sendMessageTo DataProcessingActor
       -> sendMessageTo StorageActor
          -> sendMessageTo DataBaseActor 
          -> sendMessageTo IndexActor 

这是 Akka HTTP API 的定义(伪代码):

Main {
  path("input/") { 
    post {
       dataProcessingActor forward message
    }
  }
}

以下是参与者定义(伪代码):

DataProcessingActor {
  case message => 
    message = parse message
    storageActor ! message
}


StorageActor {
  case message => 
    indexActor ! message
    databaseActor ! message
}


DataBaseActor {
  case message =>
    val c = get monogCollection
    c.store(message)
}

IndexActor {
  case message =>
    elasticSearch.index(message)
}

运行此设置后,向“input/”HTTP 端点发送多个 HTTP 请求时,我收到错误:

for( i <- 0 until 1000000) {
   post("input/", someMessage+i)
}

错误:

[ERROR] [04/22/2016 13:20:54.016] [Main-akka.actor.default-dispatcher-15] [akka.tcp://Main@127.0.0.1:2558/system/IO-TCP/selectors/$a/0] Accept error: could not accept new connection
java.io.IOException: Too many open files
    at sun.nio.ch.ServerSocketChannelImpl.accept0(Native Method)
    at sun.nio.ch.ServerSocketChannelImpl.accept(ServerSocketChannelImpl.java:422)
    at sun.nio.ch.ServerSocketChannelImpl.accept(ServerSocketChannelImpl.java:250)
    at akka.io.TcpListener.acceptAllPending(TcpListener.scala:107)
    at akka.io.TcpListener$$anonfun$bound$1.applyOrElse(TcpListener.scala:82)
    at akka.actor.Actor$class.aroundReceive(Actor.scala:480)
    at akka.io.TcpListener.aroundReceive(TcpListener.scala:32)
    at akka.actor.ActorCell.receiveMessage(ActorCell.scala:526)
    at akka.actor.ActorCell.invoke(ActorCell.scala:495)
    at akka.dispatch.Mailbox.processMailbox(Mailbox.scala:257)
    at akka.dispatch.Mailbox.run(Mailbox.scala:224)
    at akka.dispatch.Mailbox.exec(Mailbox.scala:234)
    at scala.concurrent.forkjoin.ForkJoinTask.doExec(ForkJoinTask.java:260)
    at scala.concurrent.forkjoin.ForkJoinPool$WorkQueue.runTask(ForkJoinPool.java:1339)
    at scala.concurrent.forkjoin.ForkJoinPool.runWorker(ForkJoinPool.java:1979)
    at scala.concurrent.forkjoin.ForkJoinWorkerThread.run(ForkJoinWorkerThread.java:107)

编辑 1

这是正在使用的application.conf 文件:

akka {
  loglevel = "INFO"
  stdout-loglevel = "INFO"
  logging-filter = "akka.event.slf4j.Slf4jLoggingFilter"

  actor {
    default-dispatcher {
      throughput = 10
    }
  }

  actor {
    provider = "akka.remote.RemoteActorRefProvider"
  }

  remote {
    enabled-transports = ["akka.remote.netty.tcp"]
    netty.tcp {
      hostname = "127.0.0.1"
      port = 2558
    }
  }
}

【问题讨论】:

  • 可能是mongo驱动
  • 只有一个 mongo 数据库连接被使用,它在一个 scala 对象object DB { lazy val db = mongodb.connect() } 中。基本上,在演员内部并没有建立新的联系。仅在第一次初始化 DB 对象时建立新连接,并引用 db。它认为它不可能是 MongoDB 连接。我错过了什么吗?
  • 我会附加一个调试器/分析器,看看这数千个线程在做什么。
  • 刚刚发现 ElasticSearch 是问题所在。我正在为 ElasticSearch 使用 Java API,这就是泄漏套接字。让我解决这个问题并重新检查。
  • 是的,这是 ElasticSearch 的问题。现在解决了。

标签: multithreading mongodb scala elasticsearch akka


【解决方案1】:

我发现 ElasticSearch 是问题所在。我正在为 ElasticSearch 使用 Java API,因为它是从 Java API 中使用的,所以它正在泄漏套接字。现在已按照此处所述解决。

以下是使用 Java API 的 Elastic Search 客户端服务

trait ESClient { def getClient(): Client }

case class ElasticSearchService() extends ESClient {
  def getClient(): Client = {
    val client = new TransportClient().addTransportAddress(
      new InetSocketTransportAddress(Config.ES_HOST, Config.ES_PORT)
    )
    client
  }
}

这是导致泄漏的演员:

class IndexerActor() extends Actor {

  val elasticSearchSvc = new ElasticSearchService()
  lazy val client = elasticSearchSvc.getClient()

  override def preStart = {
    // initialize index, and mappings etc.
  }

  def receive() = {
    case message => 
      // do indexing here
      indexMessage(ES.client, message)
  }
}

注意:每次创建 actor 实例时,都会建立一个新连接。

new ElasticSearchService() 的每次调用都在创建与 ElasticSearch 的新连接。我将它移到了一个单独的对象中,如下所示,并且演员也使用了这个对象:

object ES {
  val elasticSearchSvc = new ElasticSearchService()
  lazy val client = elasticSearchSvc.getClient()
}


class IndexerActor() extends Actor {

  override def preStart = {
    // initialize index, and mappings etc.
  }

  def receive() = {
    case message => 
      // do indexing here
      indexMessage(ES.client, message)
  }
}

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2014-12-10
    • 2016-06-17
    • 1970-01-01
    • 2021-07-25
    • 1970-01-01
    • 2012-06-20
    • 1970-01-01
    相关资源
    最近更新 更多