【问题标题】:How do I dynamically add Source to existing Graph?如何将 Source 动态添加到现有 Graph?
【发布时间】: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


【解决方案1】:

如果您不能将其建模为Source[Source[_,_],_],那么我会考虑使用Source.queue[Source[T,_]](queueSize, overflowStrategy)here

你必须小心的是,如果提交失败会发生什么。

【讨论】:

    【解决方案2】:

    好的,通常正确的方法是将源流加入单个源,即从Source[Source[T, _], Whatever]Source[T, Whatever]。这可以通过flatMapConcatflatMapMerge 完成。因此,如果您可以获得Source[Source[Article, NotUsed], NotUsed],您可以使用flatMap* 变体之一并获得最终的Source[Article, NotUsed]。为您的每个来源都这样做(没有双关语),然后您的原始方法应该有效。

    【讨论】:

      【解决方案3】:

      我已经根据 Vladimir Matveev 给出的答案实现了代码,并希望与其他人分享它,因为它对我来说似乎是常见的用例。

      我知道 Viktor Klang 提到的 Source.queue,但我不知道 flatMapConcat。真是太棒了。

      implicit val system = ActorSystem("root")
      implicit val executor = system.dispatcher
      implicit val materializer = ActorMaterializer()
      
      case class ImpArticle(text: String)
      case class NewsArticle(text: String)
      case class Article(text: String)
      
      val converterFlow1: Flow[ImpArticle, Article, NotUsed] = Flow[ImpArticle].map(a => Article("a:" + a.text))
      val converterFlow2: Flow[NewsArticle, Article, NotUsed] = Flow[NewsArticle].map(a => Article("a:" + a.text))
      val sinkDB: Sink[Article, Future[Done]] = Sink.foreach { a =>
        Thread.sleep(1000)
        println(a)
      }
      
      // These are being created every time I poll plugins
      val sourceContentProvider: Source[ImpArticle, NotUsed] = Source(List(ImpArticle("cp1"), ImpArticle("cp2")))
      val sourceNews: Source[NewsArticle, NotUsed] = Source(List(NewsArticle("news1"), NewsArticle("news2")))
      val sourceCit: Source[Article, NotUsed] = Source(List(Article("a1"), Article("a2")))
      
      val (queue, completionFut) = Source
        .queue[Source[Article, NotUsed]](10, backpressure)
        .flatMapConcat(identity)
        .buffer(2, OverflowStrategy.backpressure)
        .toMat(sinkDB)(Keep.both)
        .run()
      
      queue.offer(sourceContentProvider.via(converterFlow1))
      queue.offer(sourceNews.via(converterFlow2))
      queue.offer(sourceCit)
      queue.complete()
      
      completionFut.onComplete {
        case Success(res) =>
          println(res)
          system.terminate()
        case Failure(ex) =>
          ex.printStackTrace()
          system.terminate()
      }
      
      Await.result(system.whenTerminated, Duration.Inf)
      

      我仍然会检查 queue.offer 返回的 Future 是否成功,但在我的情况下,这些调用将很少见。

      【讨论】:

        猜你喜欢
        • 2014-06-14
        • 2015-05-18
        • 2011-08-22
        • 2014-09-06
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2013-06-01
        • 2018-06-07
        相关资源
        最近更新 更多