变更日志主题和重新分区主题 kafka 流

Ern*_*sto 1 apache-kafka apache-kafka-streams spring-kafka

我想问一下,如果我不使用有状态流,我的 KafkaStreamsConfiguration 中是否需要复制因子。我不使用这个RockDB。据我所知,复制因子设置适用于更改日志和重新分区主题。我理解这个变更日志主题,但是这个重新分区主题让我有点困惑...有人可以用非常基本的语言向我解释这个重新分区主题是什么,以及如果我不在流应用程序中使用状态,我是否应该关心这个复制因子?

问候

Fel*_*ipe 5

简而言之,当您更改正在处理的事件/消息的键时,Kafka Streams 中就会发生重新分区。

重新分区基本上是流处理的洗牌阶段。这种情况可能发生在 Kafka 流、Apache Spark、Flink、Storm、Hadoop 等中。这些都是分布式流处理引擎 (DSPE),旨在并行执行任务以加快处理速度。然后,当您调用map转换时,DSPE 将此逻辑转换map为并行度 X 的物理map任务(X 通常是机器的核心数量)。

因此,如果您使用mapValues并且不更改密钥,Kafka流将不会重新分区。但是如果您使用 更改密钥map,Kafka 流将重新分区。此外,如果您使用任何聚合转换(例如:reduce、join、 ...),Kafka 将执行重新分区,因为它是基于键的。

当存在聚合阶段时,会发生重新分区/洗牌阶段。假设您有一个逻辑管道:

... -> map-> reduce-> ...

引擎盖下的物理管道将如下所示:

在此输入图像描述

具有相同键的事件通过groupByKey转换进行分组并发送到相同的reduce并行任务实例。这是洗牌阶段。

对于 Kafka 流,当发生聚合时,管道会从 变为 ,KStream因为KTable消息分布在 Kafka 代理上,并且流引擎必须查找不同分区上的事件。如果您使用 IntelliJ,当管道发生变化时,它会对您产生影响。在下图中,它正在进行字数统计,并且count转换是有状态的,就像reduce.

在此输入图像描述

这是阅读有关 Kafka Stream 中重新分区的更多信息的好来源。正如我所说,其他 DSPE 也依赖于洗牌阶段的重新分区。另一个不错的来源是 Flink 的这个来源。