【发布时间】: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