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

同步场景下ArrayList.removeFirst()偶发越界,求优化单元测试

线程池偶发数组越界异常的复现与定位方案

问题描述

我实现了一个支持动态添加任务的线程池类DepleatingFiFoThreadPool,代码如下:

package xxx;

import java.util.ArrayList;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.function.Consumer;
import java.util.logging.Level;
import java.util.logging.Logger;

/**
 * 可动态添加任务的线程池,执行过程中可以继续添加新任务
 * 
 * @param <A> 任务类型
 */
public class DepleatingFiFoThreadPool<A> {

    public static final Logger LOG = Logger.getLogger(DepleatingFiFoThreadPool.class.getCanonicalName());

    public record Socket<A>(A r, String threadPostfix) {
    }

    private final Thread[] threadsRunning;
    private final ArrayList<Socket<A>> unstartedRunnables = new ArrayList<>();
    private int nextFreeSocketIndex = 0;
    private final Consumer<Throwable> errorHandler;
    private final Object lock = new Object();
    private boolean waitingForLock = false;

    private final String prefix;
    private final Consumer<A> invoker;

    /**
     * 创建线程池
     * 
     * @param threadsRunningMax 最大并发线程数,必须大于0
     * @param errorHandler      异常处理器,不能为null
     * @param prefix            线程名前缀,用于调试,不能为null
     * @param invoker           任务执行器,不能为null
     */
    public DepleatingFiFoThreadPool(final int threadsRunningMax, final Consumer<Throwable> errorHandler,
            final String prefix, final Consumer<A> invoker) {
        this.errorHandler = errorHandler;
        this.prefix = prefix;
        this.invoker = invoker;
        this.threadsRunning = new Thread[threadsRunningMax];
    }

    /**
     * 添加并启动一个任务
     * 
     * 调用{@link #executeUntilDeplated(long)}返回后添加的任务不会被执行
     * 
     * @param notRunningThread 待执行的任务
     * @param threadPostfix    线程名后缀,不能为null
     */
    public void addAndStartThread(final A notRunningThread, final String threadPostfix) {
        var mustCallLater = false;
        if (nextFreeSocketIndex < threadsRunning.length) {
            synchronized (this) {
                if (nextFreeSocketIndex < threadsRunning.length) {
                    start(notRunningThread, threadPostfix);
                } else {
                    mustCallLater = true;
                }
            }
        } else {
            mustCallLater = true;
        }
        if (mustCallLater) {
            unstartedRunnables.add(new Socket<>(notRunningThread, threadPostfix));
        }
    }

    private void finished() {
        synchronized (this) {
            threadsRunning[--nextFreeSocketIndex] = null;
            if (nextFreeSocketIndex == 0) {
                synchronized (lock) {
                    if (waitingForLock) {
                        lock.notify();
                    }
                }
            } else {
                Socket<A> e = null;
                if (!unstartedRunnables.isEmpty()) {
                    synchronized (unstartedRunnables) {
                        if (!unstartedRunnables.isEmpty()) {
                            e = unstartedRunnables.removeFirst();
                        }
                    }
                }
                if (e != null) {
                    start(e.r, e.threadPostfix);
                }
            }
        }
    }

    private void start(final A notRunningThread, final String name) {
        var socket = nextFreeSocketIndex++;
        Thread t = new Thread(new Runnable() {
            @Override
            public void run() {
                try {
                    invoker.accept(notRunningThread);
                } catch (Throwable e) { // 异常屏障
                    try {
                        errorHandler.accept(e);
                    } catch (Throwable ta) {
                        ta.addSuppressed(e);
                        LOG.log(Level.SEVERE, ta.getMessage(), ta);
                    }
                }
                finished();
            }
        }, prefix + name);
        threadsRunning[socket] = t;
        t.start();
    }

