【发布时间】: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 => SimpleActor 和 SimpleActor => upstream 一个。
我的问题是:
- 是否有任何库提供诸如此类的演员和流之间的集成?
- 有没有比从头开始实现双向
PushPullStage更简单的方法? - 是否有任何现有的测试框架可以对此类实现进行压力测试?
【问题讨论】:
-
你试过用actor ask和
mapAsync吗?如果没有,那么写PushPullStage仍然比写ActorProcessor容易得多。 -
@jrudolph 但我没有对查询的单一响应(如果是这种情况,我会使用 REST 而不是 WebSockets)。发布者可以随时发送消息。
-
我猜“随时”是指背压允许的时候?这究竟是如何工作的?如果你为每个输入元素生成一个输出流,你可以用这种方式对其进行建模,然后将流变平。
-
是的,当背压允许时。它不是查询/响应协议,客户端和服务器可以随时相互发送消息(网络允许)。
-
啊,所以输入和输出通道是完全分开的,背压也必须分开处理?
标签: akka reactive-programming akka-stream