【发布时间】: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.field,groupByKey()将起作用。 -
我错过了最明显的解决方案。多可惜!不知何故,我假设
groupByKey需要通过事先调用selectKey或groupBy来设置密钥。 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