Подключить поток производителей к графику

Я новичок, использующий потоки akka kafka (и потоки akka в целом) . Я пытаюсь построить график, чтобы опубликовать сообщение на разные темы. Как я могу подключить производителя как поток, чтобы зафиксировать обработанные сообщения? Я пытался использовать Producer.flow, но не могу получить commitScaladsl

object TestFoo {
  import akka.kafka.ProducerMessage.Message
  implicit val system = ActorSystem("test-kafka")
  implicit val materializer = ActorMaterializer()
  val evenNumbersTopic = "even_numbers"
  val allNumbersTopic = "all_numbers"
  lazy val consumerSettings = ConsumerSettings(system, new StringDeserializer(), new JsonDeserializer[Int])
    .withBootstrapServers("localhost:9092")
    .withGroupId("group1")
    .withProperty(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest")
  lazy val source =  Consumer.committableSource(consumerSettings, Subscriptions.topics(Set(evenNumbersTopic, allNumbersTopic)))
  val producerSettings = ProducerSettings(system,  new StringSerializer(), new StringSerializer())
    .withBootstrapServers("localhost:9092")
  val flow: RunnableGraph[NotUsed] = RunnableGraph.fromGraph(GraphDSL.create() { implicit b =>
    import akka.stream.scaladsl.GraphDSL.Implicits._
    type TypedMessage =  Message[String, Int,CommittableOffset]
    val bcast = b.add(Broadcast[TypedMessage](2))
    val merge = b.add(Merge[TypedMessage](2))

    val evenFilter = Flow[TypedMessage].filter (  c => c.record.value() % 2 == 0)
    val justEven = Flow[TypedMessage].map{
      case Message(pr, offset) =>
      val r = new ProducerRecord[String, Int]("general", pr.value())
      Message(r, offset)
    }
    val allNumbers = Flow[TypedMessage].map{
      case Message(pr, offset) =>
      val r = new ProducerRecord[String, Int](allNumbersTopic, pr.value())
      Message(r, offset)
    }

    val toMsg = Flow[ConsumerMessage.CommittableMessage[String, Int]].map{ msg =>
      val r = new ProducerRecord[String, Int]("general", msg.record.value())
      Message(r, msg.committableOffset)
    }
    source ~> toMsg ~> bcast

    bcast ~> evenFilter ~> justEven ~> merge
    bcast ~> allNumbers ~> merge
    merge ~> Producer.flow(producerSettings).mapAsync(producerSettings.parallelism) { result =>
      result.message.passThrough.commitScaladsl() //this doesn't compile, cannot get the .commitScaladsl()
    }
    ClosedShape 
  })}

1 ответ

Решение

Поскольку вы используете GraphDSL, компилятор не может вывести PassThrough Тип из предыдущего этапа. Попробуйте и явно передайте параметры типа Producer.flow функция, например

merge ~> Producer.flow[K, V, CommittableOffset](producerSettings).mapAsync(producerSettings.parallelism) { result =>
  result.message.passThrough.commitScaladsl()
}

я ушел K а также V в качестве несвязанного параметра, пожалуйста, поместите туда все типы ключей / значений, которые ваш продюсер обязан произвести. Если вы хотите, чтобы приведенный выше код был правильно подключен, вам необходимо producerSettings типы с тем, что исходит от стадии слияния. Вам понадобится что-то вроде:

val producerSettings = ProducerSettings(system,  new StringSerializer(), new JsonSerializer[Int])
    .withBootstrapServers("localhost:9092")
Другие вопросы по тегам