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
相关产品推荐
相关产品推荐

