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

基于Semaphore的生产者消费者模型出现竞态问题排查

生产者消费者模型中Semaphore实现的竞态问题分析

问题背景

我正在理解Semaphore在生产者消费者场景中的工作机制,实现了一个带有push、pop接口的有界队列,使用Semaphore管理多生产者与多消费者,但仍存在竞态条件,运行时抛出java.util.NoSuchElementException,相关代码及示例报错输出如下:

队列实现代码

class BoundedQueue {
    private Queue<Integer> q;
    int maxSize;
    Semaphore produce;
    Semaphore consume;
    
    BoundedQueue(int size) throws InterruptedException {
        maxSize = size;
        q = new LinkedList<>();
        produce = new Semaphore(maxSize);
        consume = new Semaphore(-1);
    }

    void push(int val) {
        try {
            produce.acquire();
            q.offer(val);
            System.out.println(Thread.currentThread().getName() + " produced " + val);
            if (q.size() > maxSize) {
                System.out.println("Max size breached");
            }
            consume.release();
        }
        catch (InterruptedException e) {
            System.out.println("Interrupted exception " + e);
            produce.release();
        }
    }

    int pop() {
        int val = -1;
        try {
            consume.acquire();
            if (q.isEmpty()) {
                System.out.println(Thread.currentThread().getName() + " is trying to consume empty queue");
                return val;
            }
            val = q.remove();
            System.out.println(Thread.currentThread().getName() + " consumed " + val);
            if (q.size() > maxSize) {
                System.out.println("Max size breached");
            }
            produce.release();
        }
        catch (InterruptedException e) {
            System.out.println("Interrupted exception " + e);
            consume.release();
        }
        return val;
    }
}

多生产者消费者测试代码

public class ProducerConsumerSemaphores extends Thread {

    private BoundedQueue que;
    private String workerType;
    
    ProducerConsumerSemaphores(BoundedQueue que, String workerType) {
        this.que = que;
        this.workerType = workerType;
    }

    public void run() {
        System.out.println("Starting worker " + workerType + " " + Thread.currentThread().getName());
        if (workerType.equals("producer")) {
            Random rand = new Random();
            int val = rand.nextInt(1, 100);
            for (int i=0; i<10; ++i) {
                que.push(val+i);
                try {
                    Thread.sleep(10);
                }
                catch (Exception e) {
                    System.out.println("Caught " + e);
                }
            }
        }
        else {
            for (int i=0; i<10; ++i) {
                que.pop();
                try {
                    Thread.sleep(100);
                }
                catch (Exception e) {
                    System.out.println("Caught " + e);
                }
            }
        }
        System.out.println("Ending woker " + workerType + " " + Thread.currentThread().getName());
    }

    public static void main(String args[]) {
        System.out.println();
        try {
            BoundedQueue que = new BoundedQueue(5);
            ProducerConsumerSemaphores p1 = new ProducerConsumerSemaphores(que, "producer");
            p1.setName("p1");
            ProducerConsumerSemaphores p2 = new ProducerConsumerSemaphores(que, "producer");
            p2.setName("p2");
            ProducerConsumerSemaphores p3 = new ProducerConsumerSemaphores(que, "producer");
            p3.setName("p3");
            p1.start();
            p2.start();
            ProducerConsumerSemaphores c1 = new ProducerConsumerSemaphores(que, "consumer");
            c1.setName("c1");
            c1.start();
            ProducerConsumerSemaphores c2 = new ProducerConsumerSemaphores(que, "consumer");
            c2.setName("c2");
            c2.start();
            p3.start();
            ProducerConsumerSemaphores c3 = new ProducerConsumerSemaphores(que, "consumer");
            c3.setName("c3");
            c3.start();

        }
        catch (Exception e) {
            System.out.println("Exception occurred");
        }
    }
}

报错输出

Starting worker producer p2
Starting worker producer p1
Starting worker consumer c1
Starting worker consumer c3
Starting worker producer p3
Starting worker consumer c2
p2 produced 87
p1 produced 61
p3 produced 90
c1 consumed 87
Exception in thread "c2" java.util.NoSuchElementException
        at java.base/java.util.LinkedList.removeFirst(LinkedList.java:274)
        at java.base/java.util.LinkedList.remove(LinkedList.java:689)     
        at BoundedQueue.pop(ProducerConsumerSemaphores.java:42)
        at ProducerConsumerSemaphores.run(ProducerConsumerSemaphores.java:85)
p2 produced 88
c3 consumed 88
p2 produced 89

问题原因分析

1. Semaphore初始化错误

consume信号量初始值设为-1是完全错误的——Semaphore的许可数不能为负数,这会直接导致后续消费逻辑的状态混乱。正确的初始值应该是0,因为队列初始为空,没有可消费的元素。

2. 队列操作未保证原子性

LinkedList本身不是线程安全容器,且pop()方法中if (q.isEmpty())和val = q.remove()之间存在竞态窗口:

  • 假设队列只剩1个元素,两个消费者线程同时通过consume.acquire()获取许可
  • 线程A先判断队列不为空,还未执行删除操作
  • 线程B紧接着也判断队列不为空
  • 线程A执行删除取走最后一个元素,队列变空
  • 线程B再执行删除时,队列已空,直接抛出NoSuchElementException

3. Semaphore与队列状态不同步

Semaphore的许可数和实际队列元素数没有绑定同步:多个生产者同时执行q.offer(val)时,可能导致队列元素超过maxSize;消费者许可数可能大于实际队列元素数,引发空队列消费的情况。

修复方案

核心修改点

  • 修正consume信号量初始值为0
  • 对队列的所有操作(添加、删除、判断空)加锁,保证原子性
  • 确保Semaphore操作和队列状态的一致性

修复后的BoundedQueue代码

class BoundedQueue {
    private Queue<Integer> q;
    int maxSize;
    Semaphore produce;
    Semaphore consume;
    private final Object lock = new Object(); // 队列操作同步锁
    
    BoundedQueue(int size) {
        maxSize = size;
        q = new LinkedList<>();
        produce = new Semaphore(maxSize);
        consume = new Semaphore(0); // 初始无消费许可
    }

    void push(int val) {
        try {
            produce.acquire();
            synchronized (lock) {
                q.offer(val);
                System.out.println(Thread.currentThread().getName() + " produced " + val);
                if (q.size() > maxSize) {
                    System.out.println("Max size breached");
                }
            }
            consume.release();
        } catch (InterruptedException e) {
            System.out.println("Interrupted exception " + e);
            produce.release();
        }
    }

    int pop() {
        int val = -1;
        try {
            consume.acquire();
            synchronized (lock) {
                if (q.isEmpty()) {
                    System.out.println(Thread.currentThread().getName() + " is trying to consume empty queue");
                    return val;
                }
                val = q.remove();
                System.out.println(Thread.currentThread().getName() + " consumed " + val);
                if (q.size() > maxSize) {
                    System.out.println("Max size breached");
                }
            }
            produce.release();
        } catch (InterruptedException e) {
            System.out.println("Interrupted exception " + e);
            consume.release();
        }
        return val;
    }
}

说明

  • 新增lock对象,用synchronized块包裹所有队列操作,确保同一时间只有一个线程能修改或访问队列状态
  • consume信号量初始化为0,匹配队列初始为空的状态
  • Semaphore的acquire()在锁外执行,避免持有锁等待许可导致性能损耗;release()在锁外执行,队列操作完成后再更新许可数

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.20 16:59:54