当前位置:

首页 > 编程开发 > Java并发消息系统会话与wait/notify详解

Java并发消息系统会话与wait/notify详解

本文将探讨如何利用Java的wait/notify机制在多线程环境中实现短信批量发送与会话重连。我们将分析常见的同步问题,特别是因不当的isEmpty()检查和共享资源访问导致的ArrayIndexOutOfBoundsException,并提供正确同步共享资源和管理线程状态的策略,以构建健壮的并发操作。

Java并发消息发送系统中的会话管理与wait/notify机制深度解析

Java并发消息发送系统中的会话管理与`wait`/`notify`机制深度解析。本文将探讨如何利用Java的`wait`/`notify`机制在多线程环境中实现短信批量发送与会话重连。我们将分析常见的同步问题,特别是因不当的`isEmpty()`检查和共享资源访问导致的`ArrayIndexOutOfBoundsException`,并提供正确同步共享资源和管理线程状态的策略,以构建健壮的并发操作。

1. 引言:并发消息发送与会话管理挑战

在企业级应用中,批量发送短信(或其他消息)是一个常见需求。为了提高吞吐量,通常会采用多线程并发发送的策略。然而,消息发送依赖于与外部服务器(如SMSC)建立的会话(SMPPSession)。这种会话可能因网络波动、服务器重启等原因中断,此时需要一个机制来重新建立会话,并在会话重连期间暂停所有发送操作,待会话恢复后再继续。

本教程将深入探讨如何使用Java的Object.wait()和Object.notifyAll()机制来协调多个消息发送线程和一个会话管理线程,以实现上述功能。我们将分析在并发场景下可能遇到的同步问题,并提供一套健壮的解决方案。

2. wait()与notify()机制详解

wait()和notify()(或notifyAll())是Java中用于线程间协作的基础机制,它们允许线程在特定条件下暂停执行并等待,直到另一个线程通知它条件满足。

  • wait(): 当一个线程调用wait()方法时,它会释放当前持有的对象锁,并进入等待状态,直到被notify()或notifyAll()唤醒,或者被中断。
  • notify(): 唤醒在该对象上等待的一个任意线程。
  • notifyAll(): 唤醒在该对象上等待的所有线程。

关键点:

  1. 必须在synchronized块内调用:wait()、notify()和notifyAll()方法必须在持有对象监视器(即synchronized块所锁定的对象)的情况下调用,否则会抛出IllegalMonitorStateException。
  2. 操作同一个监视器对象:所有等待和通知操作都必须针对同一个对象进行,这个对象充当了线程间通信的“信号量”。
  3. 虚假唤醒与条件检查:wait()方法可能会在没有收到通知的情况下被唤醒(虚假唤醒)。因此,wait()调用通常应该放在一个while循环中,不断检查等待的条件是否真正满足。

3. 原始代码中的同步问题分析

原始代码尝试使用Client.messages列表作为监视器对象来协调线程。然而,其中存在几个关键的同步问题,导致了ArrayIndexOutOfBoundsException和线程不同步。

3.1 while (!Client.messages.isEmpty())的竞态条件

在Sender和SessionProducer线程的run()方法中,外部的while (!Client.messages.isEmpty())循环条件是在synchronized (Client.messages)块外部检查的。

// Sender线程示例
while (!Client.messages.isEmpty()){ // 问题:在同步块外检查
    synchronized (Client.messages){
        // ...
    }
}

问题分析: 假设Client.messages中只剩一条消息。多个Sender线程可能同时执行到while (!Client.messages.isEmpty()),它们都发现列表不为空,然后都尝试进入synchronized (Client.messages)块。当第一个线程进入同步块并成功移除消息后,列表变为空。此时,后续进入同步块的线程在执行Client.messages.remove(0)时,就会因为列表已空而抛出ArrayIndexOutOfBoundsException。

3.2 remove(0)的并发访问问题

