【问题标题】:Kafka streams concurrency behaviorKafka 流并发行为
【发布时间】:2018-02-02 06:45:16
【问题描述】:

如果我的 kafka 流应用程序中有一个共享变量,并且在处理代码中由多个线程更新,它是如何处理的?我是否必须使该共享变量线程安全,或者 Kafka 流库如何处理?在文档的某处,我读到运行 Kafka 流应用程序时不需要在线程之间进行协调。例如,这是一个伪代码:

KStream<byte[], byte[]> input = ...;
int counter = 0;

KStream<byte[], byte[]>[] processed = input.map(
    (k, v) -> {
      ....
      ....
      //update counter by multiple threads.
);

如果此代码由来自同一应用实例的多个流任务执行,计数器会发生什么情况?变量“已处理”如何,因为它也可以由多个线程更新?这需要在普通 Java 场景中进行某种同步。我很好奇这是否由 Kafka 流库处理。

谢谢!

【问题讨论】:

    标签: apache-kafka-streams


    【解决方案1】:

    这取决于您配置了多少线程来执行任务。如果您有一个线程执行所有任务,那么您不必使该共享变量线程安全。但是,如果您有多个线程,则需要使其成为线程安全的,因为应用程序实例中的任务将分布在多个线程中。您的 Kafka Streams 应用程序只是一个以 main() 开头的正在运行的 JVM。 Kafka Streams 框架根据您指定的线程数编排处理。但它只是一个常规的 Java 运行时,并发访问仍然是并发访问。

    更多关于线程和任务的信息:Kafka Streams thread number

    有关线程和任务以及共享状态的更多信息:Kafka stream processor thread safe?

    显然,一般而言,您在代码示例中显示的模式是您可能想要避免的模式,除非它实际上只是计算应用程序本地的某些内容。在您运行多个应用程序实例的生产应用程序中,如果应用程序实例上升或下降,任务会重新分配,因此您的共享变量可能不会有用。这就是 Kafka Streams 存储机制如此有用的原因:您的状态会随着任务而变化。

    【讨论】:

    • 感谢您的回答。同意“计数器”示例,我只是以它为例。你的回答证实了我的想法。基本上,我需要在任何共享数据结构上进行同步。然后我发现了这条评论 - stackoverflow.com/questions/39985048/…。该评论似乎表明可以在线程之间共享资源。
    • 我认为它的意思是,由于一个值​​连接器最多属于一个线程(也只能由一个线程执行),它可能具有的任何状态都是线程安全的。也就是说,如果您将 value joiner 实现为具有某些本地私有属性的完整类而不是 lambda,则它的 apply 函数最多只能由一个线程调用,因此本地私有属性不必同步。但是您示例中的状态可能会被许多线程共享,因为它可以被运行在不同线程中的许多值连接器访问。所以需要同步。
    猜你喜欢
    • 1970-01-01
    • 2019-09-17
    • 2016-05-31
    • 1970-01-01
    • 2017-04-23
    • 2018-04-25
    • 1970-01-01
    • 2018-02-25
    • 1970-01-01
    相关资源
    最近更新 更多