【问题标题】:Autoscaling with KAFKA and non-transactional databases使用 KAFKA 和非事务性数据库进行自动缩放
【发布时间】:2019-05-14 16:50:10
【问题描述】:

说,我有一个从 KAFKA 读取一批数据的应用程序,它使用传入消息的键并对 HBase 进行查询(从 HBase 读取这些键的当前数据),进行一些计算并写入数据返回 HBase 获取相同的密钥集。例如

{K1, V1}, {K2, V2}, {K3, V3}(来自 KAFKA 的传入消息) --> 我的应用程序(从 HBase 读取 K1、K2 和 K3 的当前值,使用传入值 V1 , V2 和 V3 进行一些计算,并在处理完成后将 K1 (V1+x)、K2 (V2+y) 和 K3(V3+z) 的新值写回 HBase。

现在,假设我有一个 KAFKA 主题分区和 1 个消费者。我的应用程序有一个正在处理数据的消费者线程。

问题是说 HBase 出现故障,此时我的应用程序停止处理消息,并且在 KAFKA 中建立了巨大的延迟。即使我有能力增加分区数量和相应的消费者,但由于 HBase 中的 RACE 条件,我无法增加其中任何一个。 HBase 不支持行级锁定,所以现在如果我增加分区的数量,相同的键可能会转到两个不同的分区,并相应地转到两个不同的消费者,他们可能最终处于 RACE 状态,最后写入的人就是赢家。我必须等到所有消息都得到处理后才能增加分区数。

例如

HBase 出现故障 --> 最初我有一个主题分区,并且有未处理的消息 --> 分区 0 中的 {K3, V3} --> 现在我增加了分区的数量,并且带有键 K3 的消息现在是现在假设在分区 0 和 1 中 --> 然后从分区 0 消费的消费者和从分区 1 消费的另一个消费者最终将竞争写入 HBase。

有解决问题的办法吗?当然,由处理消息的消费者锁定密钥 K3 不是解决方案,因为我们正在处理大数据。

【问题讨论】:

    标签: apache-kafka kafka-consumer-api


    【解决方案1】:

    当您增加分区数量时,只有新消息会到达新添加的分区。 Kafka 负责只处理一次消息

    【讨论】:

    • 我认为这是错误的。在称为重新平衡的过程中,现有消息也会重新分配。
    • @Vassilis,如果你这么有信心,你能不能描述一下,例如,在这种情况下如何管理分区的偏移量?
    • @Vassilis,看看这个,相信你会改变你的观点。 stackoverflow.com/questions/29291051/…
    • @Vassilis 我 200% 确信在将分区添加到主题时不会重新分发消息。
    【解决方案2】:

    一条消息只会出现在一个且只有一个 kafka 分区中。它在消息上使用哈希函数以分区数为模。我相信这个保证可以解决您的问题。

    但请记住,如果您更改分区数,则可能会将相同的消息密钥分配给不同的分区。如果您关心仅保证每个分区的消息顺序,这可能很重要。如果您关心消息的顺序,则不能重新分区(例如增加分区数)。

    【讨论】:

      【解决方案3】:

      正如 Vassilis 所提到的,Kafka 保证单个 key 只会在一个分区中。 有different strategies如何在分区上分配密钥。
      当您增加分区数或更改分区策略时,可能会发生重新平衡过程,这可能会影响工作消费者。如果您停止消费者一段时间,您可以避免两个消费者处理相同密钥的可能性。

      【讨论】:

        猜你喜欢
        • 2021-06-03
        • 2015-04-02
        • 1970-01-01
        • 1970-01-01
        • 2015-01-03
        • 2020-02-07
        • 2021-05-27
        • 2014-08-10
        • 2020-09-10
        相关资源
        最近更新 更多