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