Java8如何获取Stream大小同时正常返回原始流?
问题说明
Java Stream执行count()这类终端操作后,流管线会自动关闭、数据源被消费完毕,无法二次使用,这是Stream设计上的明确约束,不存在可以直接复用已关闭流的取巧方法。结合Kafka流式传输学生数据+同步总条数的场景,有三个可落地的实现方案,根据数据量和一致性要求选择即可。
实现方案
方案1:单次流遍历边传输边计数(生产环境优先推荐)
这个方案无内存溢出风险、无数据一致性偏差,适合任意规模的数据集:
- 不提前执行count操作,直接遍历数据库返回的学生流
- 遍历过程中用原子类累计记录数,同时逐条把学生数据发往Kafka
- 等整个流遍历完成后,再把累计得到的总条数作为单独消息发给消费者,同时可以携带流传输结束的标识
示例代码:
// 原子类保证计数线程安全 AtomicLong totalCount = new AtomicLong(0); // 先发送流开始标记,告知消费者准备接收数据 kafkaTemplate.send("student_topic", "STREAM_START", ""); // 用try-with-resources保证数据库流连接正常释放,避免连接泄漏 try(Stream<Student> allStudent = studentRepo.findAll()) { allStudent.forEach(student -> { totalCount.incrementAndGet(); // 序列化学生实体发送到Kafka kafkaTemplate.send("student_topic", student.getId().toString(), JSON.toJSONString(student)); }); } // 流消费完成后发送总条数和结束标记 kafkaTemplate.send("student_topic", "STREAM_END", String.valueOf(totalCount.get()));
方案2:单独查询总条数后再获取流
如果学生表数据量不大、count查询走主键索引性能很高,且能容忍短时间窗口的数据一致性偏差,可以直接分两次数据库操作,逻辑最简洁:
// 单独执行count查询,和后续的流查询互不影响 long totalCount = studentRepo.count(); // 先把总条数发给消费者 kafkaTemplate.send("student_topic", "TOTAL_COUNT", String.valueOf(totalCount)); // 再查询学生流用于传输 Stream<Student> allStudent = studentRepo.findAll(); return allStudent;
注意:两次数据库查询之间如果有学生数据新增/删除,会出现总条数和实际传输记录数不一致的问题,一致性要求高的场景不要用
方案3:内存缓存流数据后生成新流
如果学生数据量很小(万级以内,JVM内存可以完全承载),可以先把流中的数据收集到内存集合,先计算集合大小得到总条数,再从集合生成新的流返回:
// 把流数据收集到内存List List<Student> studentCache = studentRepo.findAll().toList(); long totalCount = studentCache.size(); // 发送总条数 kafkaTemplate.send("student_topic", "TOTAL_COUNT", String.valueOf(totalCount)); // 从内存集合生成新流返回 return studentCache.stream();
注意:数据量超过10万条不要用这个方案,全量数据加载到内存很容易触发OOM
内容的提问来源于stack exchange,提问作者Ajit Barik
相关产品推荐
相关产品推荐

