【问题标题】:In Kafka Streams, how do you parallelize complex operations (or sub-topologies) using multiple topics and partitions?在 Kafka Streams 中,如何使用多个主题和分区并行化复杂操作(或子拓扑)?
【发布时间】:2023-01-09 02:56:03
【问题描述】:

我目前正在尝试了解 Kafka Streams 如何实现并行性。我主要关心的问题归结为三个问题:

  1. 多个子拓扑可以从同一个分区读取吗?
  2. 如何并行使用处理器 API 并需要阅读整个主题的复杂操作(构成子拓扑)?
  3. 多个子拓扑可以从同一个主题读取(这样可以在不同的子拓扑中运行对同一主题的独立且昂贵的操作)吗?

    作为开发人员,我们无法直接控制拓扑如何划分为子拓扑。 Kafka Streams 在可能的情况下使用主题作为“桥梁”将拓扑划分为多个子拓扑。此外,创建了多个流任务,每个任务从输入主题中读取数据的一个子集,并按分区划分。 documentation 内容如下:

    稍微简化一下,您的应用程序可以运行的最大并行度受最大流任务数的限制,而最大流任务数本身由应用程序正在读取的输入主题的最大分区数决定。


    假设有一个子拓扑读取多个分区数量不相同的输入主题。如果相信上述文档摘录,则需要将分区较少的主题的一个或多个分区分配给多个流任务(如果需要读取两个主题以使逻辑起作用)。然而,这应该是不可能的,因为据我所知,流应用程序的多个实例(每个实例共享相同的应用程序 ID)充当一个消费者组,其中每个分区只分配一次.在这种情况下,为子拓扑创建的任务数量实际上应受其输入主题的最小分区数限制,即单个分区仅分配给一个任务。

    我不确定最初的问题,即非共同分区的子拓扑是否真的会发生。如果有一个操作需要读取两个输入主题,则数据可能需要共同分区(如在联接中)。


    假设两个主题(可能由多个自定义处理器构建)之间有一个昂贵的操作,需要一个主题的数据始终完整可用。您可能希望将此操作并行化为多个任务。

    如果主题只有一个分区,并且一个分区可以被多次读取,这将不是问题。但是,如前所述,我认为这行不通。

    然后是 GlobalKTables。但是,无法将 GlobalKTables 与自定义处理器一起使用(toStream 不可用)。

    另一个想法是将数据广播到多个分区,本质上是按分区计数复制数据。这样,可以为拓扑创建多个流任务来读取相同的数据。为此,可以在给定KStream#toProduced-Instance 中指定自定义分区程序。如果可以接受这种数据重复,这似乎是实现我的想法的唯一方法。


    关于第三个问题,因为 Streams 应用程序是一个消费者组,所以我也希望这是不可能的。根据我目前的理解,这将需要将数据写入多个相同的主题(同样本质上是复制数据),以便可以创建独立的子拓扑。另一种方法是运行单独的流式应用程序(以便使用不同的消费者组)。

【问题讨论】:

    标签: apache-kafka parallel-processing apache-kafka-streams stream-processing


    【解决方案1】:

    没有看到您的拓扑定义,这是一个有点模糊的问题。您可以拥有重新分区和更改日志主题。这些来自原始输入主题的重复数据。

    但是像 map, filter 等无状态操作符传递数据通过来自相同(分配的)分区对于每个线程。

    “子拓扑”仍然只是一个 application.id 的一部分,因此是一个消费者组,所以不,它不能多次读取相同的主题分区。为此,您需要通过整个拓扑中的分支操作来实现独立的流/表,例如,按偶数和奇数过滤数字只会消耗主题一次;您不需要将记录“广播”到所有分区,我不确定这是否可能开箱即用(to 一对一发送,Produced 定义序列化,而不是多个分区)。如果你需要交叉引用不同的运算符,那么你可以使用 join/statestores/KTables。


    这些都与并行性无关。您有 num.stream.threads,或者您可以运行同一 JVM 进程的多个实例以进行扩展。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2017-06-09
      • 1970-01-01
      • 1970-01-01
      • 2016-11-05
      • 2021-10-05
      • 2020-05-04
      • 2020-04-05
      相关资源
      最近更新 更多