【问题标题】:How does stateful operations work in Kafka streams when there are multiple instances of the stream app?当流应用程序有多个实例时,有状态操作如何在 Kafka 流中工作?
【发布时间】:2019-01-18 10:12:45
【问题描述】:

状态完整操作如何在具有多个实例的 Kafka Stream 应用程序中工作? 假设我们有 2 个主题,每个 A 和 B 有 2 个分区。 我们有一个流应用程序,它从两个主题中消费,并且两个流之间存在连接。

现在我们正在运行这个流应用程序的 2 个实例。据我了解,每个实例都将被分配到每个主题的 2 个分区之一。

现在,如果要加入的消息被应用程序的不同实例消费,那么加入将如何进行?我无法理解它。

虽然我测试了一个似乎工作正常的小型流应用程序。我是否可以在不考虑流应用程序中定义的拓扑类型的情况下始终增加任何类型应用程序的实例数量?

有什么文件可以让我了解它的工作细节吗?

【问题讨论】:

  • 状态由您的 Kafka 集群维护——为此它会创建任意数量的内部主题。同一个应用程序的所有实例都使用同一个集群。

标签: apache-kafka apache-kafka-streams


【解决方案1】:

并行度的最大单位是分区数。如果您运行的实例多于分区数,则过多的实例将处于空闲状态。

连接操作应满足以下要求:

  1. 在加入时输入数据必须共同分区。这意味着要加入的输入主题应该具有相同数量的分区。

  2. 两个主题应该具有相同的分区策略,以便可以将具有相同键的记录传递到相同的分区。如果不同,就有可能丢失记录。

例子:如果topic1有2个partition,topic2有3个partition,Join(topic1,topic2)会因为partition不等而失败。重新划分主题后,让我们说 3。 现在Join(topic1, topic2) 可以工作了。此操作最多可使用 3 个任务。每个分区将以内部主题的形式在状态存储中维护其状态。默认情况下,KStream 使用 RocksDB 来存储状态。

在这里,您可以看到该过程通常如何用于有状态转换:

请参阅这些以了解详细信息:

https://cwiki.apache.org/confluence/display/KAFKA/Kafka+Streams+Internal+Data+Management https://docs.confluent.io/current/streams/developer-guide/dsl-api.html#streams-developer-guide-dsl-joins

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2016-05-06
    • 1970-01-01
    • 2022-09-27
    相关资源
    最近更新 更多