【问题标题】:Kafka consumer threads with same group id consuming same record具有相同组 ID 的 Kafka 消费者线程使用相同的记录
【发布时间】:2019-03-19 04:11:54
【问题描述】:

我需要在多个线程中使用来自 Kafka 分区的记录,并在每个线程上使用唯一的记录来处理。 我有以下代码,我不知道是什么错误

public class ConsumerThread implements Runnable {
    public String name;
    public ConsumerThread(String name){
        this.name = name;
    }
    public Properties getDefaultProperty(){
        Properties prop = new Properties();
        prop.setProperty("group.id", "4");
        prop.put("enable.auto.commit", "false");
        prop.put("auto.offset.reset", "earliest");
        prop.setProperty("bootstrap.servers", "localhost:9092");
        prop.setProperty("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
        prop.setProperty("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
        prop.setProperty("max.poll.records","150");
        return prop;
    }
    public void run() {
        TopicPartition tp = new TopicPartition("my.topic", 0);
        KafkaConsumer consumer = new KafkaConsumer(getDefaultProperty());
        ArrayList tpList = new ArrayList<TopicPartition>();
        tpList.add(tp);
        consumer.assign(tpList);
        ConsumerRecords poll = consumer.poll(1000);
        Iterator it = poll.iterator();
        consumer.commitAsync();
        while(it.hasNext()){
            ConsumerRecord cr = (ConsumerRecord) it.next();
            System.out.println("From "+this.name+" : "+cr.value());
        }
        consumer.close();
        System.out.println("Thread Exiting "+this.name);
    }
}

结果

From Thread1 : produced_0
From Thread1 : produced_1
From Thread1 : produced_2
From Thread1 : produced_3
.
.
.
From Thread1 : produced_136
From Thread2 : produced_0
From Thread2 : produced_1
From Thread2 : produced_2
From Thread2 : produced_3
.
.
.


预期:

From Thread1 : produced_0
From Thread1 : produced_1
From Thread1 : produced_2
From Thread1 : produced_3
.
.
.
From Thread1 : produced_136
From Thread2 : produced_4
From Thread2 : produced_5
From Thread2 : produced_6
From Thread2 : produced_137

【问题讨论】:

  • 看起来你有多个线程只订阅了分区 0,不保证线程消耗顺序
  • 如果您希望能够使用来自相同主题的相同记录,这应该是两个独立的组,并且具有不同的group.id 分配。如果这不是您想要完成的,请提供更多信息。

标签: java multithreading apache-kafka kafka-consumer-api producer-consumer


【解决方案1】:

只有使用 kafka 消费者的subscribe 方法才能将分区自动分配给消费者组。 但是,您将 assign 与特定主题分区一起使用,因此您负责将特定分区分配给不同的消费者(但您始终使用相同的分区 0,因此所有消费者都从同一个主题分区消费)。

【讨论】:

  • 已使用订阅,但它以某种方式同步线程1 接收所有消息或线程2 接收所有消息我希望两个线程都接收唯一消息而不是相同消息。 ``` public void run() { ArrayList ar = new ArrayList(); KafkaConsumer 消费者 = 新的 KafkaConsumer(getDefaultProperty()); ar.add("my.topic");消费者.订阅(ar); ConsumerRecords poll = consumer.poll(1000); ```
  • 这是不可能的。假设两个线程中的消费者都使用相同的group.id,并且您在max.poll.interval.ms (kafka.apache.org/documentation/#newconsumerconfigs) 内提交,这将导致在消费期间重新加入。根据您使用的版本,您受到issues.apache.org/jira/browse/KAFKA-5430 影响的可能性很小,但这种情况会发生一次,并且不容易重现。
【解决方案2】:

就像 Lior Chaga 在他的评论中所说,您手动将主题分区分配给您的消费者。这不是推荐的方法。最重要的是,您的所有消费者似乎都使用相同的确切 groupID。使用这种配置,两个线程消耗,如果至少有一个消费者收到特定消息,则 none 其他线程将收到该消息。如果您希望所有消费者线程各自获得自己的“一组”消息,而不会相互中断,那么您需要给它们不同的group.ids。

要订阅一个主题,以便它为您处理自动重新平衡,然后消费,您应该执行以下操作(取自下面链接的 KafkaConsumer javadoc):

 consumer.subscribe(Arrays.asList("foo", "bar"));
 while (true) {
     ConsumerRecords<String, String> records = consumer.poll(100);
     for (ConsumerRecord<String, String> record : records)
         System.out.printf("offset = %d, key = %s, value = %s%n", record.offset(), record.key(), record.value());
 }

Kafka 官方 javadocs 有更详细的解释: https://kafka.apache.org/20/javadoc/index.html?org/apache/kafka/clients/consumer/KafkaConsumer.html

【讨论】:

  • 如你所说。我不希望其他线程使用相同的消息。使用相同的 group.id 订阅该主题将为每个线程提供唯一的消息,对吗?
猜你喜欢
  • 2021-12-04
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2023-03-24
  • 2020-02-10
  • 1970-01-01
  • 2020-04-19
  • 1970-01-01
相关资源
最近更新 更多