当Source包含大量记录时Akka Streams无法运行的问题排查
我刚遇到过类似的Akka Streams问题,结合你描述的场景,咱们一步步拆解问题:
首先明确你的核心场景:你正在编写Akka Streams入门示例,用整数范围作为Source,通过PrimeStream类的filterPrimes Flow筛选质数输出。小范围(1000000010001000)测试正常,但把范围扩大到1000000010010000后,程序无异常无警告,直接退出且没有任何输出。非流式的测试代码却能正常运行并输出结果,而且移除mapConcat后能输出10000个整数,但替换filterPrimes为直接返回参数的方法时依旧无输出——甚至调试发现filterPrimes都没被调用。
先贴出你的核心代码方便参考:
PrimeStream类代码
import akka.NotUsed; import akka.actor.ActorSystem; import akka.stream.javadsl.Flow; import com.aparapi.Kernel; import com.aparapi.Range; import java.util.Arrays; import java.util.List; import java.util.stream.Collectors; public class PrimeStream { private final AverageRepository averageRepository = new AverageRepository(); private final ActorSystem actorSystem; public PrimeStream(ActorSystem actorSystem) { this.actorSystem = actorSystem; } public Flow<Integer, Integer, NotUsed> filterPrimes() { return Flow.of(Integer.class) .grouped(10000) .mapConcat(PrimeKernel::filterPrimes) .filter(v -> v != 0); } }
非流式测试代码
@Test public void testPlain() { List<Integer> in = IntStream.rangeClosed(10000000, 10010000) .boxed() .collect(Collectors.toList()); List<Integer> out = PrimeKernel.filterPrimes(in); System.out.println(out); }
问题根源分析
Akka Streams是惰性求值的流框架——只有当下游有需求时,上游才会推送数据。结合你的现象,最可能的几个原因:
测试代码未等待流处理完成
如果你测试Akka Streams的代码只是构建了流,但没有连接到能触发需求的Sink(比如Sink.foreach),或者没有等待流处理完成就结束了程序,JVM可能在流开始处理前就退出了。尤其是数据量变大后,流的初始化、批次处理需要更多时间,但程序没等就终止了,自然看不到输出。Aparapi异步执行未等待完成
你的PrimeKernel用了Aparapi做GPU加速,它的Kernel.execute()默认是异步的。如果filterPrimes方法没有等待Kernel执行完成就返回结果,那Akka Streams拿到的可能是未处理的全0数组,后续filter(v -> v !=0)会把所有元素过滤掉,导致无输出;甚至可能Kernel还没开始执行,程序就退出了。流的缓冲区或批次大小不匹配
grouped(10000)会把元素攒成10000个一组的批次,当总元素是10001个(10000000~10010000包含两端),会生成两个批次:10000个+1个。如果下游需求不足,或者缓冲区设置过小,第二个小批次可能无法被触发处理。
解决方案与排查步骤
1. 先确保测试代码正确等待流完成
这是最常见的新手坑!修改你的Akka Streams测试代码,连接到Sink并等待流处理结束:
@Test public void testAkkaStream() throws Exception { // 创建ActorSystem ActorSystem system = ActorSystem.create("PrimeStreamTest"); PrimeStream primeStream = new PrimeStream(system); // 构建流并运行,等待完成 Source.rangeClosed(10000000, 10010000) .via(primeStream.filterPrimes()) .runWith(Sink.foreach(System.out::println), system) .toCompletableFuture() .get(); // 阻塞等待流处理完成 // 优雅关闭ActorSystem system.terminate(); system.getWhenTerminated().toCompletableFuture().get(); }
如果之前的测试代码没有get()等待,JVM会在流刚开始初始化就退出,自然看不到任何输出。
2. 确保Aparapi Kernel执行完成
修改PrimeKernel.filterPrimes方法,等待GPU计算完成后再返回结果:
public static List<Integer> filterPrimes(List<Integer> numbers) { int[] nums = numbers.stream().mapToInt(Integer::intValue).toArray(); PrimeKernel kernel = new PrimeKernel(nums); // 执行并等待完成(Aparapi的execute默认异步,需要确保同步完成) kernel.execute(Range.create(nums.length)); kernel.getExecutionTime(); // 这一步会等待执行完成 // 转换为List返回 return Arrays.stream(nums) .boxed() .collect(Collectors.toList()); }
如果Kernel是异步执行的,mapConcat拿到的就是未处理的数组,后续过滤后没有有效元素,导致无输出。
3. 添加日志调试流的各个阶段
在filterPrimes Flow中添加日志,跟踪元素流动情况,确认哪个阶段出了问题:
public Flow<Integer, Integer, NotUsed> filterPrimes() { return Flow.of(Integer.class) .log("received-input", num -> "Got number: " + num) .grouped(10000) .log("created-batch", batch -> "Batch size: " + batch.size()) .mapConcat(PrimeKernel::filterPrimes) .log("after-filter", num -> "Post-filter value: " + num) .filter(v -> v != 0) .log("final-prime", num -> "Found prime: " + num); }
通过日志可以清楚看到:有没有元素进入流?批次是否生成?filterPrimes返回了什么?最终有没有质数输出?
4. 调整流的缓冲区或批次大小
如果是缓冲区不足导致元素积压,可以调整Flow的缓冲区设置:
public Flow<Integer, Integer, NotUsed> filterPrimes() { return Flow.of(Integer.class) .grouped(1000) // 减小批次大小,降低单次处理压力 .mapConcat(PrimeKernel::filterPrimes) .filter(v -> v != 0) .withAttributes(ActorAttributes.withInputBuffer(2048, 2048)); // 增大缓冲区 }
小批次处理可以降低GPU的单次计算负载,也更容易触发下游的需求。
总结
优先检查测试代码是否等待流处理完成——这是Akka Streams新手最容易踩的坑。如果问题还存在,再排查Aparapi的执行是否同步,最后调整流的参数。按这个步骤应该能解决你的无输出问题。
内容的提问来源于stack exchange,提问作者Jeffrey Phillips Freeman

