代码之家  ›  专栏  ›  技术社区  ›  erip Jigar Trivedi

如何在流中使用确认语义?

  •  0
  • erip Jigar Trivedi  · 技术社区  · 8 年前

    我想设计一个只有在 ActorRef 接收一些确认消息。

    我的用例是使用一些静态集群参与者来完成工作。我有一个监护人,负责维护是否有可用的工人。当参与者可用时,它将发布消息。我希望我的流程能够意识到这一信息,从而完成一项工作并将其推到下游。否则会产生反压力。

    在阅读文档时, GraphStage API似乎没有一些反应性组件,依赖于 InHandler s用于自定义处理:

    setHandler(in, new InHandler {
      override def onPush(): Unit = {
        println(grab(in))
        pull(in)
      }
    })
    

    我可以在这个方法中进行投票,以更新可用参与者的数量,但这不是被动的。

    默认情况下,Akka是否支持流中的确认?

    1 回复  |  直到 8 年前
        1
  •  1
  •   Ramón J Romero y Vigil    8 年前

    推而不拉

    Actor Source

    ActorRef

    object ActorIsFreeMessage
    
    val source : Source[ActorIsFreeMessage, ActorRef] = Source.actorRef(???, ???)
    

    然后可以附加 Flow

    type Job = ???
    
    val pullAJob : () => Job = ???
    
    val jobFlow : Flow[ActorIsFreeMessage, Job] = 
      Flow[ActorIsFreeMessage].map[Job](_ => pullAJob())
    
    val jobSource : Source[Job, ActorRef] = source via jobFlow
    

    一旦这个新的源连接到一个接收器,其他参与者就可以向物化的发送消息。

    val jobSink : Sink[Job, ActorRef] = ???
    
    val streamRef = jobSource.to(jobSink).run()
    
    //inside of the Actor
    streamRef ! ActorIsFreeMessage
    
    推荐文章