【问题标题】:Kafka offset management and sync with DBKafka 偏移管理和与数据库同步
【发布时间】:2021-03-07 01:59:50
【问题描述】:

我正在开发一个使用 Kafka 流和数据库的应用程序。

在我的应用程序中,我手动管理 Kafka 偏移量并仅在消息处理成功的情况下提交偏移量(即在处理和更新到数据库成功之后)。

但是,如果在更新数据库之后,我的应用程序在提交之前出现故障,那么当它恢复时,由于未提交的偏移量,它会导致重复写入数据库。

我想避免这些重复,同时仍确保我正在处理每条消息。这样做的正确方法是什么?

编辑:我对 DB 的更新基本上将记录的计数器增加了一些值。所以 MERGE 语句不是一个选项。

【问题讨论】:

  • 对于 RDBMS 使用 MERGE 语句。对于像 Cassandra 这样的 nosql 数据库,具有相同主键的重复行将简单地覆盖而不会出现任何错误。
  • @SaptarshiBasu 感谢您的回答。更新了我的问题,以说明为什么这些对我来说不是可行的选择。

标签: apache-kafka architecture kafka-consumer-api system-design


【解决方案1】:

这有点棘手。

Kafka 支持完全一次语义。但是,当您将数据写入外部数据存储时,您需要确保消费者端的恰好一次。

实现这一目标的一种方法(由 Jay Kreps here 提出)是将数据存储中的 Kafka 偏移量作为单个事务的一部分进行维护。因此,如果您维护每个分区的最后一个偏移量,当您收到的偏移量小于存储在数据库中的偏移量时,您总是可以忽略来自给定分区的消息。

但是,这种方法有一个警告。如果您有一个多数据中心主动-主动部署,如果主集群出现故障,消费者会回退到不同的不同数据中心集群,您不能盲目依赖偏移量。 Offset 是一个物理 id,消息在一个集群中的偏移量可以与复制消息在另一个集群上的偏移量不同。

在这种情况下,我认为正确的方法是利用 Kafka 流并维护存储在压缩 Kafka 主题中的 Kafka 表 (KTable) 中的计数。 Kafka 内部会使用生产者 id、epoch、事务 id 等来保证语义的精确。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2017-07-17
    • 1970-01-01
    • 2018-09-22
    • 2021-03-05
    • 2014-08-08
    • 2021-01-15
    • 2017-07-09
    相关资源
    最近更新 更多