【问题标题】:Streaming events from Kafka to Couchbase using Akka Stream and Kafka offset committing使用 Akka Stream 和 Kafka 偏移提交将事件从 Kafka 流式传输到 Couchbase
【发布时间】:2020-02-15 12:45:55
【问题描述】:

我正在尝试使用 Alpakka 设计一个 Akka Stream 来读取来自 kafka 主题的事件并将它们放入 Couchbase。

到目前为止,我有以下代码,它似乎以某种方式工作:

Consumer
      .committableSource(consumerSettings, Subscriptions.topics(topicIn))
      .map(profile ⇒ {
        RawJsonDocument.create(profile.record.key(), profile.record.value())
      })
      .via(
        CouchbaseFlow.upsertDoc(
          sessionSettings,
          writeSettings,
          bucketName
        )
      )
      .log("Couchbase stream logging")
      .runWith(Sink.seq)

“不知何故”是指流实际上是从主题读取事件并将它们作为 json 文档放入 Couchbase,尽管事实上我不明白如何将消费者偏移量提交给 Kafka,但它看起来甚至不错。

如果我已经清楚地理解了隐藏在 Kafka 消费者偏移量背后的主要思想,那么在发生任何故障或重新启动的情况下,流会从上次提交的偏移量中读取所有消息,并且由于我们没有提交任何消息,因此它可能再次重新读取上一个会话中正在读取的记录。

那么我的假设是否正确?如果是这样,如果从 Kafka 读取并发布到某个数据库,如何处理消费者提交? Akka Streams 官方文档提供了示例,展示了如何使用普通 Kafka Streams 处理此类情况,因此我不知道如何在我的情况下提交偏移量。

非常感谢!

【问题讨论】:

  • 尝试在所有函数调用中传递CommittableOffset,然后使用Commiter.sink 捕获偏移量并进行回写。
  • 到目前为止,我的主要想法是我必须使用toMatvia(Committer.flow) 向相关接收器提供CommitableOffset,但问题是我无法绕开我的脑袋一种可能的实现方式。
  • 为什么不使用 Couchbase 提供的 Kafka Connect 插件?
  • @cricket_007 我不知道,会看看,谢谢!

标签: scala apache-kafka akka akka-stream alpakka


【解决方案1】:

您需要在 Couchbase 中提交偏移量才能获得“恰好一次”语义。

这应该会有所帮助:https://doc.akka.io/docs/alpakka-kafka/current/consumer.html#offset-storage-external-to-kafka

【讨论】:

  • 嗨!谢谢你的回答!在这种情况下,我通常不需要“恰好一次”语义,我只是想确保我的应用在重启后不会重新读取未提交的偏移量。
猜你喜欢
  • 1970-01-01
  • 2018-11-28
  • 2020-10-10
  • 2018-03-27
  • 1970-01-01
  • 2019-03-11
  • 2018-01-19
  • 2016-03-17
  • 2023-03-26
相关资源
最近更新 更多