【问题标题】:Reading already partitioning topic in Kafka Streams DSL在 Kafka Streams DSL 中读取已经分区的主题
【发布时间】:2020-12-06 12:51:45
【问题描述】:

在 Kafka Streams 中对大容量主题进行重新分区可能会非常昂贵。一种解决方案是在生产者端通过键对主题进行分区,并在 Streams 应用中提取已分区的主题。

有没有办法告诉 Kafka Streams DSL 我的源主题已经被给定的键分区并且不需要重新分区?


让我澄清一下我的问题。假设我有一个像这样的简单聚合(为简洁起见省略了详细信息):

builder
    .stream("messages")
    .groupBy((key, msg) -> msg.field)
    .count();

鉴于此代码,Kafka Streams 将读取 messages 主题并立即将消息写回内部重新分区主题,这次以 msg.field 作为键进行分区。

使这种往返不必要的一种简单方法是首先编写由msg.field 分区的原始messages 主题。但是 Kafka Streams 对 messages 主题分区一无所知,而且我没有办法告诉它主题是如何分区而不引起真正的重新分区的。

请注意,我并不是要完全消除分区步骤,因为主题必须被分区才能计算键控聚合。我只是想将上游的分区步骤从 Kafka Streams 应用程序转移到原始主题生产者。

我要找的基本上是这样的:

builder
    .stream("messages")
    .assumeGroupedBy((key, msg) -> msg.field)
    .count();

其中assumeGroupedBy 会将流标记为已经msg.field 分区。我知道这个解决方案有点脆弱,并且会在分区键不匹配时中断,但它解决了处理大量数据时的问题之一。

【问题讨论】:

  • 感谢您更新我们的答案,鲍里斯。您是否检查过groupByKey() 函数(而不是groupBy(),它总是 会导致对其输入数据进行重新分区)?它假定输入数据已经根据现有消息键的需要进行了分区。在您的示例中,如果key == msg.fieldgroupByKey() 将起作用。
  • 我错过了最明显的解决方案。多可惜!不知何故,我假设groupByKey 需要通过事先调用selectKeygroupBy 来设置密钥。 Wish 文档明确指出这不是必需的。 @MichaelG.Noll 能否请您更新您的答案,以便我将其标记为已接受?
  • 别担心,很容易错过!我更新了我的答案。
  • 关于文档混乱:关于我们如何改进它的任何建议? kafka.apache.org/documentation/streams/developer-guide/… 今天说:“当且仅当流被标记为重新分区时才会导致数据重新分区。groupByKey 比 groupBy 更可取,因为它仅在流已标记为重新分区时才重新分区数据。但是, groupByKey 不允许您像 groupBy 那样修改密钥或密钥类型。"
  • 关于文档,这有意义吗? “当且仅当流被标记为重新分区时才会导致数据重新分区。groupByKey 比 groupBy 更可取,因为它仅在流已标记为重新分区时重新分区数据,否则它假定输入数据已经根据现有消息键根据需要进行了分区。如果您只是想聚合数据而不引起重新分区操作,那么您只需要使用 groupByKey()。但是,groupByKey 不允许您像 groupBy 那样修改密钥或密钥类型。” @MichaelG.Noll

标签: apache-kafka apache-kafka-streams


【解决方案1】:

更新问题后更新:如果您的数据已经根据需要进行了分区,并且您只想聚合数据而不需要重新分区操作(两者都适用于您的用例),那么所有您需要使用groupByKey() 而不是groupBy()。而groupBy() 总是导致重新分区,其兄弟groupByKey() 假设输入数据已经根据现有消息键的需要进行了分区。在您的示例中,如果key == msg.fieldgroupByKey() 将起作用。

原答案如下:

在 Kafka Streams 中对大容量主题重新分区可能非常昂贵。

是的,没错,它可能非常昂贵(例如,当高容量意味着每秒数百万个事件时)。

有没有办法告诉 Kafka Streams DSL 我的源主题已经被给定的键分区并且不需要重新分区?

Kafka Streams 不会重新分区数据,除非您指示它;例如,使用 KStream#groupBy() 函数。因此,无需像您在问题中所说的那样告诉它“不要分区”。

一种解决方案是在生产者端通过键对主题进行分区,并在 Streams 应用程序中提取已分区的主题。

鉴于您的这种解决方法,我的印象是您提出问题的动机是其他的(您必须考虑到特定的情况),但是您的问题文本并没有明确说明可能是什么。也许您需要更新您的问题以提供更多详细信息?

【讨论】:

  • 是的,问题是关于每秒处理数百万个事件的主题。我更新了我的问题,并试图更好地解释我想要实现的目标。
猜你喜欢
  • 1970-01-01
  • 2018-03-07
  • 1970-01-01
  • 2020-01-07
  • 1970-01-01
  • 1970-01-01
  • 2017-12-18
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多