测试类中ExecutorService意外提前终止问题成因及方案咨询
问题描述
我尝试使用多个实现ExecutorService接口的类,均出现执行器提前退出的相同问题,相关代码如下:
public class StreamConsumer { private static final Logger LOGGER = LogManager.getLogger(MethodHandles.lookup().lookupClass()); final Environment environment; final String streamName; final String shardId; public StreamConsumer(Environment environment, String streamName, String shardId) { this.environment = environment; this.streamName = streamName; this.shardId = shardId; } public void initialize() { try { final ListShardsRequest listShardsRequest = new ListShardsRequest().withStreamName(streamName); final ListShardsResult listShardsResult = environment.kinesis().listShards(listShardsRequest); final List<Shard> shards = listShardsResult.getShards(); final ThreadFactory threadFactory = new ThreadFactoryBuilder() .setDaemon(true) .setNameFormat(streamName + "-%d") .build(); final ScheduledExecutorService executor = Executors.newScheduledThreadPool(shards.size(), threadFactory); shards.stream().map(ShardConsumer::new).forEach(consumer -> consumer.start(executor)); } catch (Exception ex) { LOGGER.warn("Unable to start Stream Consumer of {}", streamName, ex); } } private class ShardConsumer { private final String shardId; private String shardIterator; ShardConsumer(Shard shard) { this.shardId = shard.getShardId(); } private String getShardIterator() { GetShardIteratorRequest shardIteratorRequest = new GetShardIteratorRequest() .withStreamName(streamName) .withShardId(shardId) .withShardIteratorType(ShardIteratorType.LATEST.name()); return environment.kinesis().getShardIterator(shardIteratorRequest).getShardIterator(); } void start(ScheduledExecutorService executor) { this.shardIterator = getShardIterator(); executor.scheduleAtFixedRate(() -> { try { LOGGER.info("Started consuming stream: {} with shardId: {}", streamName, shardId); GetRecordsRequest request = new GetRecordsRequest().withShardIterator(shardIterator); final GetRecordsResult records = environment.kinesis().getRecords(request); UserRecord.deaggregate(records.getRecords()).forEach(x -> System.out.println(x)); shardIterator = records.getNextShardIterator(); } catch (Throwable e) { LOGGER.error("Problem consuming records from {} {}", streamName, shardId, e); } }, 0, 1_000, TimeUnit.MILLISECONDS); } }
代码执行到如下行时会意外终止:final GetRecordsResult records = environment.kinesis().getRecords(request);
我还尝试创建另一个ExecutorService执行如下无限循环逻辑:
while (true) { System.out.println(x) }
该任务同样会提前意外终止,暂未定位到问题根因。
补充排查结果:在Main方法中创建的ExecutorService可以无限期正常存活,但从测试类启动时ExecutorService就会异常退出。
问题原因
核心是两个机制共同导致的:
- 你配置线程工厂时显式开启了
setDaemon(true),线程池里所有工作线程都是守护线程。JVM的退出规则是只要所有用户线程执行完毕,不管守护线程有没有在跑任务,都会直接终止进程,不会等守护线程的任务执行完成。 - 测试框架的运行逻辑和main方法不同:main方法默认的主线程是用户线程,如果你在main里启动线程池后没有立刻退出main,JVM会保持存活;但JUnit、TestNG这类测试框架,执行完所有测试用例的逻辑后,会直接结束测试进程,不会等待后台守护线程运行。
你看到代码卡在getRecords调用时终止,本质是这行是Kinesis的IO阻塞调用,耗时较长,测试主线程跑完initialize()方法后就直接退出了JVM,把还在等待IO返回的守护线程直接杀掉,不是业务代码本身抛异常导致的终止。
解决方法
- 如果是生产环境代码:去掉
.setDaemon(true)配置,把线程池工作线程设为用户线程,同时在应用关闭逻辑里显式调用executor.shutdown(),配置合理的等待时间让已提交的任务执行完,避免进程退出时任务被强制中断。 - 如果是测试场景验证消费逻辑:要么同样把线程改成用户线程,要么在测试方法里加阻塞逻辑(比如
CountDownLatch、足够时长的Thread.sleep())卡住测试主线程,等你需要验证的逻辑执行完成后再结束测试,不要让测试主线程跑完初始化就直接退出。
内容的提问来源于stack exchange,提问作者tabiodun
相关产品推荐
相关产品推荐

