【发布时间】:2021-03-07 01:59:50
【问题描述】:
我正在开发一个使用 Kafka 流和数据库的应用程序。
在我的应用程序中,我手动管理 Kafka 偏移量并仅在消息处理成功的情况下提交偏移量(即在处理和更新到数据库成功之后)。
但是,如果在更新数据库之后,我的应用程序在提交之前出现故障,那么当它恢复时,由于未提交的偏移量,它会导致重复写入数据库。
我想避免这些重复,同时仍确保我正在处理每条消息。这样做的正确方法是什么?
编辑:我对 DB 的更新基本上将记录的计数器增加了一些值。所以 MERGE 语句不是一个选项。
【问题讨论】:
-
对于 RDBMS 使用 MERGE 语句。对于像 Cassandra 这样的 nosql 数据库,具有相同主键的重复行将简单地覆盖而不会出现任何错误。
-
@SaptarshiBasu 感谢您的回答。更新了我的问题,以说明为什么这些对我来说不是可行的选择。
标签: apache-kafka architecture kafka-consumer-api system-design