【问题标题】:How to get a Map [String, String] returned by a Kafka Consumer (Alpakka)?如何获取 Kafka 消费者(Alpakka)返回的 Map [String, String]?
【发布时间】:2020-12-06 16:48:51
【问题描述】:

我应该从 Kafka 消费者那里得到一个 Map [String, String],但我真的不知道怎么做。我设法配置了消费者,它工作正常,但我不明白如何获得地图。

implicit val system: ActorSystem = ActorSystem()

    val consumerConfig = system.settings.config.getConfig("akka.kafka.consumer")

    val = kafkaConsumerSettings =
      ConsumerSettings(consumerConfig, new StringDeserializer, new StringDeserializer)
        .withBootstrapServers(localhost:9094)
        .withGroupId(group1)
    Consumer
      .plainSource(kafkaConsumerSettings, Subscriptions.topics(entity.entity_name))
      .toMat(Sink.foreach(println))(DrainingControl.apply)
      .run()

【问题讨论】:

  • 从写下你尝试过的东西和你现在得到的回报开始。一些代码会很有帮助
  • @GamingFelix 我编辑了问题,现在有代码
  • a Map 哪些键和哪些值?主题中每个键和关联值的映射?
  • 如果这就是你想要的(即你从该主题中读取有限的值),我真的会质疑 Alpakka Kafka 是否非常适合:它更面向流(即来自主题的无限数量的值)方法。

标签: scala apache-kafka akka alpakka


【解决方案1】:

Lightbend's recommendation是在反序列化来自Kafka的传入数据的同时处理字节数组

消息的反序列化的一般建议是使用字节数组(或字符串)作为值,并在 Akka Stream 的映射操作中进行反序列化,而不是直接在 Kafka 反序列化器中实现它.在 Akka Stream 中显式处理反序列化时,更容易实现所需的错误处理策略,如下例所示。

为此,您可以使用此设置设置消费者:

val consumerSettings = ConsumerSettings(consumerConfig, new StringDeserializer, new ByteArrayDeserializer)

并通过调用 Record 类中的 .value() 方法来获取结果。要反序列化它,我建议使用 circe + jawn。这段代码应该可以解决问题。

import io.circe.jawn
import io.circe.generic.auto._

val bytes = record.value()

val data = jawn.parseByteBuffer(ByteBuffer.wrap(bytes)).flatMap(_.as[Map[String, String]])

【讨论】:

    猜你喜欢
    • 2020-07-15
    • 2016-11-27
    • 2010-12-28
    • 2017-06-02
    • 2022-11-02
    • 1970-01-01
    • 2015-09-23
    • 1970-01-01
    • 2016-07-10
    相关资源
    最近更新 更多