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

如何获取已完成线程的结果列表?ThreadUnion实现问题排查

问题描述

需要实现ThreadUnion接口的两个方法:

  • Thread newThread(Runnable runnable):创建并注册线程,命名格式为[前缀]-worker-n(n为线程序号),需监控线程执行状态;
  • List<FinishedThreadResult> results():返回已完成线程的结果列表,包含线程名、结束时间戳、异常(若有),要求线程安全。

用户自行实现后遇到两个问题:

  1. 部分场景下时间戳不符合测试预期(结束时间不在指定区间内);
  2. 偶尔抛出java.util.ConcurrentModificationException。

方法定义

Thread newThread (Runnable runnable)
创建并注册新线程,线程名格式为-worker-n(n是线程编号)。ThreadUnion必须监控创建线程的执行状态,具体逻辑参考results()方法。

List results()
返回已完成线程的结果列表,未完成的线程不能出现在结果中。每个结果需包含线程名、执行结束时间戳,以及线程抛出的异常(如果有)。

注意:ThreadUnion的实现必须是线程安全的,支持并发创建线程。

public FinishedThreadResult(final String threadName) {
    this(threadName, null);
}

public FinishedThreadResult(final String threadName, final Throwable throwable) {
    this.threadName = threadName;
    this.throwable = throwable;
    this.finished = LocalDateTime.now();
}

public String getThreadName() {
    return threadName;
}

public LocalDateTime getFinished() {
    return finished;
}

public Throwable getThrowable() {
    return throwable;
}        

测试代码

