【问题标题】:What are internal topics used in Kafka?Kafka 中使用了哪些内部主题?
【发布时间】:2019-09-28 13:22:57
【问题描述】:

我们正在使用 kafka 流 api 进行聚合,其中我们也使用 group by。 我们还使用状态存储来保存输入主题数据。

我注意到的是

Kafka 内部创建了 3 种主题

  1. Changelog-<storeid>-<partition>
  2. Repartition-<storeid>-<partition>
  3. <topicname>-<partition>

我无法理解的是

  1. 为什么当我拥有<topic>-<partition> 中的所有数据时它会创建更改日志主题
  2. 重新分区主题是否包含分组后的数据。
  3. 我发现 Changelog 和 topicname-parition 的大小大致相同。

数据有什么不同,因此必须为此保存不同的文件。

【问题讨论】:

标签: apache-kafka apache-kafka-streams


【解决方案1】:

'Changelog' 和 'repartition' 内部 Kafka 主题特定于 Kafka Streams。

来自卡夫卡维基,

Kafka Streams 允许有状态的流处理,即具有内部状态的操作符。这种内部状态在所谓的状态存储中进行管理。状态存储可以是临时的(失败时丢失)或容错的(失败后恢复)。 Kafka Streams DSL 使用的默认实现是容错状态存储,使用 1. 内部创建和压缩的变更日志主题(用于容错)和 2. 一个(或多个)RocksDB 实例(用于缓存键值查找)。因此,在启动/停止应用程序和倒带/重新处理的情况下,需要正确管理这些内部数据。

变更日志主题在流上有加入/聚合操作时创建。实际上,聚合调用的结果会创建一个状态存储,并且为了容错,状态存储由 Kafka Changelog 主题备份。

聚合结果存储到这个内部主题中。当应用程序重新启动且应用程序 ID 未更改时,状态将从更改日志主题中恢复。

重新分区主题是在对流进行关键修改操作时创建的。例如,groupByKey() 操作创建重新分区主题。查看JIRA page以了解更多关于自动创建重新分区主题的信息。

这两个内部主题使 Kafka 流具有容错的状态流处理能力。

repartition topic 是否包含分组后的数据? - 是

Changelog 和 topicname-parition 的大小差不多 - 可能所有聚合操作的结果都存储在这个 topic 中。

更多详情请查看Kafka Wiki page

【讨论】:

  • 为什么grouByKey被认为是修改键操作?
【解决方案2】:

有几种类型的内部 Kafka 主题:

  • __consumer_offsets 用于存储每个主题/分区的偏移提交。
  • __transaction_state 用于使用事务语义为 Kafka 生产者和消费者保持状态。
  • Schema Registry 使用_schemas 来存储所有模式、元数据和兼容性配置。
  • 以下三个主题是 Kafka Streams 使用的内部主题示例。前两个是常规的join信息,第三个其实是一个RocksDB持久化StateStore:
    • {consumer-group}--KSTREAM-JOINOTHER-0000000005-store-changelog
    • {consumer-group}--KSTREAM-JOINTHIS-0000000004-store-changelog
    • {consumer-group}--incompleteMessageStore-changelog

这里有更多信息:

【讨论】:

  • 问题更多是关于更新日志主题和重新分区主题,我对这些主题及其内容有疑问,并询问了问题中的具体疑问,请您帮忙解决这些问题
  • @mjuarez,重新分区主题和更改日志主题在哪里?
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2022-01-15
  • 2014-05-27
  • 1970-01-01
  • 2016-10-27
相关资源
最近更新 更多