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

Java有界缓冲区生产者-消费者模型死锁原因排查求助

问题分析与解决:生产者-消费者模型死锁排查

问题描述

我实现了一个典型的有界缓冲区生产者-消费者模型,本地运行正常,但在线提交时提示未成功运行5次,疑似死锁。我认为wait()和notify()的使用正确,不清楚还有哪些原因会导致死锁,怀疑问题出在Buffer.java中,附上所有相关代码。

相关代码

Buffer.java

package HW4;

public class Buffer {
    
    private double[] array;
    private int empty;
    private int full;
    private int capacity;

    //Initialize values
    public Buffer(int capacity) {
        this.capacity = capacity;
        this.array = new double[capacity];
        this.empty = capacity;
        this.full = 0;
    }

    //Produce method
    //Check if full; add value to array; notify()
    synchronized public void produce(double value) throws InterruptedException {
        while (full == capacity) {
            wait();
        }

        array[full] = value;
        full++;
        empty--;
        notify();
    }

    //Consume method
    //Check if empty; remove value from array; notify()
    synchronized public double consume() throws InterruptedException {
        while (empty == capacity) {
            wait();
        }

        double retValue = array[capacity - empty - 1];
        array[capacity - empty - 1] = 0d;
        full--;
        empty++;
        notify();
        
        return retValue;
    }
}

Main.java

package HW4;

public class Main {
    
    public static final int iterations = 1000000;
    public static void main(String[] args) {

        //Create a buffer item along with a producer and consumer instantiation
        Buffer buffer = new Buffer(1000);
        Producer producer = new Producer(buffer, iterations);
        Consumer consumer = new Consumer(buffer, iterations);

        //Create a thread for both the producer and consumer; 1 of each
        Thread prodThread = new Thread(producer);
        Thread consThread = new Thread(consumer);

        //Start each thread
        prodThread.start();
        consThread.start();

        //Attempt to join the threads after they are done running.
        try {
            prodThread.join();
            consThread.join();
        } catch (InterruptedException e) {
            System.out.println("ERROR IN MAIN.\n");
        }
        System.out.println("Exiting!\n");
    }
}

Producer.java

package HW4;

import java.util.Random;

public class Producer implements Runnable {
    
    private Buffer buffer;
    private Random random;
    private int iterations;

    //Import buffer item and create Random item for generating doubles
    public Producer(Buffer buffer, int iterations) {
        this.buffer = buffer;
        this.random = new Random();
        this.iterations = iterations;
    }
    
    //Loop to create 1,000,000 items and add them to the buffer
    //Maintain total value for outputting
    @Override
    public void run() {
        double total = 0d;
        for (int i = 0; i < iterations; i++) {
            double bufferElement = random.nextDouble() * 100.0;
            
            try {
                buffer.produce(bufferElement);
                total += bufferElement;
            } catch (InterruptedException e) {
                System.out.println("ERROR in PRODUCER.\n");
            }

            if (i % 100000 == 0 && i != 0) {
                System.out.println("Producer: Generated " + i/1000 + ",000 items, Cumulative value of generated items=" + formatDouble(total));
            }
        }
        System.out.println("Producer: Finished generating 1,000,000 items");
    }

    private String formatDouble(double value) {
        return String.format("%.3f", value);
    }
}

Consumer.java

package HW4;

public class Consumer implements Runnable {

    private Buffer buffer;
    private int iterations;

    //Import buffer item
    public Consumer(Buffer buffer, int iterations) {
        this.buffer = buffer;
        this.iterations = iterations;
    }

    //Loop to remove 1,000,000 items from the buffer
    //Maintain total value for outputting
    @Override
    public void run() {
        double total = 0d;
        for (int i = 0; i < iterations; i++) {
            try {
                double bufferElement = buffer.consume();
                total += bufferElement;
            } catch (InterruptedException e) {
                System.out.println("ERROR in CONSUMER.\n");
            }

            if (i % 100000 == 0 && i != 0) {
                System.out.println("Consumer: Consumed " + i/1000 + ",000 items, Cumulative value of consumed items=" + formatDouble(total));
            }
        }
        System.out.println("Consumer: Finished consuming 1,000,000 items");
    }