即使isEmpty()检查放在同步块内,remove(0)操作也需要谨慎。如果多个线程同时尝试移除,且消息数量不足,仍然可能导致问题。在CopyOnWriteArrayList中,remove(0)本身是线程安全的,但它并不能阻止在列表为空时尝试移除。

3.3 不当的通知机制

Sender线程在成功发送消息后调用了Client.messages.notifyAll()。然而,此时SessionProducer线程可能正在等待会话断开,或者其他Sender线程可能正在等待消息或会话。这种通知通常是不必要的,并且可能导致不必要的线程唤醒,甚至掩盖真正的等待条件。

3.4 wait()的条件检查缺失

原始代码中的wait()没有放在while循环中检查条件。

// 原始Sender线程的等待逻辑
} else {
    try {
        Client.messages.wait(); // 问题:没有在while循环中检查条件
    } catch (InterruptedException e) {
        throw new RuntimeException(e);
    }
}

这可能导致线程在被唤醒后,其等待的条件(例如smppSession.isBind()为true或Client.messages不为空)实际上并未满足,从而导致逻辑错误或再次进入等待状态。

4. 改进方案与最佳实践

为了解决上述问题,我们需要对代码进行重构,遵循以下核心原则:

  1. 统一监视器对象:选择一个能代表共享状态的唯一对象作为所有wait()和notifyAll()操作的监视器。在本例中,SMPPSession对象本身是一个很好的选择,因为它代表了会话的绑定状态。
  2. 所有共享资源访问都需同步:包括isEmpty()、remove()等操作,都必须在持有监视器锁的同步块内执行。
  3. wait()必须在while循环中检查条件:防止虚假唤醒和条件不满足时继续执行。
  4. 明确通知时机:只有当某个线程改变了其他线程正在等待的条件时,才调用notifyAll()。

4.1 改进SMPPSession类

SMPPSession作为共享资源,其bind状态是所有线程关注的焦点。我们可以将它作为监视器对象。

public class SMPPSession {
    private boolean bind = false; // 初始状态为未绑定
    private static final Random idGenerator = new Random();

    public synchronized int sendMessage(String msg) { // 保持sendMessage同步
        try {
            Thread.sleep(100L); // 模拟发送延迟
            System.out.println("Sending message: " + msg);
            return Math.abs(idGenerator.nextInt());
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
            System.err.println("Message sending interrupted: " + e.getMessage());
        }
        return -1;
    }

    public synchronized void reBind() { // reBind方法也同步
        try {
            System.out.println("Rebinding...");
            Thread.sleep(2000L); // 模拟重连延迟
            this.bind = true;
            System.out.println("Session established!");
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
            System.err.println("Rebinding interrupted: " + e.getMessage());
        }
    }

    public synchronized boolean isBind() { // isBind方法也同步
        return this.bind;
    }

    public synchronized void setBind(boolean bind) { // 允许外部设置绑定状态
        this.bind = bind;
    }
}

4.2 改进Sender线程

Sender线程需要等待两个条件:SMPPSession已绑定,且消息队列不为空。

import java.util.concurrent.CopyOnWriteArrayList;

public class Sender extends Thread {
    private SMPPSession smppSession;
    private CopyOnWriteArrayList messages; // 引用共享消息列表
    private volatile boolean running = true; // 控制线程生命周期

    public Sender(String name, SMPPSession smppSession, CopyOnWriteArrayList messages) {
        this.setName(name);
        this.smppSession = smppSession;
        this.messages = messages;
    }

    public void terminate() {
        this.running = false;
        // 确保线程不会无限等待,如果正在wait(),需要被中断或notify
        synchronized (smppSession) {
            smppSession.notifyAll();
        }
    }

