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

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(): 唤醒在该对象上等待的所有线程。

关键点:

立即学习“Java免费学习笔记(深入)”;

必须在synchronized块内调用:wait()、notify()和notifyAll()方法必须在持有对象监视器(即synchronized块所锁定的对象)的情况下调用,否则会抛出IllegalMonitorStateException。操作同一个监视器对象:所有等待和通知操作都必须针对同一个对象进行,这个对象充当了线程间通信的“信号量”。虚假唤醒与条件检查: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)本身是线程安全的,但它并不能阻止在列表为空时尝试移除。

讯飞绘文 讯飞绘文

讯飞绘文:免费AI写作/AI生成文章

讯飞绘文 118 查看详情 讯飞绘文

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. 改进方案与最佳实践

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

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

以上就是Java并发消息发送系统中的会话管理与wait/notify机制深度解析的详细内容,更多请关注创想鸟其它相关文章!

版权声明:本文内容由互联网用户自发贡献,该文观点仅代表作者本人。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。
如发现本站有涉嫌抄袭侵权/违法违规的内容, 请发送邮件至 chuangxiangniao@163.com 举报,一经查实,本站将立刻删除。
发布者:程序猿,转转请注明出处:https://www.chuangxiangniao.com/p/897957.html

(0)
打赏 微信扫一扫 微信扫一扫 支付宝扫一扫 支付宝扫一扫
MongoDB学习(三)MongoDB shell 命令行的使用
上一篇 2025年11月28日 18:51:02
学信网学历认证报告和在线验证报告有什么区别_学信网两类报告对比
下一篇 2025年11月28日 18:51:02