    /**
     * 执行所有任务,包括执行过程中新增的任务
     * 
     * @param timeoutMs 超时时间(毫秒),应为正数
     * @return 超时返回false,否则返回true
     * @throws InterruptedException 当前线程被中断时抛出
     */
    public boolean executeUntilDeplated(final long timeoutMs) throws InterruptedException {
        AtomicBoolean resultHolder = new AtomicBoolean(true);
        if (nextFreeSocketIndex > 0) {
            Thread timeout = new Thread(() -> {
                try {
                    Thread.sleep(timeoutMs);
                    resultHolder.set(false);
                    synchronized (lock) {
                        lock.notify();
                    }
                } catch (InterruptedException e) {
                    LOG.log(Level.FINE, "未触发超时,无需看门狗", e);
                }
            });
            timeout.start();
            if (nextFreeSocketIndex > 0) {
                waitingForLock = true;
                synchronized (lock) {
                    try {
                        lock.wait(timeoutMs);
                    } finally {
                        timeout.interrupt();
                    }
                }
            }
        }
        return resultHolder.get();
    }
}

运行时偶发以下异常,每3-5周出现一次,难以调试:

Exception in thread "NN/SQL-caller" java.lang.ArrayIndexOutOfBoundsException: arraycopy: last source index 50 out of bounds for object array[49]
  at java.base/java.lang.System.arraycopy(Native Method)
  at java.base/java.util.ArrayList.fastRemove(ArrayList.java:724)
  at java.base/java.util.ArrayList.removeFirst(ArrayList.java:573)
  at xxx.DepleatingFiFoThreadPool.finished(DepleatingFiFoThreadPool.java:74)
  at xxx.DepleatingFiFoThreadPool$1.run(DepleatingFiFoThreadPool.java:100)
  at java.base/java.lang.Thread.run(Thread.java:1570)

我尝试用以下单元测试复现问题,但未能成功:

public void testMassiveQueuedFunction() throws InterruptedException {
    var dtp = new DepleatingFiFoThreadPool<Runnable>(100, x -> System.out.println(x), "testQueueFunction",
            Runnable::run);
    var l = new CountDownLatch(0);
    Runnable wait = () -> {
        try {
            Thread.sleep(2);
            l.countDown();
        } catch (InterruptedException e) {
        }
    };
    for (int i = 0; i < 9999; i++) {
        dtp.addAndStartThread(wait, i + "");
    }
    long start = System.currentTimeMillis();
    dtp.executeUntilDeplated(3000);
    long took = System.currentTimeMillis() - start;
    assert l.getCount() == 0 : "期望计数为0,实际为: " + l.getCount();
    assert took < 2900 : "执行耗时过长,实际耗时: " + took;
}

需要更有效的单元测试方案来复现并定位该偶发异常。


问题根源分析

异常出现在ArrayList.removeFirst()时的arraycopy操作,说明并发修改导致了ArrayList内部结构的不一致。虽然finished()方法中对unstartedRunnables做了同步,但addAndStartThread()方法中直接调用unstartedRunnables.add()时没有加锁——当多个线程同时调用addAndStartThread()添加任务,同时有线程在finished()中移除任务时,ArrayList的内部数组会出现并发修改冲突,最终触发越界异常。


针对性单元测试方案

1. 高并发添加任务测试

模拟多个线程同时向线程池添加任务,同时让已有任务快速完成并从队列取任务,制造并发读写场景:

import org.junit.jupiter.api.Test;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.TimeUnit;

public class DepleatingFiFoThreadPoolTest {

    @Test
    public void testConcurrentAddAndFinish() throws InterruptedException {
        int poolSize = 50;
        int taskCount = 10000;
        int addThreadCount = 20;

        CountDownLatch allTasksDone = new CountDownLatch(taskCount);
        DepleatingFiFoThreadPool<Runnable> dtp = new DepleatingFiFoThreadPool<>(
                poolSize,
                Throwable::printStackTrace,
                "ConcurrentTest-",
                Runnable::run
        );

        // 多线程并发添加任务
        ExecutorService addExecutor = Executors.newFixedThreadPool(addThreadCount);
        for (int i = 0; i < taskCount; i++) {
            int taskId = i;
            addExecutor.submit(() -> {
                dtp.addAndStartThread(() -> {
                    // 任务快速完成,立即触发finished()
                    allTasksDone.countDown();
                }, "-task" + taskId);
            });
        }
        addExecutor.shutdown();
        addExecutor.awaitTermination(1, TimeUnit.MINUTES);

        // 等待所有任务执行完毕
        boolean completed = dtp.executeUntilDeplated(10000);
        assert completed : "任务执行超时";
        assert allTasksDone.getCount() == 0 : "仍有未完成的任务";
    }
}

