【问题标题】:Akka Java File IO ThrottlingAkka Java 文件 IO 节流
【发布时间】:2018-03-14 19:06:56
【问题描述】:

我想将文件内容逐行流式传输到 Actor。我有类似的东西:

final ActorSystem system = ActorSystem.create("stream_system");
final Materializer materializer = ActorMaterializer.create(system);
final ActorRef actor = system.actorOf(Props.create(streamActor.class), "sink");

final Path file = Paths.get("path/file.txt");

Sink<ByteString, CompletionStage<Done>> printlnSink =
        Sink.<ByteString> foreach(chunk -> actor.tell(chunk.utf8String(), null));
        //Sink.<ByteString> actorRef(actor, null);

CompletionStage<IOResult> ioResult =
        FileIO.fromPath(file)
                .throttle(1, Duration.create(1, TimeUnit.SECONDS), 1, ThrottleMode.shaping())
                .to(printlnSink)
                .run(materializer);

未注释的版本有效,但它可以一次性流式传输整个文件内容。评论版本以“未知”消息结尾。

我想以几秒钟的延迟逐行发送给 Actor。任何帮助如何做到这一点?接收演员只需获取字符串消息并将其打印在输出上。

【问题讨论】:

    标签: java akka akka-stream


    【解决方案1】:

    Framing 类可以帮助您:

    CompletionStage<IOResult> ioResult =
        FileIO.fromPath(file)
              .via(Framing.delimiter(ByteString.fromString(System.lineSeparator()), 1000, FramingTruncation.ALLOW))
              .throttle(1, Duration.create(1, TimeUnit.SECONDS), 1, ThrottleMode.shaping())
              .to(printlnSink)
              .run(materializer);
    

    【讨论】:

      猜你喜欢
      • 2017-12-18
      • 1970-01-01
      • 2012-11-01
      • 2016-10-14
      • 1970-01-01
      • 2023-03-21
      • 1970-01-01
      • 1970-01-01
      • 2016-02-25
      相关资源
      最近更新 更多