【问题标题】:Akka Stream: Can not write to file sinkAkka Stream:无法写入文件接收器
【发布时间】:2016-09-19 20:43:44
【问题描述】:

我正在尝试运行一个简单的 Akka Stream File Sink 示例,但没有成功。我可以创建一个 Source,运行 Flow,然后创建一个文件,但 ByteString 没有被写入文件。而如果我尝试将流输出打印到控制台,我可以这样做。我在这里遗漏了什么吗?

import akka.stream._ 
import akka.stream.scaladsl._
import akka.{ NotUsed, Done}
import akka.actor.ActorSystem
import akka.util.ByteString
import scala.concurrent._
import scala.concurrent.duration._
import java.nio.file.Paths

object First extends App {

  val source: Source[Int, NotUsed] = Source ( 1 to 100)

  implicit val system = ActorSystem("QuickStart")
  implicit val materializer = ActorMaterializer()

  // works: prints 1-100
  //source.runForeach(println) (materializer)

  val factorials = source.scan(BigInt(1))((acc,next) => acc * next)

  // there is no content in the Sink (file)
  /**val result =
    factorials
    .map(num => ByteString(s"${num}\n"))
    .runWith(FileIO.toPath(Paths.get("factorials.txt")))
**/

  def lineSink(fileName: String): Sink[String, Future[IOResult]] =
    Flow[String]
    .map(s => ByteString(s + "\n"))
    .toMat(FileIO.toPath(Paths.get(fileName))) (Keep.right)

  //There is no content in the Sink.
  factorials.map(_.toString).runWith(lineSink("factorials.txt"))

system.terminate()

}

build.sbt 有:

name := "akkaGuide"
    version := "1.0"
    scalaVersion := "2.11.8"
    libraryDependencies ++= Seq(
      "com.typesafe.akka" %% "akka-stream" % "2.4.10"
    )

提前感谢您的宝贵时间。

【问题讨论】:

    标签: scala akka akka-stream


    【解决方案1】:

    我认为您可能只是过早地终止了。尝试等待 Future 完成:

    val result = factorials.map(_.toString).runWith(lineSink("factorials.txt"))
    import system.dispatcher
    result.onComplete { _ => system.terminate() }
    

    【讨论】:

      【解决方案2】:

      看看下面这个工作示例:

      package ru.io
      
      import java.io.File
      
      import akka.actor.ActorSystem
      import akka.stream.scaladsl._
      import akka.stream.{ActorMaterializer, ClosedShape}
      import akka.util.ByteString
      
      import scala.util.{Failure, Success}
      
      object WriteStreamApp extends App {
        implicit val actorSystem      = ActorSystem()
        implicit val flowMaterializer = ActorMaterializer()
        import actorSystem.dispatcher
      
        // Source
        val source = Source(1 to 10000).filter(isPrime)
      
        // Sink
        val sink = FileIO.toFile(new File("src/main/resources/prime.txt"))
      
        // output for file
        val fileSink = Flow[Int]
          .map(i => ByteString(i.toString))
          .toMat(sink)((_, bytesWritten) => bytesWritten)
      
        val consoleSink = Sink.foreach[Int](println)
      
        // using Graph API send the integers to both skins: file and console
        val graph = GraphDSL.create(fileSink, consoleSink)((file, _) => file) { implicit builder => (file, console) =>
          import GraphDSL.Implicits._
      
          val broadCast = builder.add(Broadcast[Int](2))
      
          source ~> broadCast ~> file
          broadCast ~> console
          ClosedShape
        }
      
        val materialized = RunnableGraph.fromGraph(graph).run()
      
        // make sure the system is terminated
        materialized.onComplete {
          case Success(_) =>
            actorSystem.terminate()
          case Failure(e) =>
            println(s"Failure: ${e.getMessage}")
            actorSystem.terminate()
        }
      
        def isPrime(n: Int): Boolean = {
          if (n <= 1) false
          else if (n == 2) true
          else !(2 until n).exists(x => n % x == 0)
        }
      }
      

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 2016-09-05
        • 2018-02-07
        • 1970-01-01
        • 1970-01-01
        • 2021-03-14
        • 1970-01-01
        • 1970-01-01
        相关资源
        最近更新 更多