【问题标题】:How set different state directory for multiple instances of the same Kafka Streams application on a single machine如何在一台机器上为同一个 Kafka Streams 应用程序的多个实例设置不同的状态目录
【发布时间】:2023-01-17 06:19:22
【问题描述】:

从 2.6.0 版本开始,带状态的 KafkaStreams 锁定了状态目录目录,如文档所述

州名录。 Kafka Streams 将本地状态保存在状态目录下。每个应用程序在其主机上都有一个子目录,该子目录位于状态目录下。子目录的名称是应用程序 ID。与应用程序关联的状态存储在该子目录下创建。在一台机器上运行同一应用程序的多个实例时,此路径对于每个此类实例必须是唯一的。

在单机运行同一个应用的多个实例的场景下, 路径不能是像这样的随机路径/state/dir/{uuid}因为这个解决方案绕过了KAFKA-10716问题。

我的解决方案是有一个像这样的目录/state/dir具有顺序子目录,例如 0,1,2... 并且启动时的每个实例都会从 0 开始检查此子目录并找到第一个未锁定的子目录并将该目录用于状态目录.结果,进程 ID 从图元文件中读取,之前的任务将正确分配给新进程。

这是一个正确的解决方案吗?

在一台机器上为每个实例设置不同路径的最佳做法是什么?

【问题讨论】:

  • 一个实例可以使用多个num.stream.threads进行并发。为什么要运行多个 JVM 实例?
  • @OneCricketeer 如果应用程序崩溃,一个实例会出于任何原因导致单点故障。除此之外,在 KafkaStreams 有 30 个任务(每个线程一个)的场景中,出于上下文切换和 cpu 使用原因,最好使用多处理而不是多线程。
  • 如果 JVM 崩溃,可能有充分的理由(例如 OOM、NPE)。否则,可以将异常处理程序添加到流处理器。您始终可以使用进程调度程序来重新启动失败的进程,因此它不是真正的 SPoF
  • 你是对的,但当任务是 cpu 密集型时,多处理编程仍然有三个好处,例如,更好地使用多个 cpu 内核、更小的堆大小和 gc 时间、更短的上下文切换时间、线程等待时间。此外,如果由于任何未知原因任务进入关闭状态(线程未处理的异常),则只会重新启动一小部分任务。正如卡夫卡文件所说状态目录他们通过为每个实例设置一个唯一的目录来预测它,我们不能为所有规模扩展多线程编程,它只适用于小主题分区。
  • 无论如何,回到问题。除了唯一性之外,该文档没有规定任何解决方案。在运行时创建序号目录对我来说真的没有意义,因为你需要跟踪/检查锁,就像你说的那样。总的来说,您确实需要一些流程监督来确保每个实例以其正确的状态目录重新启动,这将在 Kafka api 之外完成。否则,您只需设置一个硬编码目录,在每个实例中都是唯一的,可以使用 supervisord 来模板化进程号

标签: apache-kafka apache-kafka-streams


【解决方案1】:

我有同样的问题,我也有一个类似于你的解决方案:

我已经创建了一个服务注册表。每个kafka streams实例在启动时都会请求一个instance-id。然后,服务注册表将返回一个从 0 开始的整数。如果出现第二个实例,它将获得 id 1。如果 instance-0 关闭并重新启动,它将再次获得 id 0。 instance-id 用于设置group.instance.id 和state.dir。

为了使其更可靠,每个实例都会定期向服务注册中心发送心跳请求。这需要使实例 ID 再次可用,以防实例出现故障。

【讨论】:

    猜你喜欢
    • 2010-11-26
    • 1970-01-01
    • 2019-11-30
    • 1970-01-01
    • 2017-04-26
    • 2020-02-12
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多