    private String formatDouble(double value) {
        return String.format("%.3f", value);
    }
}

问题排查与解决

1. 线程中断处理不当导致永久阻塞

这是导致线上“死锁”的核心原因:

  • 当生产者/消费者线程在wait()时被中断,会抛出InterruptedException,当前代码仅打印错误就继续循环,导致迭代计数i仍会递增,最终生产/消费的元素数量少于预期的1000000个。
  • 例如生产者少生产1个元素,消费者会一直卡在consume()的wait()中等待不存在的元素,表现为线程永久阻塞,被判定为死锁。

修复方案:
修改生产者和消费者的循环逻辑,确保只有成功完成生产/消费后才递增计数,同时捕获中断时恢复线程中断状态并退出循环:

修改后的Producer.run()

@Override
public void run() {
    double total = 0d;
    int i = 0;
    while (i < iterations) {
        double bufferElement = random.nextDouble() * 100.0;
        try {
            buffer.produce(bufferElement);
            total += bufferElement;
            i++;
        } catch (InterruptedException e) {
            System.out.println("ERROR in PRODUCER.\n");
            Thread.currentThread().interrupt(); // 恢复中断状态
            break;
        }

        if (i % 100000 == 0 && i != 0) {
            System.out.println("Producer: Generated " + i/1000 + ",000 items, Cumulative value of generated items=" + formatDouble(total));
        }
    }
    System.out.println("Producer: Finished generating " + i + " items");
}

修改后的Consumer.run()

@Override
public void run() {
    double total = 0d;
    int i = 0;
    while (i < iterations) {
        try {
            double bufferElement = buffer.consume();
            total += bufferElement;
            i++;
        } catch (InterruptedException e) {
            System.out.println("ERROR in CONSUMER.\n");
            Thread.currentThread().interrupt(); // 恢复中断状态
            break;
        }

        if (i % 100000 == 0 && i != 0) {
            System.out.println("Consumer: Consumed " + i/1000 + ",000 items, Cumulative value of consumed items=" + formatDouble(total));
        }
    }
    System.out.println("Consumer: Finished consuming " + i + " items");
}

2. Buffer实现逻辑错误(非死锁原因,但不符合需求)

原Buffer的消费逻辑是后入先出(LIFO),而非生产者-消费者模型通常要求的先入先出(FIFO),虽然不会导致死锁,但会打乱数据顺序。建议改用环形队列实现,逻辑更清晰且不易出错:

修复后的Buffer.java

package HW4;

public class Buffer {
    
    private double[] array;
    private int head; // 消费指针,指向当前可消费的元素位置
    private int tail; // 生产指针,指向当前可生产的元素位置
    private int count; // 当前缓冲区元素数量
    private int capacity;

    public Buffer(int capacity) {
        this.capacity = capacity;
        this.array = new double[capacity];
        this.head = 0;
        this.tail = 0;
        this.count = 0;
    }

    synchronized public void produce(double value) throws InterruptedException {
        while (count == capacity) {
            wait();
        }

        array[tail] = value;
        tail = (tail + 1) % capacity; // 环形指针移动
        count++;
        notify();
    }

    synchronized public double consume() throws InterruptedException {
        while (count == 0) {
            wait();
        }

        double retValue = array[head];
        head = (head + 1) % capacity; // 环形指针移动
        count--;
        notify();
        
        return retValue;
    }
}

总结

线上的死锁假象本质是线程中断处理不当导致的生产/消费数量不匹配,修复中断处理逻辑即可解决;同时优化Buffer实现为FIFO环形队列,符合生产者-消费者模型的标准需求。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 07:37:01