【问题标题】:Kafka consumer api failed to subscribe to topicKafka 消费者 api 订阅主题失败
【发布时间】:2017-07-31 03:12:56
【问题描述】:

我正在使用简单的 Kafka 客户端 API。据我所知,消费者消息有两种方式,订阅主题和分配分区给消费者。

但是第一种方法不起作用。消费者poll() 将永远挂起。它仅适用于assign

    // common config for consumer
    Map<String, Object> config = new HashMap<>();
    config.put("bootstrap.servers", bootstrap);

    config.put("group.id", KafkaTestConstants.KAFKA_GROUP);
    config.put("enable.auto.commit", "true");
    config.put("auto.offset.reset", "earliest");
    config.put("key.deserializer", StringDeserializer.class.getName());
    config.put("value.deserializer", StringDeserializer.class.getName());
    StringDeserializer deserializer = new StringDeserializer();
    KafkaConsumer<String, String> consumer = new KafkaConsumer<>(config, deserializer, deserializer);

    // subscribe does not work, poll() hangs
    consumer.subscribe(Arrays.asList(KafkaTestConstants.KAFKA_TOPIC));

这是分配分区的代码。

    // assign works
    TopicPartition tp = new TopicPartition(KafkaTestConstants.KAFKA_TOPIC, 0);
    List<TopicPartition> tps = Arrays.asList(tp);
    consumer.assign(tps);

因为我想使用自动提交功能,根据这个post,它应该只适用于消费者组管理。为什么subscribe() 不起作用?

【问题讨论】:

  • 不能同时使用带assign(Collection)的手动分区分配和带subscribe的组分配。
  • @amethystic 我并不是要让两者都起作用。只是想让subscribe 工作就可以了。
  • 你能删除consumer.assign(tps)然后重试吗?
  • @amethystic 很抱歉造成混乱。最后两个代码块不同时存在。只是想说subscribe 方法不起作用,而最后三行起作用
  • 你在服务器端和客户端使用什么版本的Kafka?并且你也可以尝试一个新的 group.id 重新运行程序,检查服务端/客户端是否有异常抛出?

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


【解决方案1】:

我遇到了同样的问题。 我使用的是 kafka_2.12 jar 版本,当我将其降级到 kafka_2.11 时它可以工作。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2017-05-16
    • 2018-01-19
    • 2015-08-16
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多