基于wait/notify的批量短信发送线程同步问题求助
一、wait/notify/notifyAll 工作原理
- wait():当前线程会立刻释放持有的对象锁,进入该对象的等待队列,直到被
notify()或notifyAll()唤醒。唤醒后不能直接执行,必须重新竞争对象锁,抢到锁后才会继续执行。注意:必须在synchronized代码块/方法内调用,否则会抛出IllegalMonitorStateException。 - notify():唤醒对象等待队列中的任意一个线程,被唤醒的线程不会立刻执行,要等当前线程释放锁后,和其他线程竞争锁资源。同样必须在
synchronized块内调用。 - notifyAll():唤醒对象等待队列中的所有线程,所有被唤醒的线程都会参与锁竞争,抢到锁的线程继续执行,未抢到的则继续等待锁。
二、问题排查
你遇到的ArrayIndexOutOfBoundsException本质是线程竞态条件导致的,具体原因如下:
- 循环判断与同步块的间隙:
Sender线程的while (!Client.messages.isEmpty())判断在synchronized块之外。当队列最后一条消息被某线程移除后,其他线程已经通过了while判断,进入同步块时队列已经为空,执行remove(0)自然报错。 - CopyOnWriteArrayList误用:虽然它是线程安全集合,但你用它作为同步锁对象,且结合外层非同步的while判断,依然会出现竞态问题。另外
CopyOnWriteArrayList.remove(0)效率极低,每次删除都会复制整个数组。 - SessionProducer逻辑冗余:它的while循环基于消息队列是否为空,但会话重建只需要在连接断开时触发,不需要常驻循环监听;当前逻辑会导致会话正常时它一直wait,浪费资源。
- notifyAll时机不合理:
Sender每次发送完消息都调用notifyAll(),会唤醒所有等待线程(包括SessionProducer),完全没必要,反而增加锁竞争开销。
三、优化解决方案
1. 核心优化思路
- 用阻塞队列替代手动同步的List,利用阻塞队列自带的线程安全阻塞特性,避免手动wait/notify的繁琐和错误。
- 采用集中式会话管理,由单独的管理器负责会话的有效性检查和重建,避免多个线程重复处理会话逻辑。
2. 具体代码实现
修正后的SMPPSession(增加模拟断开逻辑)
import java.util.Random; public class SMPPSession { private volatile boolean bind; private static final Random idGenerator = new Random(); public SMPPSession() { this.bind = false; } public int sendMessage(String msg) throws Exception { if (!bind) { throw new IllegalStateException("Session not bound"); } try { Thread.sleep(100L); System.out.println("Sending message: " + msg); // 模拟随机断开连接 if (idGenerator.nextInt(10) == 0) { bind = false; throw new IllegalStateException("Session disconnected unexpectedly"); } return Math.abs(idGenerator.nextInt()); } catch (InterruptedException e) { Thread.currentThread().interrupt(); throw new Exception("Send interrupted", e); } } public void reBind() throws InterruptedException { System.out.println("Rebinding session..."); Thread.sleep(1000L); this.bind = true; System.out.println("Session reconnected successfully"); } public boolean isBind() { return bind; } }
会话管理器SessionManager
public class SessionManager { private final SMPPSession session; private final Object lock = new Object(); public SessionManager() { this.session = new SMPPSession(); // 初始绑定会话 try { session.reBind(); } catch (InterruptedException e) { Thread.currentThread().interrupt(); throw new RuntimeException("Failed to initialize session", e); } } public SMPPSession getValidSession() throws InterruptedException { synchronized (lock) { while (!session.isBind()) { session.reBind(); } return session; } } public void invalidateSession() { synchronized (lock) { session.bind = false; } } }
发送线程Sender(基于阻塞队列)
import java.util.concurrent.BlockingQueue; public class Sender extends Thread { private final BlockingQueue<String> messageQueue; private final SessionManager sessionManager; public Sender(String name, BlockingQueue<String> messageQueue, SessionManager sessionManager) { super(name); this.messageQueue = messageQueue; this.sessionManager = sessionManager; } @Override public void run() { while (!Thread.currentThread().isInterrupted()) { try { String msg = messageQueue.take(); // 队列为空时自动阻塞 SMPPSession session = sessionManager.getValidSession(); int msgId = session.sendMessage(msg); System.out.println(getName() + " sent msg: " + msg + ", msgId: " + msgId); } catch (InterruptedException e) { Thread.currentThread().interrupt(); break; } catch (Exception e) { System.out.println(getName() + " failed to send msg: " + e.getMessage()); // 标记会话无效,触发重建 sessionManager.invalidateSession(); // 将失败消息放回队列(可选,根据业务调整重发策略) try { messageQueue.put(msg); } catch (InterruptedException ie) { Thread.currentThread().interrupt(); break; } } } } }
主类Client
import java.util.concurrent.ArrayBlockingQueue; import java.util.concurrent.BlockingQueue; public class Client { public static void main(String[] args) { BlockingQueue<String> messageQueue = new ArrayBlockingQueue<>(10); // 填充消息 messageQueue.add("msg1"); messageQueue.add("msg2"); messageQueue.add("msg3"); messageQueue.add("msg4"); messageQueue.add("msg5"); messageQueue.add("msg6"); SessionManager sessionManager = new SessionManager(); // 创建发送线程 Sender sender1 = new Sender("Sender1", messageQueue, sessionManager); Sender sender2 = new Sender("Sender2", messageQueue, sessionManager); Sender sender3 = new Sender("Sender3", messageQueue, sessionManager); Sender sender4 = new Sender("Sender4", messageQueue, sessionManager); sender1.start(); sender2.start(); sender3.start(); sender4.start(); // 等待所有消息发送完成后中断线程 try { Thread.sleep(8000); sender1.interrupt(); sender2.interrupt(); sender3.interrupt(); sender4.interrupt(); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } } }
3. 优化点说明
- BlockingQueue替代手动同步:
take()方法会在队列为空时自动阻塞,彻底避免了手动wait/notify的竞态问题。 - 集中式会话管理:
SessionManager统一处理会话有效性检查和重建,避免多个线程重复执行重建逻辑。 - volatile修饰bind状态:确保bind变量的修改对所有线程可见,避免读取过期状态。
- 异常处理与消息重发:发送失败时将消息放回队列,保证消息不丢失(可根据业务调整重发次数或策略)。
- 优雅线程终止:用
interrupt()终止线程,避免线程无限等待。
内容的提问来源于stack exchange,提问作者Wrapper
相关产品推荐
相关产品推荐

