- 文档
- 教程
- 后端
【免费下载链接】CodeGuide
:books: 本代码库是作者小傅哥多年从事一线互联网 Java 开发的学习历程技术汇总,旨在为大家提供一个清晰详细的学习教程,侧重点更倾向编写Java核心内容。如果本仓库能为您提供帮助,请给予支持(关注、点赞、分享)!
本文基于 CodeGuide 仓库中的《本地任务消息组件》文档体系(主文档 及同目录下的 6 节课程文档),系统讲解该组件如何解决"本地数据库事务 + 外部 MQ/HTTP 调用"的最终一致性难题:通过本地消息表在同一事务内落库、Spring Event 事件驱动异步通知、策略模式分发 HTTP/RabbitMQ 通道、门牌号分片的定时扫描补偿,以及自定义注解 AOP 的轻量化接入方式。读完后,你将掌握一套可直接借鉴到业务系统中的最终一致性组件设计范式与完整功能链路。
一、问题背景:为什么需要本地任务消息组件
在业务功能开发中,存在一个非常普遍的场景:在完成一次数据库写事务的同时,还需要向外部系统发送一条 MQ 消息,或发起一次 HTTP 远程调用。例如订单落库后要通知结算系统、库存扣减后要同步给搜索服务。
问题在于,MQ 消息发送和 HTTP 调用无法与数据库写操作纳入同一个事务:
- 如果先写库、再发消息,消息发送失败(网络超时、服务宕机、流量洪峰、线程阻塞等),业务数据已经落库但外部系统收不到通知,数据不一致;
- 如果先发消息、再写库,写库回滚后消息已经发出,下游会处理一条"幽灵消息";
- 分布式事务(XA/TCC/Saga)引入成本高、性能开销大,多数业务场景并不需要这么重。
在没有统一组件支持之前,各业务系统只能各自实现:自己建一张本地消息表、自己维护消息表的写入、自己写定时任务扫描补偿、自己处理 HTTP/MQ 的发送与重试——重复开发且难以维护,正如 主文档 中所述:"为了完成业务流程的同时,在发送一个 MQ 消息或者远程调用 HTTP 操作,都需要自己写一个本地消息表,之后还要维护消息表的扫描补偿。"
本地任务消息组件(Local Task Message)正是为了把这套通用能力从各业务系统中抽取出来,凝练成一个可被上游系统以 jar 包方式引入的领域服务内核,统一解决这一共性难题。
二、产品概述:组件解决什么问题、如何工作
按照 主文档 的产品方案描述,本地任务消息组件基于 Spring 框架能力实现通用功能内核,便于集成到各类业务系统中,核心工作链路如下:
- 事务内写库:在业务系统的事务内,完成业务数据写库的同时,写入一条本地消息记录(上游系统需在自身数据库中创建符合组件规范的本地消息表,且与业务数据同库,保证同一个事务);
- 事件驱动异步通知:写入完成后,组件同步推送 Spring 事件(
ApplicationEvent),触发事务外的异步处理(@EventListener+@Async),执行 MQ 消息发送或 HTTP 回调; - 分片扫描补偿:即使异步处理失败,组件内置的本地消息表定时任务会持续检测并重试通知,且支持自定义配置"门牌号"(houseNumber)多任务并行扫描,提升扫描吞吐量,确保消息最终一致性和业务流程的可靠执行。
使用方式上,组件提供两种接入途径,用户可以选择:
- 注解方式:对目标方法配置自定义注解
@LocalTaskMessage,切面会自动获取入参,入参需要为TaskMessageEntityCommand对象(也可以是某个入参对象中携带该对象),之后通过req.command配置也可以获取; - 编程方式:直接调用组件内核服务
ILocalTaskMessageHandleService的acceptTaskMessage(taskMessageEntityCommand)方法受理任务消息。
两种方式的好处是:业务项目工程不再需要自己维护本地消息表的写入,以及 MQ/HTTP 的处理和补偿逻辑,全部由组件内核完成。
三、技术架构:DDD 分层与端口-适配器模式
从 主文档 的技术架构章节可以明确组件的架构定位:
- Local Task Message 组件是为解决"本地数据库事务与外部 MQ/HTTP 调用一致性"问题而设计的领域服务内核,让上游业务系统不必在每个流程中做大量重复编码,通过注解或直接调用组件内核服务即可完成消息通知操作;
- 它不是一个单纯的工具性功能,而是一个具备完整领域功能的内核:具备操作数据库表的能力,接收 Spring Event 事件,对接 MQ、HTTP 完成与外部的交互处理;
- 上游系统使用时,只需要配置好对应的本地消息表(一个事务下连的同一个库)、引入组件并完成 yml 配置,即可直接使用。
在工程结构上,组件采用 DDD 分层与端口-适配器模式,清晰划分domain / infrastructure / trigger / config模块(对应仓库 第2节课程文档 中的说明:"内核有点类似把一个业务项目中,从上到下一整套流程,被单独提取出来做成一个独立的小项目,之后引入到上游使用方的系统,就可以运行")。各层职责在后续 6 节课程文档中均有体现:
| 模块 | 职责 | 课程文档佐证 |
|---|---|---|
| domain(领域层) | 任务消息领域服务,按notifyType分发不同的通知操作(http、mq) | 第4节 |
| infrastructure(基础设施层) | 数据库访问(原生 JDBC)、HTTP 调用(Retrofit2)、MQ 推送(RabbitTemplate) | 第3节、第4节 |
| trigger(触发层) | Spring Event 事件监听,接收事件后调用领域层的通知服务方法 | 第2节、第4节 |
| config(配置层) | 自定义注解的切面逻辑、@ConfigurationProperties动态调度配置 | 第6节 |
本仓库(CodeGuide)为文档仓库,组件的工程代码以课程配套项目形式逐步拉分支开发,本节及以下各节的实现描述均以仓库中对应课程文档为准,可通过 docs/md/project/local-task-message/ 目录下的各节文档继续深入。
四、本地消息任务表设计与数据写入
这是组件一致性的基石,对应 第3节课程文档。
1. 设计原则
- 任务表由上游系统自行在数据库中配置:组件不绑定任何特定库表命名,引入组件的上游系统在自己数据库中建表,之后调用组件时以同一个数据源(
DataSource)对库表进行操作,从而保证"业务数据 + 任务消息记录"在同一个本地事务中提交; - 直接使用原生 JDBC 操作数据库,不引入 MyBatis 框架。这样做的目的是:避免上游系统使用组件时的版本兼容问题——最原始的 JDBC 方式,兼容性最好;
- 数据访问以 DAO 封装落地,完成插入、状态更新、分片条件查询、最小游标查询等操作。
2. 关键 DAO 能力
结合 主文档 的学习要点与 第5节课程文档 的补偿流程,任务表的数据操作可以归纳为四类:
- 插入:受理任务消息时,与业务数据同事务写入一条待通知的任务记录;
- 状态更新:HTTP/MQ 通知执行完成后,无论成功还是失败都要回写任务记录状态;
- 分片条件查询:按"门牌号"(houseNumber)分片条件拉取待处理任务列表;
- 最小游标查询:根据条件获取满足条件的最小 id,再以
id > 最小id limit x的方式批量拉取数据,保证扫描推进有序、可续扫。
五、Spring Event 事件消息:解耦事务内写入与事务外通知
对应 第2节课程文档,本节完成了组件的事件驱动骨架:
- 构建本地任务消息组件的工程框架与对应的测试工程服务;
- 使用 Spring Event 事件消息完成"行为触达的通知和监听",用于后续处理外部 HTTP、MQ 的调用操作;
- 把本地消息组件构建成一个 jar 包,让测试工程通过 pom 依赖方式引入使用。
从事件链路上看,组件使用ApplicationEvent+@EventListener+@Async的组合实现解耦通知链路:业务事务提交前,任务消息已随业务数据落库;事件被发布后,监听器在事务之外异步执行外部通知。这一设计的关键收益是——即使通知动作失败,也不会影响业务事务本身;反之,事务未提交的数据,外部也不会提前感知。
六、通知策略处理:HTTP 与 RabbitMQ 双通道
对应 第4节课程文档,本节为组件增加了 HTTP 与 MQ(RabbitMQ)两种通知通道,采用策略模式完成可插拔的通知能力:
- trigger 层监听:事件消息的监听放在 trigger 层,监听后调用领域层的通知服务方法(领域层新增一个按通知类型分发的领域方法);
- 领域服务按
notifyType分发:如http、mq,将来想扩展其他通知通道(如 Kafka、gRPC),只需在此处添加新的策略实现; - 基础设施层完成具体调用:
- HTTP:使用Retrofit2(底层OkHttp3)框架统一封装 HTTP 网关,掌握动态 URL、Header、Body 的组合与异常处理;
- MQ:RabbitMQ 直接使用
RabbitTemplate模板 push 消息。
此外还有一个重要的工程细节:RabbitMQ 事件发布采用可选依赖注入方式,避免上游系统未配置 MQ 时,因强依赖导致应用启动失败。也就是说组件支持"只用 HTTP 不用 MQ"或"只用 MQ 不用 HTTP"的部署形态。
七、动态任务补偿处理:门牌号分片扫描保证最终一致
对应 第5节课程文档,这是组件"最终一致"承诺的兜底机制:
- 先更新状态,再定时补偿:无论 MQ 发送还是 HTTP 调用都有可能失败(网络超时、服务宕机、线程阻塞、流量洪峰等),因此在完成 MQ/HTTP 处理后,首先更新数据库任务表状态(成功或失败),之后再由定时扫描任务做补偿处理;
- 门牌号(houseNumber)分片扫描:为了提高整体扫描效率,设计了"门牌号"机制——可以配置多个定时任务,每个任务只扫描自己门牌号范围内的记录,多个任务并行扫描从而提升扫描吞吐量;
- 游标推进拉取:扫描库表时,先根据条件获取一个符合条件的最小 id,再以
id > 最小id limit x获取数据列表,实现高效的拉取与顺序处理; - 补偿重复与幂等:由于"从 Spring Event 接收消息 → 执行通知(http/mq)→ 更新数据库"这些步骤都不在同一个事务中(即从事务——业务数据 + 任务表数据写入——往后,都是有可能失败的),补偿就可能出现重复,比如 HTTP 被重复调用一次、MQ 被重复发送一次。因此课程文档明确强调:业务方对接这些通知时,一定要做幂等操作(比如以 OrderId 做唯一索引处理)。
八、切面拦截任务操作:@LocalTaskMessage 注解的轻量化接入
对应 第6节课程文档。由于这是一个通用组件项目,使用方式必须足够轻量,因此对编程式调用handleService.acceptTaskMessage(taskMessageEntityCommand)的编码方式,提供了更优雅的注解方案:
- 添加自定义注解
@LocalTaskMessage,并在 config 配置层编写切面逻辑,核心是获取配置了该注解的方法入参,从中拿到TaskMessageEntityCommand任务消息对象; - 事务边界处理:切面会判断当前是否已有事务操作——如果没有则开启一个新事务,如果有则使用同一个事务,完成数据库表数据的插入,随后推送 Spring Event 事件消息,后续流程与编程式调用完全一致。
切面方式的优点在于:更优雅简洁,用户不需要自己维护调用关系;同时结合事务边界的统一处理,保证了"业务数据 + 消息记录"的原子写入。
九、配置驱动:多任务组动态调度与线程池化管理
结合 主文档 的学习要点,组件的调度与配置体系具备以下能力:
@ConfigurationProperties驱动的多任务组动态调度配置:通过 yml 配置即可定义多组扫描任务,每组绑定自己的门牌号(houseNumber)分片范围;- 两种触发方式:支持
cron(固定时间表达式)与fixedDelay(固定延迟)两种触发方式; - 批次大小 limit 可配:每轮扫描拉取的任务数量通过 limit 参数控制,配合"最小 id 游标 + limit"策略,避免单次拉取过多造成压力;
ThreadPoolTaskScheduler线程池化调度管理:合理设置线程名与池大小,提升任务调度的可观测性与稳定性(线程名便于日志定位,池大小决定并发调度能力)。
此外,组件还通过示例命令对象TaskMessageEntityCommand的构建与调用,展示了入参约定、枚举策略(通知类型枚举)与配置对象之间的协作方式;整体上也要求使用者具备异常、日志与枚举的综合使用能力,建立稳定的错误处理机制。
十、能力清单小结
以下是从 主文档 完整继承的"能学到啥"清单,也是本组件覆盖的技术面全景:
- 【架构】掌握 DDD 分层与端口-适配器模式,清晰划分 domain/infrastructure/trigger/config 模块,提升可维护性与扩展性;
- 【后端】学习注解 + AOP 方式受理任务消息,结合事务边界进行统一处理,理解
@LocalTaskMessage与切面配合的落地实践; - 【后端】掌握本地消息表设计与分片扫描策略(按门牌号 houseNumber 分片),实现高效拉取与顺序处理,提升系统可靠性;
- 【后端】熟悉 Spring Event 事件驱动与异步消费,使用
ApplicationEvent+EventListener+Async实现解耦通知链路; - 【后端】实践策略模式实现可插拔通知能力,支持 HTTP 与 RabbitMQ 两种通知通道,并在成功/失败时更新任务状态;
- 【后端】熟练使用 OkHttp3 与 Retrofit2 统一封装 HTTP 网关,掌握动态 URL、Header、Body 的组合与异常处理;
- 【后端】了解 RabbitMQ 事件发布的可选依赖注入方式,避免未配置 MQ 时的强依赖导致应用启动失败;
- 【配置】掌握
@ConfigurationProperties驱动的多任务组动态调度配置,支持 cron 与 fixedDelay 两种触发方式,并可配置批次大小 limit; - 【运维】学习
ThreadPoolTaskScheduler的线程池化调度管理,合理设置线程名与池大小,提升任务调度的可观测性与稳定性; - 【数据】掌握原生 JDBC 访问与 DAO 封装,完成插入、状态更新、分片条件查询、最小游标查询等落地实现;
- 【测试】通过示例命令对象
TaskMessageEntityCommand的构建与调用,理解入参约定、枚举策略与配置对象的协作; - 【实践】提升异常、日志与枚举的综合使用能力,建立稳定的错误处理。
十一、延伸阅读:仓库中的配套课程文档
本组件的完整开发过程在 CodeGuide 仓库中以 6 节课程文档逐步展开,建议按顺序阅读以还原从需求分析到注解拦截的完整实现链路:
- 第1节:组件需求分析——聚焦"数据库事务 + 对外发送消息(MQ)/发起 HTTP 调用"的最终一致性问题,提炼通用技术解决方案;
- 第2节:SpringEvent事件消息——搭建组件工程框架,用 Spring Event 完成事件发布与监听,把组件构建成 jar 供测试工程 pom 引入;
- 第3节:任务表设计和数据写入——设计通用本地消息任务表,引入 DataSource 数据源,以原生 JDBC 完成插入处理,规避 MyBatis 版本兼容问题;
- 第4节:通知策略处理(HTTP&MQ)——接收 Spring Event 监听后,以通知行为策略完成 HTTP 远程调用(Retrofit2)和 MQ 消息推送(RabbitTemplate);
- 第5节:动态任务补偿处理——通知完成后更新任务表状态,定时扫描任务按门牌号分片补偿,并强调下游幂等设计;
- 第6节:切面拦截任务操作——增加自定义注解,config 层编写切面获取入参中的任务消息对象,结合现有事务或新事务完成数据插入并推送 Spring Event。
十二、总结
本地任务消息组件给出的是一套"本地消息表"经典模式的组件化落地:
- 一致性边界清晰——业务数据与消息记录同库同事务写入,事务提交是"消息一定会被投递"的唯一起点;
- 通知链路解耦——Spring Event 事件驱动 + 异步消费,通知动作的失败不影响业务主流程;
- 通道可插拔——策略模式分发 HTTP(Retrofit2/OkHttp3)与 RabbitMQ(RabbitTemplate)两种通道,扩展新通道只需增加策略实现,RabbitMQ 可选注入避免强依赖;
- 补偿可水平扩展——门牌号分片 + 多任务组并行扫描 + 最小 id 游标 + limit 批次控制,保证扫描吞吐量与顺序性;
- 接入足够轻量——
@LocalTaskMessage注解切面自动解析入参、自动识别事务边界,编程式调用作为兜底方案; - 工程兼容性优先——原生 JDBC 而非 ORM 框架、jar 包内核引入、yml 配置驱动,让上游系统"建表 + 引包 + 配置"三步即可使用。
这套设计不依赖任何重量级中间件,仅用 Spring 事件、JDBC、Retrofit2 与 RabbitTemplate 即可完成最终一致性的完整闭环,非常适合作为业务团队自研"业务型通用组件"的参考范式。
- 文档
- 教程
- 后端
【免费下载链接】CodeGuide
:books: 本代码库是作者小傅哥多年从事一线互联网 Java 开发的学习历程技术汇总,旨在为大家提供一个清晰详细的学习教程,侧重点更倾向编写Java核心内容。如果本仓库能为您提供帮助,请给予支持(关注、点赞、分享)!
相关推荐
my-tv 我的电视:装完就能看的电视直播软件,4 个遥控器按键搞定换台
my tv 我的电视:装完就能看的电视直播软件,4 个遥控器按键搞定换台 晚上回家想把老电视当直播入口,却发现市面上的电视直播软件要么绑会员、要么广告铺满界面,
音视频直播分布式事务完全指南:XA、TCC、Saga、本地消息表与可靠消息最终一致性方案对比(doocs/advanced-java)
分布式事务完全指南:XA、TCC、Saga、本地消息表与可靠消息最终一致性方案对比(doocs/advanced java) 分布式事务是微服务与分布式系统面试
文档教程后端gh_mirrors/ps/psr7与Doctrine集成:数据库事务中的HTTP消息
gh_mirrors/ps/psr7与Doctrine集成:数据库事务中的HTTP消息 在Web应用开发中,处理HTTP请求(Request)和响应(Respo
后端
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考