【问题标题】:Kafka stream processor thread safe?Kafka流处理器线程安全吗?
【发布时间】:2017-11-05 07:43:38
【问题描述】:

我知道之前有人问过这个问题:Kafka Streaming Concurrency?

但这对我来说很奇怪。根据文档(或者我可能遗漏了一些东西),每个分区都有一个任务,这意味着不同的处理器实例,并且每个任务都由不同的线程执行。但是当我测试它时,我看到不同的线程可以获得不同的处理器实例。因此,如果您想在处理器中保持任何内存状态(老式方式),您必须锁定?

示例代码:

public class SomeProcessor extends AbstractProcessor<String, JsonObject> {

   private final String ID = UUID.randomUUID().toString();

   @Override
   public void process(String key, JsonObject value) {
     System.out.println("Thread id: " + Thread.currentThread().getId() +" ID: " + ID);

输出:

线程 ID:88 ID:26b11094-a094-404b-b610-88b38cc9d1ef

线程 ID:88 ID:c667e669-9023-494b-9345-236777e9dfda

线程 ID:88 ID:c667e669-9023-494b-9345-236777e9dfda

线程 ID:90 ID:0a43ecb0-26f2-440d-88e2-87e0c9cc4927

线程 ID:90 ID:c667e669-9023-494b-9345-236777e9dfda

线程 ID:90 ID:c667e669-9023-494b-9345-236777e9dfda

有没有办法强制每个实例线程?

【问题讨论】:

    标签: java multithreading apache-kafka-streams


    【解决方案1】:

    每个实例的线程数是一个配置参数(num.stream.threads,默认值为1)。因此,如果您启动单个 KafkaStreams 实例,您将获得 num.stream.threads 线程。

    任务将工作拆分为并行单元(基于您输入的主题分区),并将分配给线程。因此,如果您有多个任务和一个线程,则所有任务都将分配给该线程。如果您有两个线程(所有 KafkaStreams 实例的总和),每个线程执行大约 50% 的任务。

    注意:因为 Kafka Streams 应用程序本质上是分布式的,所以如果您运行具有多个线程的单个 KafkaStreams 实例,或者每个具有一个线程的多个 KafkaStreams 实例,则没有区别。任务将分布在应用程序的所有可用线程上。

    如果您想在任务之间共享任何数据结构并且您有多个线程,则您有责任同步对该数据结构的访问。请注意,任务到线程的分配可以在运行时更改,因此,所有访问都必须同步。 但是,不建议使用这种模式,因为它会限制可扩展性。您应该设计没有共享数据结构的程序! 这样做的主要原因是,您的程序通常分布在多台机器上,因此不同的KafkaStreams 实例无论如何都无法访问共享数据结构。共享数据结构只能在单个 JVM 中工作,但使用单个 JVM 会阻止应用程序横向扩展。

    【讨论】:

    • 感谢@matthias-j-sax 的回复。我实际上知道这些选项。问题是我不会将我的流写为“非线程安全”并强制使用一个线程,因为这是不好的做法。我只是假设 Kafka 流将支持进程/标点 API 上的“线程安全”,类似于 AKKA 演员模式。但我想我错了。无论如何,现在我明白我除了使用老式锁定之外别无他法。再次感谢
    • 嗨@mathhias-j-sax - 我们也遇到了这个问题 - 处理器实例是否有可能同时在不同线程上运行 process() 和 punctuate()?
    • 这是不可能的。 process() 方法和注册的Punctuators 在一个线程中执行。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多