【发布时间】:2016-10-22 23:39:33
【问题描述】:
动态变化的运行图有什么替代方法?这是我的情况。我有将文章摄入数据库的图表。文章来自 3 个不同格式的插件。因此我有几个流程
val converterFlow1: Flow[ImpArticle, Article, NotUsed]
val converterFlow2: Flow[NewsArticle, Article, NotUsed]
val sinkDB: Sink[Article, Future[Done]]
// These are being created every time I poll plugins
val sourceContentProvider : Source[ImpArticle, NotUsed]
val sourceNews : Source[NewsArticle, NotUsed]
val sourceCit : Source[Article, NotUsed]
val merged = Source.combine(
sourceContentProvider.via(converterFlow1),
sourceNews.via(converterFlow2),
sourceCit)(Merge(_))
val res = merged
.buffer(10, OverflowStrategy.backpressure)
.toMat(sinkDB)(Keep.both)
.run()
问题是我每 24 小时从内容提供商获取一次数据,每 2 小时从新闻获取一次数据,最后一个来源可能随时出现,因为它来自人类。
我意识到图表是不可变的,但我如何可以定期将 Source 的新实例附加到我的图表,以便对摄取过程进行单点限制?
更新:你可以说我的数据是Source-s 的流,在我的例子中是三个来源。但我无法改变这一点,因为我从外部类(所谓的插件)中获得了 Source 的实例。这些插件独立于我的摄取类工作。我不能将它们组合成一个巨大的类来拥有单个 Source。
【问题讨论】:
-
不清楚为什么需要附加新的来源。您说您有来自内容提供商、新闻和手动输入的数据;因此,你有三个来源,不多也不少。所以你的代码对我来说很好。
-
我会定期在任意时间获取需要摄取的
Source类的新实例。因此,我想避免在从内容提供商那里摄取 10K 文章时出现这种情况,并且在其中我从 2K 项目的新闻中得到Source。我希望它们同时被摄取并尊重我的单一限制规则。 -
我建议不要将您的数据流建模为
Sources 的序列,而是将其建模为单个Source,它会按顺序生成所有数据。那么Merge组合子就足够了。顺便说一句,我不确定您的设计如何避免您描述的情况。 -
@VladimirMatveev 不确定我理解。我更新了问题。
-
@VladimirMatveev 可能还有其他方法可以使此代码正常工作,但我对提出的问题非常感兴趣,因为它还有许多其他应用程序。
标签: scala akka akka-stream