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

Java多客户端TCP服务器:线程间非空值变为空的问题排查

嘿,这个并发问题太典型了!我之前做高并发TCP服务的时候也踩过类似的坑,咱们来拆解一下问题根源,再给你几个靠谱的解决方案。

首先,你虽然给messagesToSend加了volatile,但这解决不了根本问题——volatile只能保证变量引用的可见性,但管不了集合内部的操作安全。

为什么这么说?

  • volatile的作用是让其他线程能立刻看到这个变量的最新引用,但ArrayList的add操作是修改它内部的数组(不是改变messagesToSend的引用),所以volatile管不到内部元素的变化;
  • 更关键的是,ArrayList本身不是线程安全的,add、clear这些操作都不是原子的,并发调用时可能会导致内部数据结构混乱(比如扩容时的数组不一致、size计数错误),这就会出现你看到的“一个线程里非空,另一个线程里为空”的诡异情况。

接下来给你几个可行的解决思路,按推荐程度排序:

1. 用BlockingQueue实现生产者-消费者模型(最推荐)

Java的java.util.concurrent包提供的BlockingQueue天生就是为这种场景设计的,它自带线程安全的入队/出队操作,还支持阻塞等待(空队列时自动挂起线程,有消息再唤醒),完全不用自己处理锁和等待逻辑,代码简洁又健壮。

修改后的代码大概是这样:

public class MessageSender implements Runnable{
    private OutputStream output;
    private BlockingQueue<byte[]> messageQueue;

    public MessageSender(OutputStream _output) {
        output = _output;
        // 用LinkedBlockingQueue,容量可选,默认无界
        messageQueue = new LinkedBlockingQueue<>();
    }

    // 其他线程调用这个方法添加消息
    public void sendMessage(byte[] message) throws InterruptedException {
        // put()会阻塞直到队列有空间,也可以用offer()非阻塞返回是否成功
        messageQueue.put(message);
    }

    @Override
    public void run() {
        while (!Thread.currentThread().isInterrupted()) {
            try {
                // take()会阻塞直到队列有消息
                byte[] msg = messageQueue.take();
                output.write(msg);
                output.flush();
            } catch (IOException e) {
                // 处理客户端断开等IO异常,比如退出线程
                e.printStackTrace();
                break;
            } catch (InterruptedException e) {
                // 响应中断,退出循环
                Thread.currentThread().interrupt();
                break;
            }
        }
    }
}

2. 使用线程安全的集合类

如果不想改架构,也可以把ArrayList换成线程安全的集合:

  • CopyOnWriteArrayList:适合读多写少的场景,每次写操作都会复制底层数组,保证读操作永远不会抛出ConcurrentModificationException,但写操作性能稍差;
  • Collections.synchronizedList(new ArrayList<>()):给ArrayList的所有方法套上同步锁,保证线程安全,但并发高的时候锁竞争会影响性能。

比如用CopyOnWriteArrayList的示例:

public class MessageSender implements Runnable{
    private OutputStream output;
    private CopyOnWriteArrayList<byte[]> messagesToSend;

    public MessageSender(OutputStream _output) {
        output = _output;
        messagesToSend = new CopyOnWriteArrayList<>();
    }

    public void addMessage(byte[] message) {
        messagesToSend.add(message);
    }

    @Override
    public void run() {
        while (!Thread.currentThread().isInterrupted()) {
            // 遍历是线程安全的,不会被并发修改打断
            for (byte[] msg : new ArrayList<>(messagesToSend)) {
                try {
                    output.write(msg);
                    output.flush();
                    messagesToSend.remove(msg);
                } catch (IOException e) {
                    e.printStackTrace();
                    break;
                }
            }
            // 加个短睡眠避免空轮询浪费CPU
            try {
                Thread.sleep(10);
            } catch (InterruptedException e) {
                Thread.currentThread().interrupt();
                break;
            }
        }
    }
}

3. 手动加锁保证原子性和可见性

如果坚持用ArrayList,那必须给所有访问messagesToSend的操作加锁,比如用synchronized块或者ReentrantLock,确保同一时间只有一个线程能修改或读取集合:

public class MessageSender implements Runnable{
    private OutputStream output;
    private ArrayList<byte[]> messagesToSend;

    public MessageSender(OutputStream _output) {
        output = _output;
        messagesToSend = new ArrayList<>();
    }

    public synchronized void addMessage(byte[] message) {
        messagesToSend.add(message);
        // 唤醒等待的发送线程,避免空轮询
        notify();
    }

    @Override
    public void run() {
        while (!Thread.currentThread().isInterrupted()) {
            ArrayList<byte[]> messages = null;
            synchronized (this) {
                // 等待直到有消息
                while (messagesToSend.isEmpty()) {
                    try {
                        wait();
                    } catch (InterruptedException e) {
                        Thread.currentThread().interrupt();
                        break;
                    }
                }
                // 复制一份消息,避免持有锁时发送数据阻塞其他线程
                messages = new ArrayList<>(messagesToSend);
                messagesToSend.clear();
            }
            // 发送消息,此时不需要持有锁
            for (byte[] msg : messages) {
                try {
                    output.write(msg);
                    output.flush();
                } catch (IOException e) {
                    e.printStackTrace();
                    break;
                }
            }
        }
    }
}

最后再提醒一句:并发编程里,不要依赖volatile来保证集合的线程安全,它只对单个变量的引用有效,集合的线程安全必须靠同步机制或者线程安全的集合类来实现。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.26 10:24:56