Kotlin如何将Sequence或Stream转为对应固定数量明细的汇总序列
Kotlin Sequence 实现方案
Kotlin 标准库原生提供chunked()扩展函数满足需求,该函数会按传入的固定大小将Sequence惰性拆分为多个子集合,随后对子集合做汇总映射即可,完全符合你提到的「非一一映射、不合并为单个累加器、输出仍为Sequence」的特性。
用法示例如下:
// 示例固定24条小时明细汇总为1条日汇总 val hourlyInfo: Sequence<HourlyData> = ... val dailyInfo: Sequence<DailySummary> = hourlyInfo // 参数为每个汇总对应的明细数量 .chunked(24) // 可选:如果要求每个块必须满固定数量,丢弃不足的尾块 // .filter { it.size == 24 } .map { hourList -> // 此处实现明细转汇总的业务逻辑 DailySummary( date = hourList.first().date, totalFlow = hourList.sumOf { it.flow }, maxTemp = hourList.maxOf { it.temperature } ) }
chunked默认保持Sequence的惰性求值特性,不会一次性加载全量数据到内存,适合处理大数据量场景。
Java Stream 实现方案
Java 标准库未提供原生的Stream分块方法,你可以通过封装自定义Spliterator实现,通用工具方法及用法如下:
通用分块工具方法
import java.util.*; import java.util.function.Consumer; import java.util.stream.Stream; import java.util.stream.StreamSupport; public class StreamUtils { public static <T> Stream<List<T>> chunked(Stream<T> source, int chunkSize) { if (chunkSize <= 0) { throw new IllegalArgumentException("块大小必须大于0"); } Spliterator<T> sourceSpliterator = source.spliterator(); Spliterator<List<T>> chunkedSpliterator = new Spliterators.AbstractSpliterator<List<T>>( (sourceSpliterator.estimateSize() + chunkSize - 1) / chunkSize, sourceSpliterator.characteristics() & ~Spliterator.SIZED ) { private final Iterator<T> sourceIterator = Spliterators.iterator(sourceSpliterator); @Override public boolean tryAdvance(Consumer<? super List<T>> action) { if (!sourceIterator.hasNext()) { return false; } List<T> chunk = new ArrayList<>(chunkSize); for (int i = 0; i < chunkSize && sourceIterator.hasNext(); i++) { chunk.add(sourceIterator.next()); } action.accept(chunk); return true; } }; return StreamSupport.stream(chunkedSpliterator, source.isParallel()) .onClose(source::close); } }
用法示例
Stream<HourlyData> hourlyInfo = ... Stream<DailySummary> dailyInfo = StreamUtils.chunked(hourlyInfo, 24) // 可选:过滤不足24条的尾块 // .filter(hourList -> hourList.size() == 24) .map(hourList -> { // 实现明细汇总逻辑 LocalDate date = hourList.get(0).getDate(); int totalFlow = hourList.stream().mapToInt(HourlyData::getFlow).sum(); double maxTemp = hourList.stream().mapToDouble(HourlyData::getTemperature).max().orElse(0); return new DailySummary(date, totalFlow, maxTemp); });
如果项目依赖Guava,可直接调用Guava提供的Iterables.partition方法简化实现,无需自定义工具类。
内容的提问来源于stack exchange,提问作者Ted Hopp
相关产品推荐
相关产品推荐

