我有一个卡夫卡消费者的问题。我使用新的 Kafka 和新的 Consumer Java API。从快速入门开始,它是最简单的 Kafka 和 Zookeeper 。
我启动应用程序,在我的消费者从主题消费消息几次后,它停止接收。
import java.util.Arrays;
import java.util.List;
import java.util.Properties;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import org.apache.kafka.common.TopicPartition;
public class MyKC{
public MyKC(){
Properties config = new Properties();
config.put("zookeeper.connect", "localhost:2181");
config.put("group.id", "default");
config.put("bootstrap.servers", "localhost:9092");
config.put("enable.auto.commit", "true");
config.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
config.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
KafkaConsumer<String, String> consumer = new KafkaConsumer<String, String>(config);
TopicPartition tp = new TopicPartition("connect-test", 0);
List<TopicPartition> ltp = Arrays.asList(tp);
consumer.assign(ltp);
consumer.seekToEnd(ltp);
ConsumerRecords<String, String> records;
while(true){
records = consumer.poll(1000);
for (ConsumerRecord<String, String> record : …Run Code Online (Sandbox Code Playgroud)