【问题标题】:Play + ReactiveMongo: capped collection and tailable cursorPlay + ReactiveMongo:上限集合和可尾游标
【发布时间】:2015-07-30 14:53:36
【问题描述】:

我正在使用带有 Scala、Akka 和 ReactiveMongo 的 Play Framework。我想将 MongoDB 中的集合用作循环队列。多个参与者可以将文档插入其中;一个参与者在这些文档可用时立即检索它们(一种发布-订阅系统)。 我正在使用上限集合和可尾光标。每次我检索一些文档时,我都必须运行命令 EmptyCapped 来刷新上限集合(不可能从中删除元素),否则我总是检索相同的文档。有替代解决方案吗?例如,有没有办法在不删除元素的情况下滑动光标?或者在我的情况下最好不要使用上限集合?

object MexDB {

def db: reactivemongo.api.DB = ReactiveMongoPlugin.db
val size: Int = 10000

// creating capped collection
val collection: JSONCollection = {

    val c = db.collection[JSONCollection]("messages")

    val isCapped = coll.convertToCapped(size, None)

    Await.ready(isCapped, Duration.Inf)

    c
}

def insert(mex: Mex) = {

    val inserted = collection.insert(mex)

    inserted onComplete {
      case Failure(e) =>
        Logger.info("Error while inserting task: " + e.getMessage())
        throw e

      case Success(i) =>
        Logger.info("Successfully inserted task")
    }

}


def find(): Enumerator[Mex] = {

  val cursor: Cursor[Mex] = collection
    .find(Json.obj())
    .options(QueryOpts().tailable.awaitData)
    .cursor[Mex]

    // meaning of maxDocs ???
    val maxDocs = 1
    cursor.enumerate(maxDocs)
}


def removeAll() = {
    db.command(new EmptyCapped("messages"))
}

}

/*** part of receiver actor code ***/

// inside preStart
val it = Iteratee.fold[Mex, List[Mex]](Nil) {
    (partialList, mex) => partialList ::: List(mex)
}

// Inside "receive" method
case Data =>

  val e: Enumerator[Mex] = MexDB.find()

  val future = e.run(it)

  future onComplete {
    case Success(list) =>
      list foreach { mex =>
        Logger.info("Mex: " + mex.id)
      }
      MexDB.removeAll()
      self ! Data

    case Failure(e) => Logger.info("Error:  "+ e.getMessage())
  }

【问题讨论】:

    标签: mongodb scala playframework reactivemongo capped-collections


    【解决方案1】:

    您的可尾光标在每个找到的文档后关闭为maxDocs = 1。要使其无限期保持打开状态,您应该省略此限制。

    使用awaitData.onComplete 只有在您明确关闭 RM 时才会被调用。

    您需要从光标使用一些流式处理函数,例如.enumerate 并处理每个新的步骤/结果。见https://github.com/sgodbillon/reactivemongo-tailablecursor-demo/

    【讨论】:

    • 谢谢。因此,如果我省略该限制,我可以只创建一次枚举器并重新使用它,对吗?有没有另一种方法可以像循环队列一样从有上限的集合中检索文档,而无需每次都运行 EmptyCapped?
    • Max 1 只读取一个文档。删除限制,光标在其当前“位置”保持打开状态,而不是回到第一个文档。
    • 我删除了限制,我只创建了一次光标和枚举器。我定期执行: val future = enum.run(it) future onComplete { *** } 我在集合中定期插入 mex 但 enum 从未检索到一些 mex ...我做错了什么?
    • 我需要将 Input.Done 注入我的枚举器吗?
    • 对于awaitDataonComplete 只会在您明确关闭 RM 时被调用。您需要使用光标中的一些流功能,例如enumerate 并处理每个新步骤/结果。见github.com/sgodbillon/reactivemongo-tailablecursor-demo/blob/…
    猜你喜欢
    • 2012-08-18
    • 2014-01-09
    • 2012-09-17
    • 2020-08-10
    • 1970-01-01
    • 2019-06-11
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多