代码之家  ›  专栏  ›  技术社区  ›  Funzo

如何从操作列表中创建接收器

  •  0
  • Funzo  · 技术社区  · 7 年前

    我想在阿克卡河上建造一个水槽,它由许多操作组成。 例如地图,过滤,折叠,然后下沉。 目前我能做的最好的事情是: 我不喜欢它,因为我必须指定广播,即使我只允许一个值通过。 有人知道更好的方法吗?

    def kafkaSink(): Sink[PartialBatchProcessedResult, NotUsed] = {
        Sink.fromGraph(GraphDSL.create() { implicit b =>
        import GraphDSL.Implicits._
        val broadcast = b.add(Broadcast[PartialBatchProcessedResult](1))
        broadcast.out(0)
        .fold(new BatchPublishingResponseCollator()) { (c, e) => c.consume(e) }
        .map(_.build())
        .map(a =>
          FunctionalTesterResults(sampleProjectorConfig, 0, a)) ~> Sink.foreach(new KafkaTestResultsReporter().report)
      SinkShape(broadcast.in)
    })
    

    }

    1 回复  |  直到 7 年前
        1
  •  0
  •   Ramón J Romero y Vigil    7 年前

    记住一个关键点 akka-stream 那有多少 Flow 价值加上 Sink 值仍然是一个接收器。

    演示此属性的几个示例:

    val intSink : Sink[Int, _] = Sink.head[Int]
    
    val anotherSink : Sink[Int, _] = 
      Flow[Int].filter(_ > 0)
               .to(intSink)
    
    val oneMoreSink : Sink[Int, _] = 
      Flow[Int].filter(_ > 0)
               .map(_ + 4)
               .to(intSink)
    

    因此,您可以实现 map filter 作为流动。这个 fold 你问的问题可以用 Sink.fold .