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

基于锁机制实现双线程同步:死锁问题排查与解决

问题分析与解决方案

死锁根源

你的代码出现死锁的核心原因有两个:

  1. handleEventData方法的synchronized锁与ReadWriteLock形成嵌套阻塞:串口事件线程进入synchronized的handleEventData后,提交任务到单线程执行器并调用future.get()阻塞,执行器线程尝试获取graphReadLock;若此时另一个线程调用closeGraphInputStream尝试获取graphWriteLock,读写锁的排他性会导致执行器线程等待写锁,串口事件线程等待执行器任务完成,最终形成死锁。
  2. graphInputStreamOpen的判断与锁获取非原子:先判断状态再拿锁,会出现判断为true后流被关闭的竞态条件,导致写流时抛出异常。

修复方案

1. 移除不必要的synchronized锁

由于使用newSingleThreadExecutor(),所有提交的任务会串行执行,无需额外的synchronized同步,移除该修饰符可减少锁竞争。

2. 原子化状态判断与锁获取

将graphInputStreamOpen的判断移到读锁内部,确保判断和写操作的原子性,彻底消除竞态条件。

3. 优雅处理流关闭异常

在写PipedOutputStream时捕获IOException,处理流已关闭的场景。

4. 完善资源清理逻辑

在close方法中主动关闭graph输入流和执行器,避免资源泄漏。

修改后的代码

public class AppSerialPort implements SerialPortEventListener, Closeable {
    private final SerialPort serialPort;
    private final ExecutorService executor = Executors.newSingleThreadExecutor();
    private final ReadWriteLock graphLock = new ReentrantReadWriteLock();
    private final Lock graphWriteLock = graphLock.writeLock();
    private final Lock graphReadLock = graphLock.readLock();

    private PipedOutputStream graphPacketOutputStream;
    private PipedInputStream graphPacketInputStream;
    private volatile boolean graphInputStreamOpen;
    private volatile boolean connectionClosed;

    public AppSerialPort(String portIdentifier) throws NoSuchPortException, PortInUseException, TooManyListenersException, UnsupportedCommOperationException, IOException {
        CommPortIdentifier portId = CommPortIdentifier.getPortIdentifier(portIdentifier);
        this.serialPort = (SerialPort) portId.open("MyApp", 0);
        this.initListener();
    }

    public void serialEvent(SerialPortEvent event) {
        try {
            if (event.getEventType() == SerialPortEvent.DATA_AVAILABLE && !connectionClosed) {
                int available = serialPort.getInputStream().available();
                if (available > 0) {
                    byte[] buffer = new byte[available];
                    serialPort.getInputStream().read(buffer);
                    handleEventData(buffer);
                }
            }
        } catch (IOException ex) {
            // 建议添加日志记录,而非直接忽略
        }
    }

    private void handleEventData(final byte[] buffer) {
        Callable<Void> eventHandlerTask = () -> {
            graphReadLock.lock();
            try {
                if (graphInputStreamOpen) {
                    graphPacketOutputStream.write(buffer);
                }
            } catch (IOException ex) {
                // 流已关闭,无需额外处理,或记录日志
            } finally {
                graphReadLock.unlock();
            }
            return null;
        };

        Future<Void> future = executor.submit(eventHandlerTask);
        try {
            future.get();
        } catch (InterruptedException | ExecutionException ex) {
            serialPort.removeEventListener();
            executor.shutdownNow();
        }
    }

    private void initListener() throws TooManyListenersException {
        serialPort.addEventListener(this);
        serialPort.notifyOnDataAvailable(true);
    }

    public InputStream openGraphInputStream() throws IOException {
        graphWriteLock.lock();
        try {
            graphPacketOutputStream = new PipedOutputStream();
            graphPacketInputStream = new PipedInputStream(15000);
            graphPacketInputStream.connect(graphPacketOutputStream);
            graphInputStreamOpen = true;
            return graphPacketInputStream;
        } finally {
            graphWriteLock.unlock();
        }
    }

    public void closeGraphInputStream() throws IOException {
        graphWriteLock.lock();
        try {
            if (graphInputStreamOpen) {
                graphInputStreamOpen = false;
                graphPacketInputStream.close();
                graphPacketOutputStream.close();
            }
        } finally {
            graphWriteLock.unlock();
        }
    }

    @Override
    public void close() throws IOException {
        connectionClosed = true;
        serialPort.removeEventListener();
        serialPort.close();
        closeGraphInputStream();
        executor.shutdown();
        try {
            if (!executor.awaitTermination(1, TimeUnit.SECONDS)) {
                executor.shutdownNow();
            }
        } catch (InterruptedException ex) {
            executor.shutdownNow();
        }
    }
}

关键修改说明

  • 将graphInputStreamOpen和connectionClosed标记为volatile,确保多线程下的状态可见性。
  • openGraphInputStream加入写锁保护,避免多线程调用时的资源冲突。
  • close方法中主动关闭graph输入流并优雅关闭执行器,确保所有资源被清理。
  • 将状态判断移到锁内部,彻底消除竞态条件,同时避免流关闭时的写操作异常。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.21 18:37:35