Akka Streams:为何在onPull(_,_)外调用push(_,_)会阻塞流?
我来帮你拆解这个自定义Akka Streams Source的逻辑,这类Source的核心是通过Actor来驱动元素的生成——只有当下游(比如Sink或者中间流)准备好接收元素时,它才会向指定Actor请求数据,咱们一步步理清楚:
核心设计思路
这个ActorSource是典型的**“需求驱动”**源实现:它不会主动推送所有元素,而是等下游“要数据”的时候,才去问Actor要下一个元素,非常适合从异步Actor系统里拉取数据的场景。
代码逐段拆解(结合你给出的不完整代码补全逻辑)
1. 类定义与流形状
class ActorSource[T](context: ActorRefFactory, actor: ActorRef) extends GraphStage[SourceShape[T]] { val out: Outlet[T] = Outlet("actor-source") // 你代码里大概率漏了这行:定义流的形状,Source只有一个输出口 override val shape = SourceShape(out)
GraphStage[SourceShape[T]]:自定义流算子必须继承GraphStage,并指定它的“形状”——这里是SourceShape,说明这是一个只有输出端口(Outlet)的源。Outlet[T]:这是Source对外输出元素的端口,名字"actor-source"是给监控、调试用的标识,方便排查问题时定位。
2. 核心运行逻辑:GraphStageLogic
这部分是整个Source的灵魂,createLogic返回的GraphStageLogic负责处理流运行时的所有行为:
override def createLogic(inheritedAttributes: Attributes): GraphStageLogic = { new GraphStageLogic(shape) { setHandler(out, new OutHandler { // 你代码里的"receivin..."应该是`onPull()`方法,这是下游拉取元素时的核心回调 override def onPull(): Unit = { // 当下游准备好接收元素了,就给指定的Actor发请求消息 actor ! RequestNextElement // 这里的消息是自定义的,比如你可以定义case object RequestNextElement } // 处理Actor返回的元素:必须用异步回调保证线程安全 private val pushElement = createAsyncCallback[T] { element => push(out, element) // 把Actor给的元素推送给下游 } // 初始化时,创建一个和流绑定的StageActor,用来接收目标Actor的消息 override def preStart(): Unit = { getStageActor { case (_, element: T) => pushElement.invoke(element) } } }) } }
关键逻辑细节:
onPull():下游(比如Sink)准备好处理元素时,会触发这个方法。这时候咱们向目标Actor发送“要下一个元素”的请求,实现按需拉取,避免数据积压。createAsyncCallback:Actor的消息线程和流的运行线程是分开的,直接在Actor线程里操作流的端口会有线程安全问题,所以用AsyncCallback来安全地把元素推送到流里。getStageActor:创建一个和当前流逻辑绑定的小型Actor,专门用来接收目标Actor返回的元素,然后通过回调把元素推给下游。
完整运行流程(按顺序走一遍)
- 流启动后,下游(比如
Sink.foreach)会主动向Source发送pull信号,触发onPull()。 onPull()里给指定的actor发请求消息(比如RequestNext)。- 目标Actor收到请求,生成/获取一个元素,发送回
ActorSource的StageActor。 StageActor收到元素后,通过AsyncCallback调用push(out, element),把元素交给下游。- 下游处理完元素,再次发送
pull信号,重复上述循环。 - 如果Actor没有更多元素了,可以发送一个自定义的
StreamComplete消息,ActorSource就调用complete(out)来结束整个流;如果Actor出问题了,发送错误消息,调用fail(out, exception)终止流并抛出异常。
要注意的坑点
- 必须处理流的终止:别只想着推送元素,也要处理Actor的结束/错误信号,否则流会一直挂着。
- 绝对不能直接在Actor线程里操作
push(out, element):一定要用AsyncCallback或者StageActor来中转,不然会出现线程安全问题,导致流状态混乱。 - 请求消息要幂等:如果下游因为某些原因多次触发
pull,Actor收到多个请求,要避免重复发送相同元素,可以用请求ID或者状态标记来处理。
内容的提问来源于stack exchange,提问作者cokeSchlumpf
相关产品推荐
相关产品推荐

