【发布时间】:2016-06-29 16:54:33
【问题描述】:
我有一个队列(它恰好是 Kafka,但我不确定这是否重要),我正在从中读取消息。我想创建一个流来表示这些数据。
我使用(Kafka)队列的伪代码如下所示:
List<Message> messages = new ArrayList<>();
while (true) {
ConsumerRecords<String, Message> records = kafkaConsumer.poll(100);
messages.add(recordsToMessages(records));
if (x) {
break;
}
}
return messages.stream();
使用此伪代码,直到 while 循环被破坏,即直到所有队列都被读取后,才会返回流。
我希望能够立即返回流,即可以将新消息添加到流在返回之后。
我觉得我需要使用 Stream.generate 但我不确定如何使用,或者我需要一个拆分器?
我还想稍后在代码中关闭流。
谢谢!
【问题讨论】:
-
你不能使用 do-while 循环吗?
-
不幸的是,这只会运行一次循环。一旦退出 while 循环,就不会再向流中添加任何值了。
标签: java java-8 apache-kafka java-stream kafka-consumer-api