Подключить поток производителей к графику
Я новичок, использующий потоки 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")