【发布时间】: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捕获偏移量并进行回写。 -
到目前为止,我的主要想法是我必须使用
toMat或via(Committer.flow)向相关接收器提供CommitableOffset,但问题是我无法绕开我的脑袋一种可能的实现方式。 -
为什么不使用 Couchbase 提供的 Kafka Connect 插件?
-
@cricket_007 我不知道,会看看,谢谢!
标签: scala apache-kafka akka akka-stream alpakka