【问题标题】:How to create Source from Flow in Akka-stream? (Programming Reactive Systems activity)如何在 Akka-stream 中从 Flow 中创建 Source? (编程反应系统活动)
【发布时间】: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


    【解决方案1】:

    请注意 - SO 不是此类问题的最佳资源。你应该使用相应edx课程中的discussion部分


    关于你的问题 - 我不会给你明确的答案,只是一些提示。

    在 akka-streams 中,你不能只用 Flow 创建一个 SourceFlow 负责转换,而Source 创建新事件。在您的作业中,您只是忘记使用可用值之一。

    1. 仔细阅读class Server(不是object)中的cmets。
    2. 仔细查看val (inboundSink, broadcastOut) = ... 并尝试找出每个vals 的用途以及它们之间的关系以及与应用程序本身的关系。了解它们的类型会很有帮助

    这些提示应该足以理解如何实现outgoingFlow,也就是Source[ByteString, NotUsed]

    【讨论】:

    • 感谢您的回答。我没有使用聊天,因为我认为,由于课程不在第一个提供的版本中,所以不会有用户查看聊天。我意识到我的错误是尝试使用followersFlow 代替broadcastOut。过滤原始源并用 mapConcat 渲染每个 Event 解决了这个问题。
    猜你喜欢
    • 2020-08-14
    • 2015-12-02
    • 2017-08-31
    • 1970-01-01
    • 2018-09-17
    • 2016-01-21
    • 1970-01-01
    • 2019-12-05
    • 1970-01-01
    相关资源
    最近更新 更多