【问题标题】:Kafka consumer returns no recordsKafka 消费者不返回任何记录
【发布时间】:2019-05-11 09:30:48
【问题描述】:

我正在尝试使用 Kafka 制作一个小型 PoC。但是,在 java 中创建消费者时,此消费者没有收到任何消息。即使当我使用相同的 url/topic 启动 kafka-console-consumer.sh 时,我也会收到消息。有谁知道我可能做错了什么?此代码由 GET API 调用。

public List<KafkaTextMessage> receiveMessages() {
    log.info("Retrieving messages from kafka");
    val props = new Properties();
    // See https://kafka.apache.org/documentation/#consumerconfigs
    props.put("bootstrap.servers", "my-cluster-kafka-bootstrap:9092");
    //props.put("client.id", "my-topic consumer");
    props.put("group.id", "test");
    props.put("enable.auto.commit", "false");
    props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
    props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
    ImmutableList.Builder<KafkaTextMessage> builder;
    try (KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props)) {
        consumer.subscribe(Collections.singletonList(TEXT_MESSAGE_TOPIC));
        ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
        builder = ImmutableList.builder();
        for (ConsumerRecord<String, String> record : records) {
            builder.add(new KafkaTextMessage(record.value()));
            log.info("We got at position: {} key:{} value: {}", record.offset(), record.key(), record.value());
            consumer.commitSync();
        }
    }
    return builder.build();
}

【问题讨论】:

  • 你使用的是什么版本的kafka?
  • 更改此 group.id 并重试props.put("group.id", "test");
  • 您确定订阅了正确的主题吗?
  • 100ms的轮询时间够吗?
  • 您找到解决方案了吗?我也有同样的问题。一定要喜欢无声的失败。

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


【解决方案1】:

尝试在您的消费者属性中添加auto.offset.reset=earliest。默认值设置为latest。我建议这样做是因为我看到您的 group.id 设置为 test,您可能已经在之前的测试中使用了该值。

【讨论】:

    猜你喜欢
    • 2017-09-22
    • 1970-01-01
    • 1970-01-01
    • 2018-06-20
    • 2019-11-21
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多