【问题标题】:Integrating an ack-based Actor with akka-stream将基于 ack 的 Actor 与 akka-stream 集成
【发布时间】:2015-05-27 10:26:51
【问题描述】:

我有一个被设计用于 akka-io acking 的 Actor, 这样它会在向上游发送消息时等待 Ack(到 网络)。这个actor是异步应用程序的接口 后端。

我想要一个允许我转换它的包装层 演员进入 akka-streams Flow[Incoming, Outgoing, ???] 以便 它可以与期望这样的新库集成 签名。

(来自上游的传入消息很少,所以我们也不在乎 很多关于那里的背压,但这并不是一件坏事 拥有它。)

sealed trait Incoming //... with implementations
sealed trait Outgoing //... with implementations
object Ack

// `upstream` is an akka-io connection actor that will send Ack
// when it writes an Outgoing message to the socket
class SimpleActor(upstream: Actor) extends Actor {
  def receive = {
    case in: Incoming if sender() == upstream =>
       // does some work in response to upstream
    case other =>
       // does some work in response to downstream
       // including sending messages to upstream and
       // `becoming` a stashing state waiting for Ack
       // to `unbecome`, then sending Ack downstream
       // (which will respect the backpressure).
  }
}

我从 akka-user 邮件列表中获得了良好的授权 akka-streams 中没有代码将参与者与 流,并且为了将 Actor 插入 Stream 并保留 基于 Ack 的背压,必须实现 PushPullStage.

看来我们实际上需要两个 PushPullStages... 一个 upstream => SimpleActorSimpleActor => upstream 一个。

我的问题是:

  1. 是否有任何库提供诸如此类的演员和流之间的集成?
  2. 有没有比从头开始实现双向PushPullStage 更简单的方法?
  3. 是否有任何现有的测试框架可以对此类实现进行压力测试?

【问题讨论】:

  • 你试过用actor ask和mapAsync吗?如果没有,那么写PushPullStage 仍然比写ActorProcessor 容易得多。
  • @jrudolph 但我没有对查询的单一响应(如果是这种情况,我会使用 REST 而不是 WebSockets)。发布者可以随时发送消息。
  • 我猜“随时”是指背压允许的时候?这究竟是如何工作的?如果你为每个输入元素生成一个输出流,你可以用这种方式对其进行建模,然后将流变平。
  • 是的,当背压允许时。它不是查询/响应协议,客户端和服务器可以随时相互发送消息(网络允许)。
  • 啊,所以输入和输出通道是完全分开的,背压也必须分开处理?

标签: akka reactive-programming akka-stream


【解决方案1】:

我认为 akka-stream 的理念是提供低级积木并在其之上构建高级工具。 如果您查看我们最近发布的开源库https://github.com/MfgLabs/akka-stream-extensions,您会发现我们确实做到了这一点。我们提供了一些有用的结构,使管理速率限制器、状态处理器、惰性和生成器等变得更容易...... 对于actor集成,我认为应该可以创建某种帮助器,以便更容易地将actor与akka-stream集成,试图传播背压。 Akka-Stream 还很年轻,生态系统还在不断发展;)

【讨论】:

    【解决方案2】:

    是的,您可以将演员与流集成。
    为此目的有一些特殊的演员:演员发布者和演员订阅者。

    都在这里:http://doc.akka.io/docs/akka-stream-and-http-experimental/1.0-RC3/scala/stream-integrations.html

    当然,您必须以这样一种方式编写actor,它可以与流背压一起工作。但是你不需要推拉阶段。

    【讨论】:

    • 这听起来确实比引入新阶段并实施其方法要容易得多。我希望避免编写一个中间角色,但看起来这是唯一的出路——同时实现发布者和订阅者。这将是 很多 样板文件:-/
    猜你喜欢
    • 1970-01-01
    • 2019-12-19
    • 1970-01-01
    • 2018-01-14
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2023-04-02
    • 2016-02-06
    相关资源
    最近更新 更多