void testResults() {
    ...
    final ThreadUnion threadUnion = ThreadUnion.newInstance(unionName);

    final Set<RuntimeException> exceptionsToThrow = ...;
    final int fineThreadsCount = 3;
    final int longThreadsCount = 2;
    final int shortThreadsCount = fineThreadsCount + exceptionsToThrow.size();
    final int allResultsCount = fineThreadsCount + longThreadsCount + exceptionsToThrow.size();

    final List<Thread> exceptionThreads = exceptionsToThrow.stream()
            .map(e -> threadUnion.newThread(() -> throwException(e)))
            .collect(toList());

    final List<Thread> fineThreads = Stream.generate(() -> threadUnion.newThread(this::printThreadName))
            .limit(fineThreadsCount)
            .collect(toList());

    CountDownLatch latch = new CountDownLatch(1);
    final List<Thread> longThreads = Stream.generate(() -> threadUnion.newThread(() -> awaitLatch(latch)))
            .limit(longThreadsCount)
            .collect(toList());

    final LocalDateTime beforeStart = LocalDateTime.now();
    final LocalDateTime afterStart;
    final LocalDateTime afterShortThreadsFinish;
    final LocalDateTime afterLongThreadsFinish;

    final List<FinishedThreadResult> shortResults;

    try {
        startThreads(exceptionThreads, fineThreads, longThreads);
        afterStart = LocalDateTime.now();

        joinThreads(exceptionThreads, fineThreads);
        afterShortThreadsFinish = LocalDateTime.now();

        shortResults = threadUnion.results();
        assertEquals(shortThreadsCount, shortResults.size());
        assertEquals(expectedNames(unionName, 0, shortThreadsCount), resultThreadNames(shortResults));
        // 测试失败点
        shortResults.stream()
                .map(FinishedThreadResult::getFinished)
                .forEach(finished -> assertAll(
                        () -> {
                            System.out.println("before "+beforeStart);
                            System.out.println("thread " + finished.toString());
                            assertFalse(finished.isBefore(beforeStart));
                        },
                        () -> {
                            System.out.println("after "+afterShortThreadsFinish);
                            assertFalse(finished.isAfter(afterShortThreadsFinish));
                        }
                ));

        assertEquals(exceptionsToThrow, collectThrowables(shortResults));  

用户实现代码

private final String nameThread;
private final Collection<Thread> syncResultsThreadList = Collections.synchronizedCollection(new ArrayList<>());
private final Collection<FinishedThreadResult> resultList = Collections.synchronizedCollection(new ArrayList<>());
private final AtomicInteger numberThreads = new AtomicInteger(0);
private final Map<Thread, FinishedThreadResult> resultMap = Collections.synchronizedMap(new HashMap<>());
...
@Override
public Thread newThread(Runnable r) {

    Thread thread = new Thread( r, nameThread + "-worker-" + numberThreads.get());
    numberThreads.incrementAndGet();

    syncResultsThreadList.add(thread);
    thread.setUncaughtExceptionHandler((t, throwable) -> {
            resultList.add(new FinishedThreadResult(thread.getName(), throwable));
        });
    return thread;
 }

@Override
public List<FinishedThreadResult> results() {
     List<FinishedThreadResult> results = new ArrayList<>();
    for (Thread tr : syncResultsThreadList) {
            resultMap.put(tr, new FinishedThreadResult(tr.getName()));
        }

    if (resultList.size() > 0) {
        for (FinishedThreadResult ftr : resultList) {
            resultMap.forEach((t, f) -> {
               if (Objects.equals(f.getThreadName(), ftr.getThreadName())) {
                    resultMap.replace(t, f, ftr);
                }
            });
        }
    }

    resultMap.forEach((t, f) -> {
        if (!t.isAlive()) {
            results.add(f);
       }
    });

    return Collections.unmodifiableList(results);
}

问题根源与修复方案

问题1:时间戳不符合预期的原因

用户的实现中,正常线程的FinishedThreadResult是在调用results()方法时才创建的,此时的时间戳是调用results()的时间,而非线程实际结束的时间。而异常线程的结果是在UncaughtExceptionHandler中创建的(时间正确),这就导致正常线程的时间戳晚于线程实际结束时间,超出了测试的afterShortThreadsFinish区间。

问题2:ConcurrentModificationException的原因

虽然用了Collections.synchronizedCollection包装集合,但在results()遍历集合的同时,可能有新线程通过newThread()添加到集合中,遍历过程中集合结构被修改,直接触发并发修改异常。另外遍历resultList时如果有新异常结果被添加,也会引发同样问题。

修复后的实现

import java.time.LocalDateTime;
import java.util.ArrayList;
import java.util.Collections;
import java.util.List;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.atomic.AtomicInteger;

public class ThreadUnionImpl implements ThreadUnion {
    private final String unionName;
    private final AtomicInteger threadCounter = new AtomicInteger(0);
    private final ConcurrentMap<Thread, FinishedThreadResult> threadResults = new ConcurrentHashMap<>();

    public ThreadUnionImpl(String unionName) {
        this.unionName = unionName;
    }

    @Override
    public Thread newThread(Runnable runnable) {
        int threadNum = threadCounter.getAndIncrement();
        String threadName = unionName + "-worker-" + threadNum;
        Thread thread = new Thread(() -> {
            FinishedThreadResult result;
            try {
                runnable.run();
                // 线程正常结束时立即记录时间
                result = new FinishedThreadResult(threadName);
            } catch (Throwable t) {
                // 捕获异常时立即记录时间和异常信息
                result = new FinishedThreadResult(threadName, t);
                throw t;
            } finally {
                // 无论正常/异常结束,都把结果存入map
                threadResults.put(thread, result);
            }
        }, threadName);

        // 兜底处理未捕获的异常
        thread.setUncaughtExceptionHandler((t, throwable) -> {
            threadResults.put(t, new FinishedThreadResult(t.getName(), throwable));
        });
        return thread;
    }

    @Override
    public List<FinishedThreadResult> results() {
        List<FinishedThreadResult> finishedResults = new ArrayList<>();
        // 遍历并发map,筛选已结束的线程结果
        for (var entry : threadResults.entrySet()) {
            if (!entry.getKey().isAlive()) {
                finishedResults.add(entry.getValue());
            }
        }
        return Collections.unmodifiableList(finishedResults);
    }
}

修复说明

  1. 时间戳问题解决:

    • 在线程执行的finally块中创建FinishedThreadResult,确保时间戳是线程实际结束的时间,而非results()的调用时间。
    • 异常场景下,无论是捕获的异常还是未捕获的异常,都在异常发生时立即记录时间戳。
  2. 并发修改异常解决:

    • 用ConcurrentHashMap替代同步集合,它的遍历是弱一致性的,不会因为并发修改抛出异常。
    • 移除了多集合同步的复杂逻辑,只用一个ConcurrentMap存储线程结果,避免多集合操作的并发冲突。
    • 线程创建时直接绑定结果记录逻辑,results()只做读取筛选,不修改集合结构,彻底避免遍历与修改的冲突。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.27 07:14:59