    @Override
    public void run() {
        while (running) {
            synchronized (smppSession) { // 使用smppSession作为监视器
                // 等待条件:会话未绑定 或 消息队列为空
                while (!smppSession.isBind() || messages.isEmpty()) {
                    // 如果消息已全部发送且会话已绑定,则此Sender可以退出
                    if (messages.isEmpty() && smppSession.isBind()) {
                        System.out.println(getName() + ":所有消息已发送完毕,线程退出。");
                        running = false; // 标记为停止
                        smppSession.notifyAll(); // 通知其他可能等待的线程
                        break; // 跳出内部while循环
                    }
                    try {
                        System.out.println(getName() + ":等待中... 会话绑定状态: " + smppSession.isBind() + ", 消息队列是否为空: " + messages.isEmpty());
                        smppSession.wait(); // 等待在smppSession对象上
                    } catch (InterruptedException e) {
                        System.out.println(getName() + ":被中断,线程退出。");
                        Thread.currentThread().interrupt();
                        running = false; // 标记为停止
                        break; // 跳出内部while循环
                    }
                }
                if (!running) { // 如果在等待过程中被标记为停止,则退出外部while循环
                    break;
                }

                // 条件满足:smppSession已绑定且messages不为空
                final String msg = messages.remove(0); // 安全移除消息
                final int msgId = smppSession.sendMessage(msg);
                System.out.println(Thread.currentThread().getName() + " 发送消息并收到ID: " + msgId + "。剩余消息数:" + messages.size());
                // 发送消息后,如果消息队列变空,可能需要通知其他Sender线程退出
                // 或者如果Producer在等待消息队列状态,则需要通知
                // 这里暂时不需要notifyAll,因为发送消息通常不改变Producer的等待条件
                // 但如果messages.isEmpty()是Producer的等待条件之一,则需要
            }
            // 考虑在发送消息后短暂休眠,避免发送过快
            try {
                Thread.sleep(50);
            } catch (InterruptedException e) {
                Thread.currentThread().interrupt();
                running = false;
            }
        }
    }
}

4.3 改进SessionProducer线程

SessionProducer线程负责在会话未绑定时进行重连,并在重连成功后通知所有Sender线程。

import java.util.concurrent.CopyOnWriteArrayList;

public class SessionProducer extends Thread {
    private SMPPSession smppSession;
    private CopyOnWriteArrayList messages; // 引用共享消息列表,用于判断是否还有消息需要发送
    private volatile boolean running = true;

    public SessionProducer(String name, SMPPSession smppSession, CopyOnWriteArrayList messages) {
        this.setName(name);
        this.smppSession = smppSession;
        this.messages = messages;
    }

    public void terminate() {
        this.running = false;
        synchronized (smppSession) {
            smppSession.notifyAll();
        }
    }

