Java中distinct().limit(3).forEach()引发线程阻塞问题排查
将服务部署在ECS Fargate环境中,一段包含distinct().limit(3).forEach()的流操作链间歇性导致线程阻塞,进而使容器健康检查失败,处于不健康状态。
相关代码
@Builder(toBuilder = true) @ToString @Getter @EqualsAndHashCode public class TestPojo { private String pojoId; private String score; @EqualsAndHashCode.Exclude private boolean isRandomOrdered; } private List<TestPojo> randomizeOrder(final List<TestPojo> testObjects) { // Picking testObjects randomly only from top 25 final int topLimit = Math.min(testObjects.size(), 25); final Set<TestPojo> newOrderedPojos = new LinkedHashSet<>(); // First 2 position are fixed newOrderedPojos.add(testObjects.get(0)); // Validations are present in calling function to make sure that get(0) and get(1) don't throw out of bound error. newOrderedPojos.add(testObjects.get(1)); new Random().ints(2, topLimit).distinct().limit(3) .forEach(i -> newOrderedPojos.add(testObjects.get(i).toBuilder() .isRandomOrdered(true) .build())); // Add remaining objects. Since this is a set only objects not already present will be inserted. newOrderedPojos.addAll(asins); return ImmutableList.copyOf(newOrderedPojos); }
报错栈信息
11 Jul 2023 16:19:12,942 [33m[WARN][m (abc-periodic-metrics-0) com.abc.xyz.bobcat.StuckThreadDetectingValve: Thread Bobcat-2 is detected as stuck as it has been blocked for 4442432 ms com.abc.xyz.bobcat.StuckThreadDetectingValve$StuckThreadException: Check stacktrace to see where the thread is stuck at java.util.HashMap.hash(HashMap.java:340) at java.util.HashMap.containsKey(HashMap.java:597) at java.util.HashSet.contains(HashSet.java:204) at java.util.stream.DistinctOps$1$2.accept(DistinctOps.java:173) at java.util.stream.IntPipeline$4$1.accept(IntPipeline.java:250) at java.util.Random$RandomIntsSpliterator.tryAdvance(Random.java:1029) at java.util.stream.IntPipeline.forEachWithCancel(IntPipeline.java:162) at java.util.stream.AbstractPipeline.copyIntoWithCancel(AbstractPipeline.java:499) at java.util.stream.AbstractPipeline.copyInto(AbstractPipeline.java:486) at java.util.stream.AbstractPipeline.wrapAndCopyInto(AbstractPipeline.java:472) at java.util.stream.ForEachOps$ForEachOp.evaluateSequential(ForEachOps.java:150) at java.util.stream.ForEachOps$ForEachOp$OfInt.evaluateSequential(ForEachOps.java:188) at java.util.stream.AbstractPipeline.evaluate(AbstractPipeline.java:234) at java.util.stream.IntPipeline.forEach(IntPipeline.java:427) at com.abc.service.xyz.randomizeAsinOrder(TestClass.java:123)
查阅了distinct、limit和forEach的Java官方文档,未发现提及线程不安全或此类阻塞场景的说明,请问该问题可能的原因是什么?
随机数范围过小导致distinct()无限循环
从报错栈可见,线程卡在随机数生成与流的去重逻辑中。当testObjects数量较少时,topLimit = Math.min(testObjects.size(),25)的取值会很小,比如testObjects.size()为4时,随机数范围是[2,4),仅包含2、3两个可选值。此时调用distinct().limit(3),流会持续尝试生成新的不重复随机数,但可选范围内根本凑不够3个不同数值,导致流进入无限循环,线程长期阻塞。流的终止条件设计缺陷是核心问题
报错栈中的HashSet操作是distinct()内部维护的去重集合,并非代码中的newOrderedPojos。真正的问题在于流逻辑未考虑随机数候选池大小小于所需去重数量的场景,导致limit(3)永远无法满足,线程一直卡在循环生成随机数的过程中。ECS Fargate环境触发健康检查失败
本地测试时testObjects数量通常足够大(远大于25),不会出现候选数不足的情况,因此问题间歇性出现。而在ECS Fargate容器中,线程长期阻塞会触发健康检查机制,判定服务处于不健康状态。
内容的提问来源于stack exchange,提问作者Belphegor21

