小编Val*_*riy的帖子

Apache Kafka Consumer 停止消费消息

我有一个卡夫卡消费者的问题。我使用新的 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)

java apache-kafka kafka-consumer-api

6
推荐指数
0
解决办法
8829
查看次数

标签 统计

apache-kafka ×1

java ×1

kafka-consumer-api ×1