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

基于wait/notify的批量短信发送线程同步问题求助

一、wait/notify/notifyAll 工作原理
  • wait():当前线程会立刻释放持有的对象锁,进入该对象的等待队列,直到被notify()或notifyAll()唤醒。唤醒后不能直接执行,必须重新竞争对象锁,抢到锁后才会继续执行。注意:必须在synchronized代码块/方法内调用,否则会抛出IllegalMonitorStateException。
  • notify():唤醒对象等待队列中的任意一个线程,被唤醒的线程不会立刻执行,要等当前线程释放锁后,和其他线程竞争锁资源。同样必须在synchronized块内调用。
  • notifyAll():唤醒对象等待队列中的所有线程,所有被唤醒的线程都会参与锁竞争,抢到锁的线程继续执行,未抢到的则继续等待锁。
二、问题排查

你遇到的ArrayIndexOutOfBoundsException本质是线程竞态条件导致的,具体原因如下:

  1. 循环判断与同步块的间隙:Sender线程的while (!Client.messages.isEmpty())判断在synchronized块之外。当队列最后一条消息被某线程移除后,其他线程已经通过了while判断,进入同步块时队列已经为空,执行remove(0)自然报错。
  2. CopyOnWriteArrayList误用:虽然它是线程安全集合,但你用它作为同步锁对象,且结合外层非同步的while判断,依然会出现竞态问题。另外CopyOnWriteArrayList.remove(0)效率极低,每次删除都会复制整个数组。
  3. SessionProducer逻辑冗余:它的while循环基于消息队列是否为空,但会话重建只需要在连接断开时触发,不需要常驻循环监听;当前逻辑会导致会话正常时它一直wait,浪费资源。
  4. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.10 12:55:25