【发布时间】:2019-11-13 10:40:36
【问题描述】:
我正在尝试完成 EFPL - EDx 平台上的反应式系统编程课程中的最后一个作业(名为反应式追随者)。
我能够完成除outgoingFlow之外的所有功能。
在我看来,我应该以某种方式从现有的 Flow 中创建一个新的 Source,但经过一些阅读,我仍然没有意识到如何执行 Flow 来为新的 Source 生成元素。
我尝试使用mapConcat,但没有成功。
我认为现有的流程是这样的:
eventParserFlow
.via(followersFlow)
.filter(p => isNotified(userId)(p))
现有Flows 的类型和我的暂定实现outgoingFlow 可以在这里看到:
val eventParserFlow: Flow[ByteString, Event, NotUsed]
val followersFlow: Flow[Event, (Event, Followers), NotUsed]
def outgoingFlow(userId: Int): Source[ByteString, NotUsed] = {
eventParserFlow
.via(followersFlow)
.filter(p => isNotified(userId)(p))
.mapConcat { case (e, _) => e.render }
???
}
谁能给我指点一些关于如何在 Akka 中解决类似问题的阅读或示例?
【问题讨论】:
标签: scala akka-stream