【问题标题】:How can I split data from Slick 3.1(Scala) into 4 parts如何将 Slick 3.1(Scala) 中的数据拆分为 4 个部分
【发布时间】:2017-01-30 18:49:08
【问题描述】:

为了使 lucene 索引 (v6.1) 更快,我想将 Slick 3.1(Scala) 中的数据拆分为多个部分(块),以便在线程中传递不同的数据集以加快索引过程。我在 Scala 中编写了以下代码来从 MySQL 获取数据。

class NotesService(val databaseService: DatabaseService)(implicit executionContext: ExecutionContext) extends NoteEntityTable {    
  import databaseService._
  import databaseService.driver.api._
  import com.github.t3hnar.bcrypt._    
  def getNotes(): Future[Seq[NoteEntity]] = db.run(notes.result)    
}
case class NoteEntity(id: Option[Long] = None, title: String, teaser: String, description: String)

NotesService 代码

class NotesService(val databaseService: DatabaseService)(implicit executionContext: ExecutionContext) extends NoteEntityTable {

  import databaseService._
  import databaseService.driver.api._
  import com.github.t3hnar.bcrypt._

  def getNotes(): Future[Seq[NoteEntity]] = db.run(notes.result)

}

从我使用过的 NotesService 中获取数据:

def setI = {
    val NUM_THREADS = Runtime.getRuntime().availableProcessors()
    val IndexStoreDir = Paths.get("/var/www/html/Index")
    val analyzer = new StandardAnalyzer()
    val writerConfig = new IndexWriterConfig(analyzer)
    writerConfig.setOpenMode(OpenMode.CREATE_OR_APPEND)
    writerConfig.setRAMBufferSizeMB(500)
      .setMaxBufferedDocs(2)
      .setMergeScheduler(new ConcurrentMergeScheduler())
    val directory = FSDirectory.open(IndexStoreDir)
    var writer = new IndexWriter(directory, writerConfig)

    val threads = Array.ofDim[IndexTh](NUM_THREADS)
    val notes = notesService.getNotes()

    for (i <- 0 until NUM_THREADS){

      threads(i) = new IndexTh(notesService, writer)
      //here on this line I want to pass different sets of data to thread.
    }
    for (i <- 0 until NUM_THREADS) {
      threads(i).start()
      println("Thread " + i + " Started!")
    }
  }

这里在线:

threads(i) = new IndexTh(notesService, writer)

如何从 notesService 拆分数据以传递给线程? 如何将笔记中的数据拆分为多个块? 我想要这样的数据:

假设 notesService.getNotes() 检索 20000 行数据。现在我想把这些行分成4000行的5个部分,这样每4000行数据可以传递给不同的线程。

【问题讨论】:

  • 块是什么意思?分页?
  • 我想在多个线程中传递不同的数据集(拆分主数据集)
  • 嗯,我仍然不完全理解问题是什么。也许您需要的输出示例会有所帮助。
  • notes.map(_.grouped(4)) 应该可以工作。不是吗?
  • 我已经编辑了这个问题以便更好地理解。你能复习一下吗?

标签: mysql scala lucene slick-3.0


【解决方案1】:

研究了很久终于找到了答案:

使用线程:

def setI = {
    val NUM_THREADS = Runtime.getRuntime().availableProcessors()
    val curNotes = notesService.getNotes()

    val totalRows = Await.result(curNotes, Duration.Inf).length
    var totalPages =  totalRows / NUM_THREADS
    if(totalPages != totalPages.toInt){
      totalPages = totalPages + 1
    }
    var tmp = Await.result(curNotes, Duration.Inf).grouped(totalPages).toList
    val rows = tmp(tmp.length-2) ++ tmp.last
    val threads = Array.ofDim[Index](NUM_THREADS)

    val IndexStoreDir = Paths.get("/var/www/html/LuceneIndex")
    val analyzer = new StandardAnalyzer()
    val writerConfig = new IndexWriterConfig(analyzer)
    writerConfig.setOpenMode(OpenMode.CREATE_OR_APPEND)
    writerConfig.setRAMBufferSizeMB(500)
      .setMaxBufferedDocs(10)
      .setMergeScheduler(new ConcurrentMergeScheduler())
    val directory = FSDirectory.open(IndexStoreDir)
    val writer = new IndexWriter(directory, writerConfig)
    var count = 0

    for(i <- 0 until tmp.length - 2){
      count = i
      threads(i) = new Index(tmp(i), writer, i)
    }
    count = count + 1
    threads(count) = new Index(rows, writer, count)

    for (i <- 0 until NUM_THREADS) {
      println("Thread :" + threads(i).getName + " => " + (i + 1) + " Started!")
      threads(i).start()
    }
  }

使用 Scala 未来:

def setFutureIndex = {
    val IndexStoreDir = Paths.get("/var/www/html/LuceneIndex")
    val analyzer = new StandardAnalyzer()
    val writerConfig = new IndexWriterConfig(analyzer)
    writerConfig.setOpenMode(OpenMode.CREATE)
    writerConfig.setRAMBufferSizeMB(500)
    val directory = FSDirectory.open(IndexStoreDir)
    val writer = new IndexWriter(directory, writerConfig)
    val notes = notesService.getNotes() //Gets all notes from slick. Data is coming in getNotes()
    var doc = new Document()

    def indexingFuture = {
      val list = Seq (
        notes.map(_.foreach {
          case (note) =>
            writeToDoc(note, writer)
        })
      )
      Future.sequence(list)
    }

    Await.result(indexingFuture, Duration.Inf)

    /*indexingFuture.onComplete {
      case Success(value) => println(value)
      case Failure(e) => e.printStackTrace()
    }*/
  }

  def writeToDoc(note: NoteEntity, writer: IndexWriter) = Future {
    println("*****Indexing: " + note.id.get)
    var doc = new Document()
    var field = new TextField("title", " {##" + note.id.get + "##} " + note.title, Field.Store.YES)
    doc.add(field)

    field = new TextField("teaser", note.teaser, Field.Store.YES)
    doc.add(field)

    field = new TextField("description", note.description, Field.Store.YES)
    doc.add(field)

    writer.addDocument(doc)

    writer.commit()
    println("*****Completed: " + note.id.get)
    var status = "*****Completed: " + note.id.get
  }

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2015-12-19
    • 1970-01-01
    • 2013-09-29
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2010-12-05
    • 1970-01-01
    相关资源
    最近更新 更多