【发布时间】:2016-08-29 00:12:28
【问题描述】:
我需要在晚上执行一个作业,它将获取 kafka 队列中的所有消息并使用它们执行一个进程。我能够收到消息,但 kafka 流正在等待更多消息,我无法继续我的流程。我有以下代码:
...
private ConsumerConnector consumerConnector;
private final static String TOPIC = "test";
public MessageStreamConsumer() {
Properties properties = new Properties();
properties.put("zookeeper.connect", "localhost:2181");
properties.put("group.id", "test-group");
ConsumerConfig consumerConfig = new ConsumerConfig(properties);
consumerConnector = Consumer.createJavaConsumerConnector(consumerConfig);
}
public List<String> getMessages() {
Map<String, Integer> topicCountMap = new HashMap<String, Integer>();
topicCountMap.put(TOPIC, new Integer(1));
Map<String, List<KafkaStream<byte[], byte[]>>> consumerMap = consumerConnector
.createMessageStreams(topicCountMap);
KafkaStream<byte[], byte[]> stream = consumerMap.get(TOPIC).get(0);
ConsumerIterator<byte[], byte[]> it = stream.iterator();
List<String> messages = new ArrayList<>();
while (it.hasNext())
messages.add(new String(it.next().message()));
return messages;
}
代码能够获取消息,但是当它处理最后一条消息时,它会留在行中:
while (it.hasNext())
问题是,我怎样才能从 kafka 获取所有消息,停止流并继续我的其他任务。
希望你能帮到我
谢谢
【问题讨论】:
-
但我不认为这样做的最佳做法是等到抛出异常。如果我的进程花费的时间超过了配置的超时时间怎么办
-
你不应该在这里使用直接的 Kafka Consumer,而不是 KafkaStream 吗?流自然会保持活力。
标签: java apache-kafka kafka-consumer-api