【问题标题】:Akka stream actor-conflation-ratelimit-actor drops first few messages (sometimes)Akka 流 actor-conflation-ratelimit-actor 丢弃前几条消息(有时)
【发布时间】:2019-08-20 14:41:19
【问题描述】:

一个简单的合并组合(如下)有时会在 staartup 打印一条调试消息,说明由于零需求它正在丢弃消息。 我希望合并阶段能够提供无限的需求,所以上述情况绝不应该是这样。我错过了什么?

val sourceRef = Source.actorRef[KeyedHighFreqEvent](0, OverflowStrategy.fail)
.conflateWithSeed(...into hash map...)
.throttle(8, per = 1.second, maxBurst=24, ThrottleMode.shaping)
.mapConcat(...back to individual KeyedHighFreqEvent...)
.groupedWithin(1024, 1.millisecond)
.to(Sink.actorRef(networkPublisher, Nil))
.run()

system.eventStream.subscribe(sourceRef, classOf[KeyedHighFreqEvent])

【问题讨论】:

    标签: scala akka akka-stream reactive-streams backpressure


    【解决方案1】:

    Source.actorRef 的文档对此非常清楚:

    可以使用bufferSize of 0 禁用缓冲区,然后如果没有需求,则丢弃接收到的消息 从下游。当 bufferSize 为 0 时,overflowStrategy 无关紧要。之后添加了一个异步边界 这个来源;因此,假设下游总是会产生需求是不安全的。

    问题在于源和合并阶段之间的异步边界。合并阶段确实提供了无限的需求,但异步边界类型使其传播到源的速度很慢。

    您可以在源中使用缓冲区(增加 bufferSize),或者在适当的情况下使用其他源,例如 Source.queue,因为它不会引入异步边界

    【讨论】:

    • 谢谢,非常有道理,我忽略了“添加了异步边界”部分。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2011-12-02
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2023-04-02
    • 2013-04-30
    • 1970-01-01
    相关资源
    最近更新 更多