【问题标题】:Akka streams - shutdown stream with grouping without losing dataAkka 流 - 关闭流并分组而不丢失数据
【发布时间】:2017-09-29 08:40:00
【问题描述】:

我有一个将元素分组的源和一个发出批处理请求的接收器, 我使用 KillSwitch 能够在任意时间点关闭图表。调用switch.shutdown()时源输出的最新不完整批次记录丢失的问题

val source = Source.tick(10.millis, 10.millis, "tick").grouped(500)

val (switch, _) = source.viaMat(KillSwitches.single)(Keep.right)
.toMat(sink)(Keep.both).run()

Thread.sleep(3000) // wait some arbitrary time

switch.shutdown()

有没有办法在关机发生时“清除”不完整的批次?

【问题讨论】:

    标签: scala akka-stream


    【解决方案1】:

    根据其文档,终止开关关闭的行为是位置性的

    在调用 [[UniqueKillSwitch#shutdown()]] 的运行实例后 [[FlowShape]] 的 [[Graph]] 实现为 [[UniqueKillSwitch]] 将完成其下游并取消其 上游(除非已经完成或失败,在这种情况下 命令被忽略)。

    另请参阅更多文档here

    现在grouped 阶段只会在完成时发出部分填充的组,但在取消时不会发出。

    这意味着下面的图表(分组之前 killswitch)将表现得像你观察到的那样

      val switch = 
        Source.tick(10.millis, 175.millis, "tick")
              .grouped(10)
              .viaMat(KillSwitches.single)(Keep.right)
              .toMat(Sink.foreach(println))(Keep.left)
              .run()
    

    而下图(分组 killswitch之后)将在完成时向下游发出部分组

      val switch =
        Source.tick(10.millis, 175.millis, "tick")
              .viaMat(KillSwitches.single)(Keep.right)
              .grouped(10)
              .toMat(Sink.foreach(println))(Keep.left)
              .run()
    

    【讨论】:

    • @stanislav.chetvertkov 不客气 :) 你能接受答案吗?
    • 感谢您的解释。很好的答案。
    猜你喜欢
    • 2021-05-25
    • 2021-01-24
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-10-31
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多