    @Override
    public void run() {
        while (running) {
            synchronized (smppSession) { // 使用smppSession作为监视器
                // 如果会话已绑定,且所有消息都已发送完毕,则Producer可以退出
                if (smppSession.isBind() && messages.isEmpty()) {
                    System.out.println(getName() + ":所有消息已发送完毕,会话已绑定,线程退出。");
                    running = false;
                    smppSession.notifyAll(); // 通知所有线程可以退出
                    break;
                }

                // 等待条件:会话已绑定 或 消息队列为空(如果 Producer 也需要关注消息队列状态)
                // 这里主要关注会话绑定状态
                while (smppSession.isBind() && !messages.isEmpty()) { // 如果会话已绑定且还有消息要发,Producer等待
                    try {
                        System.out.println(getName() + ":等待中... 会话已绑定,等待会话断开或所有消息发送完毕。");
                        smppSession.wait(); // 等待在smppSession对象上
                    } catch (InterruptedException e) {
                        System.out.println(getName() + ":被中断,线程退出。");
                        Thread.currentThread().interrupt();
                        running = false;
                        break;
                    }
                }
                if (!running) {
                    break;
                }

                // 此时,会话可能未绑定,或者消息队列为空(如果上面条件包含)
                if (!smppSession.isBind()) { // 如果会话未绑定,则进行重连
                    smppSession.reBind();
                    System.out.println(Thread.currentThread().getName()
本文内容来源于互联网,如有侵权请联系删除。
作者最新文章
编程开发
相关文章 更多
using namespace 使用中遇到的问题怎么解决
using namespace 使用中遇到的问题怎么解决

命名空间的基本概念与常见引入问题在C++等编程语言中,命名空间(namespace)是一种将代码标识符(如变量、函数、类名)封装在特定名称下的机制,其主要目的是避免命名冲突,尤其是在大型项目或使用多个第三方库时。使用“using namespace”指令可以将指定命名空间中的所有名称引入当前作用域,

c语言函数递归 实操经验总结:这些技巧很实用
c语言函数递归 实操经验总结:这些技巧很实用

理解递归的基本原理在C语言中,递归是一种函数调用自身的编程技术。要掌握它,首先需要理解其核心思想:将一个复杂的大问题,分解为一个或几个与原问题相似但规模更小的子问题,直到子问题足够简单,可以直接求解。这个过程通常包含两个关键部分:递归出口和递归体。递归出口定义了问题何时不再继续分解,即最简单、可直接

c语言函数递归 怎么选?常见方案对比分析
c语言函数递归 怎么选?常见方案对比分析

递归函数的基本概念与适用场景在C语言编程中,递归是一种函数调用自身的编程技巧。它并非适用于所有问题,但在处理某些具有自相似结构的问题时,能提供极其清晰和优雅的解决方案。递归的核心思想是将一个大规模问题分解为一个或多个同类型但规模更小的子问题,直到子问题简单到可以直接求解。典型的适用场景包括树形结构的

Objective-C 内存管理入门:从 alloc 到 dealloc 的生命周期详解
Objective-C 内存管理入门:从 alloc 到 dealloc 的生命周期详解

理解内存管理的基石在Objective-C的编程世界中,内存管理是开发者必须掌握的核心技能之一。它直接关系到应用的性能、稳定性与资源利用效率。与一些采用自动垃圾回收机制的语言不同,Objective-C在很长一段时间里,依赖一套基于引用计数的、需要开发者部分介入的管理规则。这套规则的核心思想是明确的

如何正确使用 dealloc 以避免 iOS 应用中的内存泄漏
如何正确使用 dealloc 以避免 iOS 应用中的内存泄漏

理解 dealloc 的角色与时机在 iOS 应用开发中,内存管理是保障应用性能与稳定性的基石。dealloc 方法是 Objective-C 中对象生命周期结束时的关键回调,它标志着对象即将被系统回收内存。正确理解其触发时机至关重要:当一个对象的引用计数降为零时,运行时系统会自动调用该对象的 de

深入理解 Objective-C 中的 dealloc 方法:内存管理核心机制
深入理解 Objective-C 中的 dealloc 方法:内存管理核心机制

内存管理的基石在Objective-C的世界里,内存管理是开发者必须掌握的核心技能之一。作为一门在手动引用计数(MRC)时代诞生的语言,Objective-C要求程序员对对象的生命周期有清晰的认识。dealloc方法正是这一生命周期中至关重要的终点站。它是一个实例方法,当对象的引用计数降为零时,系统

理解 native2ascii:Java 国际化开发中的字符编码工具
理解 native2ascii:Java 国际化开发中的字符编码工具

native2ascii 工具的基本定位在Ja va应用程序的国际化与本地化开发过程中,处理非拉丁字符集是一个常见且关键的环节。Ja va内部使用Unicode字符集来统一表示全球各种语言的文字,但其属性文件(.properties)在历史上要求使用ASCII编码,或者更准确地说,要求非ASCII字

如何使用 native2ascii 转换中文字符为 Unicode 转义序列
如何使用 native2ascii 转换中文字符为 Unicode 转义序列

理解 native2ascii 工具的基本用途在软件开发,特别是涉及国际化处理的场景中,开发者常常需要处理不同编码的文本资源。native2ascii 是 Ja va 开发工具包(JDK)中提供的一个命令行实用程序,其主要功能是将包含本地字符编码(非ASCII字符)的文件,转换为包含 Unicode

Java native2ascii 命令详解:解决属性文件乱码问题
Java native2ascii 命令详解:解决属性文件乱码问题

native2ascii 命令的由来与作用在Ja va开发中,处理国际化资源文件是一个常见需求。资源文件通常以.properties格式存储,用于支持多语言界面。然而,Ja va属性文件默认采用ISO-8859-1字符集编码,这导致了一个直接的问题:当文件中包含非拉丁字符(如中文、日文、韩文等)时,

一个 memwatch 实战案例:定位野指针问题
一个 memwatch 实战案例:定位野指针问题

内存监控工具的价值与挑战在软件开发,尤其是使用C/C++这类手动管理内存的语言时,内存错误是程序员最常遭遇的难题之一。其中,野指针问题因其隐蔽性和破坏性,往往成为最难定位的“幽灵”缺陷。它可能潜伏在代码中,在特定条件下才被触发,导致程序崩溃、数据损坏或难以预测的行为。传统的调试手段,如打印日志或使用

查看更多
精品专题 更多
装机必备
装机必备

正软商城装机必备专区,精选办公、浏览器、安全防护、影音播放、压缩解压、设计创作和系统工具等电脑常用正版软件,帮助用户快速完成新电脑软件配置。

Windows
Windows

正软商城Windows软件专区,汇集适用于Windows电脑的办公、设计、安全防护、影音播放、开发工具和系统优化软件,提供软件介绍、系统要求、正版授权及购买下载服务。

macOS软件
macOS软件

正软商城macOS软件专区,精选适用于Mac电脑的办公、设计、影音、效率、开发和系统工具,提供软件功能介绍、macOS兼容版本、正版授权及购买下载服务。

Mac软件 更多
灵活计算器
灵活计算器
macOS/iOS/Android

灵活计算器是一款笔记式算数应用,支持实时计算、动态关联和云端同步功能。记录、整理和输出之间的过渡会更自然,适合长期写作、做笔记或持续沉淀个人内容。

赤友清理大师
赤友清理大师
macOS

赤友清理大师是一款为 Mac 设计的智能清理优化工具,可精准扫描垃圾、大文件、重复文件等,释放磁盘空间。做扫描整理、文字提取和表格转换时,它能把识别后的处理步骤接得更顺,资料录入这类场景会省下不少时间。

极度公式
极度公式
Windows/macOS/Linux

极度公式是一款跨平台专业LaTeX公式识别编辑软件,支持OCR公式识别和多平台编辑。和使用说明,避免使用,享受完整功能与稳定支持。做扫描整理、文字提取和表格转换时,它能把识别后的处理步骤接得更顺,资料录入这类场景会省下不少时间。

WINDOWS 更多
Windows 10
Windows 10
Windows

Windows 10 是一款微软推出的经典操作系统,拥有硬件兼容性与多任务处理能力。它更偏向把系统状态查看和常用调节动作放在一起,适合需要持续观察和微调设备状态的场景。

极度公式
极度公式
Windows/macOS/Linux

极度公式是一款跨平台专业LaTeX公式识别编辑软件,支持OCR公式识别和多平台编辑。和使用说明,避免使用,享受完整功能与稳定支持。做扫描整理、文字提取和表格转换时,它能把识别后的处理步骤接得更顺,资料录入这类场景会省下不少时间。

密码键盘
密码键盘
Windows/macOS/iOS/Android

密码键盘是一款兼具安全性与便捷性的高效密码管理器。日常使用里的持续防护和信息管理会更突出,适合把安全控制放进长期使用流程中的场景。