【问题标题】:What, exactly happens when a repartition occurs in a kafka stream?当在 kafka 流中发生重新分区时,究竟会发生什么?
【发布时间】:2019-07-29 20:10:51
【问题描述】:

假设我有一组员工,以empId 为关键字,其中还包括departmentId。 我想按部门汇总。所以我做了一个selectKey(mapper 来获取departmentId),然后是groupByKey()(或者我可以只做一个groupBy(...),我假设),然后,比如说,count()。究竟会发生什么?我认为它会进行“重新分区”。我认为发生的事情是它写入“内部”主题,我只是一个具有派生名称的常规主题,自动创建。也就是说,由流的所有实例共享,而不仅仅是一个(即不是本地的)。所以聚合是跨越所有新键,而不仅仅是来自源流实例的那些消息(我认为)。对吗?

我没有找到关于重新分区的全面描述。有人能指点我这方面的好文章吗?

【问题讨论】:

  • 我不知道你在哪里听说过“重新分区”这个词。在我看来,分区存储实际的消息。因此,重新分区对我来说听起来很可怕。
  • Kafka 流的输入是一个主题。在主题中,生产者将数据推送到分区。生产者将具有相同密钥的消息发送到同一分区。现在,流 API 为您提供数据存储来存储操作结果。 kafka.apache.org/20/documentation/streams/developer-guide/…
  • @JRibkr:首先,KStream javadoc 中提到了 88 次“重新分区”。我想我已经掌握了它的要点,但我还没有看到任何详细的描述,并且“内部”主题的范围可能可以解释。此外,您的链接指向交互式查询,这不是我要说的。
  • 你是对的。有趣的。这为我打开了另一扇门。我正在阅读有关重新分区的信息。收到后会及时更新。
  • kafka.apache.org/20/javadoc/org/apache/kafka/streams/kstream/… 有关于重新分区的简要说明。 Kafka 将为流处理 ${applicationId}-XXX-repartition 创建重新分区的主题。您还可以配置要保留该内部主题的时间。

标签: apache-kafka apache-kafka-streams


【解决方案1】:

你所描述的正是正在发生的事情。

重新分区步骤与through() 相同(自动插入到处理拓扑中)是to("topic") 加上builder.stream("topic") 的快捷方式。

这篇博文中也有说明和解释:https://www.confluent.io/blog/data-reprocessing-with-kafka-streams-resetting-a-streams-application/

【讨论】:

  • 谢谢。我在搜索中看到过那篇文章,但由于专注于重置,我没有深入了解它的相关性。尽管如此,对于我们这些新手来说,很高兴看到一篇关于重新分区的帖子。当它发生时,性能影响,(有点令人惊讶)security implications
猜你喜欢
  • 1970-01-01
  • 2023-04-01
  • 2015-12-25
  • 2018-12-01
  • 1970-01-01
  • 2023-03-29
  • 1970-01-01
  • 2015-10-23
  • 2011-01-18
相关资源
最近更新 更多