2. 任务执行时动态添加新任务测试

模拟任务执行过程中新增任务,让队列的读写操作高度交织,放大并发冲突:

@Test
public void testDynamicTaskAddDuringExecution() throws InterruptedException {
    int poolSize = 20;
    int initialTaskCount = 100;
    int dynamicAddCountPerTask = 5;

    CountDownLatch allTasksDone = new CountDownLatch(initialTaskCount + initialTaskCount * dynamicAddCountPerTask);
    DepleatingFiFoThreadPool<Runnable> dtp = new DepleatingFiFoThreadPool<>(
            poolSize,
            Throwable::printStackTrace,
            "DynamicAddTest-",
            Runnable::run
    );

    // 添加初始任务,每个任务执行时动态添加新任务
    for (int i = 0; i < initialTaskCount; i++) {
        dtp.addAndStartThread(() -> {
            // 动态添加新任务
            for (int j = 0; j < dynamicAddCountPerTask; j++) {
                dtp.addAndStartThread(allTasksDone::countDown, "-dynamic-" + j);
            }
            allTasksDone.countDown();
        }, "-initial-" + i);
    }

    boolean completed = dtp.executeUntilDeplated(15000);
    assert completed : "任务执行超时";
    assert allTasksDone.getCount() == 0 : "仍有未完成的任务";
}

3. 超时场景下的并发测试

模拟超时触发时,任务仍在执行并修改队列的场景:

@Test
public void testTimeoutWithConcurrentQueueOperations() throws InterruptedException {
    int poolSize = 30;
    int taskCount = 500;

    CountDownLatch runningTasks = new CountDownLatch(taskCount);
    DepleatingFiFoThreadPool<Runnable> dtp = new DepleatingFiFoThreadPool<>(
            poolSize,
            Throwable::printStackTrace,
            "TimeoutTest-",
            Runnable::run
    );

    // 添加需要长时间运行的任务
    for (int i = 0; i < taskCount; i++) {
        dtp.addAndStartThread(() -> {
            runningTasks.countDown();
            try {
                // 任务持续运行,直到超时触发
                Thread.sleep(5000);
            } catch (InterruptedException ignored) {}
        }, "-longRunning-" + i);
    }

    // 等待所有任务开始运行
    runningTasks.await();

    // 设置短超时,触发超时逻辑
    boolean completed = dtp.executeUntilDeplated(1000);
    assert !completed : "未触发超时";

    // 继续添加任务,模拟超时后仍有队列操作
    for (int i = 0; i < 100; i++) {
        dtp.addAndStartThread(() -> {}, "-postTimeout-" + i);
    }

    // 等待所有任务完成
    Thread.sleep(6000);
}

代码修复建议

要彻底解决这个问题,需要确保对unstartedRunnables的所有操作都线程安全:

  1. 替换为线程安全队列:用java.util.concurrent.LinkedBlockingQueue替代ArrayList,它本身就是线程安全的,无需手动加锁。
  2. 同步所有队列操作:如果坚持使用ArrayList,则在addAndStartThread()的unstartedRunnables.add()处也加上synchronized(unstartedRunnables)锁。

示例修复(替换为LinkedBlockingQueue):

// 替换原ArrayList定义
private final BlockingQueue<Socket<A>> unstartedRunnables = new LinkedBlockingQueue<>();

// 修改addAndStartThread中的添加逻辑
if (mustCallLater) {
    unstartedRunnables.offer(new Socket<>(notRunningThread, threadPostfix));
}

// 修改finished中的取任务逻辑
Socket<A> e = unstartedRunnables.poll();
if (e != null) {
    start(e.r, e.threadPostfix);
}

这样可以完全避免并发修改导致的数组越界问题。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.17 06:47:02