You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

测试类中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就会异常退出。


问题原因

核心是两个机制共同导致的:

  1. 你配置线程工厂时显式开启了setDaemon(true),线程池里所有工作线程都是守护线程。JVM的退出规则是只要所有用户线程执行完毕,不管守护线程有没有在跑任务,都会直接终止进程,不会等守护线程的任务执行完成。
  2. 测试框架的运行逻辑和main方法不同:main方法默认的主线程是用户线程,如果你在main里启动线程池后没有立刻退出main,JVM会保持存活;但JUnit、TestNG这类测试框架,执行完所有测试用例的逻辑后,会直接结束测试进程,不会等待后台守护线程运行。
    你看到代码卡在getRecords调用时终止,本质是这行是Kinesis的IO阻塞调用,耗时较长,测试主线程跑完initialize()方法后就直接退出了JVM,把还在等待IO返回的守护线程直接杀掉,不是业务代码本身抛异常导致的终止。
解决方法
  • 如果是生产环境代码:去掉.setDaemon(true)配置,把线程池工作线程设为用户线程,同时在应用关闭逻辑里显式调用executor.shutdown(),配置合理的等待时间让已提交的任务执行完,避免进程退出时任务被强制中断。
  • 如果是测试场景验证消费逻辑:要么同样把线程改成用户线程,要么在测试方法里加阻塞逻辑(比如CountDownLatch、足够时长的Thread.sleep())卡住测试主线程,等你需要验证的逻辑执行完成后再结束测试,不要让测试主线程跑完初始化就直接退出。

内容的提问来源于stack exchange,提问作者tabiodun

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.30 09:09:20