相关推荐

  • 谷歌浏览器官方在线访问 最新版Chrome官网登录

    谷歌浏览器官方在线访问入口是https://www.google.cn/chrome/,提供简洁界面、跨设备同步、高效内核、安全防护和丰富扩展生态。 谷歌浏览器官方在线访问入口在哪里?这是不少网友都关注的,接下来由PHP小编为大家带来最新版Chrome官网登录地址,想要获取纯净浏览体验的网友一起随小…

    2026年9月21日
    100
  • Java Collections.singletonList如何创建单元素集合

    Collections.singletonList(T item) 返回只含一个元素的不可变列表,传入指定对象后生成轻量级只读集合,适用于需高效传递单元素场景。该列表禁止修改操作,否则抛出异常,允许 null 元素,内部优化减少内存开销,常用于 API 参数传递或流处理中的临时数据构造。 Java …

    2026年9月21日
    000
  • JavaScript中的模块联邦如何实现微前端的代码共享?

    模块联邦通过运行时动态加载实现微前端代码共享,无需打包公共依赖。使用 ModuleFederationPlugin 配置 name、remotes、exposes 和 shared,使应用可暴露或引入远程模块,支持组件、工具函数及状态管理共享,提升复用性并减少冗余。 模块联邦通过在构建时让不同应用直…

    2026年9月21日
    200
  • Swoole如何实现一个UDP服务器

    答案:使用Swoole可轻松创建高性能UDP服务器。通过new SwooleServer()设置UDP套接字,监听Packet事件接收数据,利用sendto()回复客户端;结合set()配置worker_num等参数优化性能,配合PHP UDP客户端测试通信,适用于高并发、低延迟场景。 使用Swoo…

    2026年9月21日
    000
  • MySQL执行计划中的Extra字段代表什么_怎么看优化空间?

    MySQL执行计划中的Extra字段代表什么_怎么看优化空间?MySQL执行计划中的Extra字段代表什么_怎么看优化空间?MySQL执行计划中的Extra字段代表什么_怎么看优化空间?MySQL执行计划中的Extra字段代表什么_怎么看优化空间?

    在 mysql 查询优化中,执行计划的 extra 字段用于说明查询执行时的额外操作,常见的值包括:1. using filesort 表示需要额外排序,应尽量通过建立索引避免;2. using temporary 表示使用了临时表,常见于 group by 或复杂 join,需优化减少其使用;3.…

    2026年9月21日 用户投稿
    000
  • 如何通过tracert命令追踪数据包从本地到目标服务器的完整路径?

    打开命令提示符,输入cmd并回车;2. 执行tracert 目标地址命令追踪路径;3. 查看每跳响应时间与IP,分析延迟变化定位网络瓶颈;4. 注意部分节点可能因防火墙不响应导致超时。 使用 tracert(Windows 系统)命令可以追踪数据包从你的计算机到目标服务器所经过的每一跳网络节点,帮助…

    2026年9月21日
    900
  • 如何在Java中理解Java I/O与NIO机制

    传统I/O是阻塞式流模型,适用于低并发场景;NIO基于缓冲区与通道,支持非阻塞和多路复用,适合高并发网络应用,核心区别在于线程模型与资源利用率。 Java中的I/O(输入/输出)与NIO(New I/O)是处理数据读写的核心机制,理解它们的区别和使用场景对开发高性能应用至关重要。传统I/O基于流模型…

    2026年9月21日
    100
  • UC浏览器网页上的文字无法选中复制怎么办 UC浏览器解决网页文字禁止复制问题

    答案:可通过开发者工具、阅读模式、打印预览、OCR识别或自定义脚本解除UC浏览器网页复制限制。具体操作依次为:开启开发者工具并执行JavaScript代码解除限制;启用阅读模式净化页面内容;使用打印预览重新渲染页面以选中文字;对截图应用OCR技术提取文本;添加书签脚本自动移除禁用选择的代码,从而实现…

    2026年9月21日
    100
  • JavaScript中的尾调用优化(TCO)在ES6中如何工作?

    尾调用是指函数的最后一个动作调用另一个函数,ES6引入尾调用优化以重用栈帧、避免内存溢出,支持真正的尾递归,如阶乘函数通过累积参数实现。 尾调用优化(Tail Call Optimization, TCO)是ES6引入的一项语言特性,目的是在特定条件下重用函数调用栈帧,避免不必要的内存增长,从而支持…

    2026年9月21日
    200
  • 抖音蝴蝶号无人直播带货操作流程及注意事项

    抖音蝴蝶号无人直播带货操作流程及注意事项抖音蝴蝶号无人直播带货操作流程及注意事项抖音蝴蝶号无人直播带货操作流程及注意事项抖音蝴蝶号无人直播带货操作流程及注意事项

    “抖音蝴蝶号无人直播带货”是一种通过自动化或半自动化技术实现的直播销售模式。①其核心在于摆脱真人主播限制,实现24小时不间断直播,提升效率与流量利用率;②关键步骤包括明确账号定位与商品选择、准备高质量且丰富的内容素材、利用虚拟人或预录内容实现直播推流、结合智能客服模拟评论区互动;③优势在于降低人力成…

    2026年9月21日 用户投稿
    500
  • Java语法基础有哪些新手必学的核心知识

    掌握Java基本数据类型与变量声明,如int、double、char和boolean,并理解强类型语言特性;2. 熟悉运算符与表达式,包括算术、比较和逻辑运算符,奠定程序逻辑基础。 Java语法基础是每个初学者必须掌握的内容,只有打好根基,才能顺利进阶面向对象编程和实际项目开发。以下是新手必学的核心…

    2026年9月21日
    300
  • 音乐文件占用空间太多怎么办_音乐文件占用空间太多如何整理详细指南

    解决音乐文件占空间问题的关键是压缩与整理:先用软件或在线工具降低比特率压缩体积,再按场景分类、利用元数据自动归集,并通过听歌片段和BPM判断保留内容,避免重复与误删。 音乐文件占空间太多,核心解决办法就两条:一是压缩单个文件体积,二是通过有效分类管理提升使用效率。直接删歌不是长久之计,学会整理和优化…

    2026年9月21日
    000
  • 升级X86架构性能大提升!极空间Z2 Ultra图赏

    升级X86架构性能大提升!极空间Z2 Ultra图赏升级X86架构性能大提升!极空间Z2 Ultra图赏升级X86架构性能大提升!极空间Z2 Ultra图赏升级X86架构性能大提升!极空间Z2 Ultra图赏

    10月23日,极空间正式推出全新双盘位nas产品——极空间z2 ultra,官方售价为1899元,参与国家补贴后仅需1457元,性价比进一步提升。 此次发布的Z2 Ultra最大的亮点在于采用X86架构处理器,相较以往使用的ARM平台,性能实现飞跃式提升,运行速度显著加快。更重要的是,新架构对Doc…

    2026年9月21日 用户投稿
    300
  • 数据库分库分表(Sharding)策略

    在现代应用程序中,随着数据量的增长,单一数据库的性能和容量往往难以满足需求。这时,数据库分库分表(Sharding)策略就成了一个关键的解决方案。那么,如何设计和实现一个有效的分库分表策略呢?让我们深入探讨一下。 在我的职业生涯中,我曾多次参与大型项目的数据库优化,其中分库分表是常见的挑战之一。我记…

    2026年9月21日
    000
  • 如何在Java中实现个人财务管理工具

    首先设计Transaction、FinanceManager和Budget核心类,实现交易记录、统计分析与预算控制功能,通过ArrayList管理数据,使用LocalDate处理日期,结合ObjectOutputStream持久化存储,初期采用Scanner构建控制台菜单实现增删查改与报表展示,后期…

    2026年9月21日
    100
  • 事务隔离级别在mysql中如何应用

    MySQL提供四种事务隔离级别:READ UNCOMMITTED、READ COMMITTED、REPEATABLE READ(默认)、SERIALIZABLE,依次增强数据一致性,分别用于平衡并发性能与脏读、不可重复读、幻读等问题;通过SELECT @@tx_isolation等命令可查看级别,S…

    2026年9月21日
    300
  • X旗下Grok上线即时语音搜索,挑战Google引领搜索新方向

    近日,x平台旗下的ai助手grok正式推出了“即时语音搜索”功能。用户现在可以通过语音直接提问,触发实时网页检索,并迅速获得整合后的精准答案。此举意在优化信息获取流程,推动人机交互向更自然、高效的方向演进。 该语音搜索模式实现了“即说即搜即答”的流畅体验。例如,当用户提出“星舰发射的具体时间是什么?…

    2026年9月21日
    100
  • Laravel应用的安全审计(Security Audit)方法

    进行安全审计对laravel应用至关重要,因为它能发现并修复安全漏洞,提升整体安全性和用户信任度。具体方法包括:1. 代码审查,确保无未过滤输入和弱密码;2. 配置文件安全性,保护敏感信息;3. 依赖管理,更新第三方包;4. 用户认证和授权,防止未授权访问;5. 日志和监控,检测异常行为。 在讨论L…

    2026年9月21日
    100
  • Laravel 8 登录后重定向到仪表盘的全面指南

    本文深入探讨了 Laravel 8 中用户登录后重定向到仪表盘的多种策略。我们将详细解析默认的重定向机制,包括 LoginController 和 RedirectIfAuthenticated 中间件,并重点介绍如何通过自定义登录逻辑实现精确的重定向控制,同时提供示例代码和常见问题排查建议,确保用…

    2026年9月21日
    000
  • Guava Multimap:高效获取并打印指定键的所有关联值

    guava multimap是处理一键多值映射关系的强大工具。要获取特定键的所有关联值,应直接使用其提供的`multimap#get(k)`方法。该方法会返回一个包含所有匹配值的`collection`,即使键不存在,也会返回一个空集合而非`null`,从而简化了值检索和空值处理逻辑,是比手动迭代键…

    2026年9月21日
    100

发表回复

登录后才能评论
关注微信