在订单、支付、结算、库存等业务中,经常会遇到这样的操作:
先修改本地数据库;
再发送 MQ 消息,或者调用外部 HTTP 服务;
两个动作必须在业务上保持一致。
问题在于,数据库事务只能可靠地约束数据库操作,无法天然覆盖 RabbitMQ、Kafka 或 HTTP 请求。业务数据已经提交,但外部通知失败,就会出现“库里成功、下游没收到”的不一致;反过来,如果先通知外部系统,再提交数据库,又可能出现“下游已处理、数据库却回滚”的情况。
本文实现一个可复用的 Local Task Message 本地任务消息组件 :在业务事务中同时写入业务数据和本地任务消息,事务提交后通过 Spring Event 触发快速通知;通知失败时,再由定时任务扫描本地消息表并持续补偿,最终达到业务一致。
本文所有架构图和流程图均使用 Mermaid 绘制。Hexo 主题需要开启 Mermaid 支持,否则代码块只会按文本显示。
一、问题本质:数据库事务管不到外部世界 假设创建订单后,需要通知库存服务:
1 2 3 4 5 @Transactional public void createOrder (CreateOrderCommand command) { orderRepository.insert(command.toOrder()); rabbitTemplate.convertAndSend("order.exchange" , "order.created" , command); }
这段代码看起来很自然,但存在多个失败窗口:
数据库写入成功,MQ 发送失败;
MQ 已经接收消息,但本地事务随后回滚;
MQ 客户端超时,生产者不知道 Broker 到底有没有收到;
HTTP 返回超时,但远端实际上已经完成处理;
应用在事务提交后、通知发送前宕机。
因此,这类问题不能简单依赖一次方法调用解决。更合理的思路是:
将“需要对外执行的动作”先可靠地记录到本地数据库,再异步执行;失败后可以从数据库恢复并重试。
本地任务消息组件提供的是 至少一次投递与最终一致性 ,并不承诺严格的全局强一致,也无法天然保证外部系统只执行一次。因此,下游接口仍然必须支持幂等。
二、组件目标 该组件主要解决以下问题:
业务数据与任务消息在同一个本地事务中写入;
支持 MQ 和 HTTP 两种通知方式,并可继续扩展 Kafka、RPC、Webhook 等通道;
支持直接调用组件服务,也支持通过自定义注解无侵入接入;
使用 Spring Event 作为事务提交后的快速触发通道;
使用本地任务表作为可靠事实来源,应用重启后仍可恢复;
使用“门牌号”对任务分片,允许多个 Job 并行扫描;
通知失败后记录状态并定时补偿;
上游业务系统只需要建表、引入依赖并完成配置。
不适合直接使用该方案的场景包括:
跨多个数据库要求原子提交;
要求严格的同步强一致;
不允许重复调用,且外部系统又无法提供幂等能力;
每秒几十万级消息,需要专业消息日志或流处理平台承载。
三、总体架构 flowchart LR
subgraph UP["上游业务系统"]
BIZ["业务方法<br/>createOrder / settleBill"]
AOP["@LocalTaskMessage<br/>AOP 切面"]
API["ILocalTaskMessageHandleService<br/>直接调用入口"]
end
subgraph CORE["Local Task Message 组件"]
HANDLE["任务受理服务"]
REPO["任务仓储"]
PUB["Spring Event 发布器"]
LISTENER["事务提交后事件监听器"]
STRATEGY["通知策略路由"]
HTTP["HTTP Gateway"]
MQ["RabbitMQ / Kafka Adapter"]
JOB["分组补偿 Job"]
end
subgraph STORE["业务数据库"]
BIZDB[("业务表")]
TASKDB[("local_task_message")]
end
subgraph EXT["外部系统"]
HTTPAPI["外部 HTTP API"]
BROKER["MQ Broker"]
end
BIZ --> BIZDB
BIZ --> AOP
BIZ --> API
AOP --> HANDLE
API --> HANDLE
HANDLE --> REPO
REPO --> TASKDB
HANDLE --> PUB
PUB --> LISTENER
LISTENER --> STRATEGY
STRATEGY --> HTTP
STRATEGY --> MQ
HTTP --> HTTPAPI
MQ --> BROKER
JOB --> TASKDB
JOB --> STRATEGY
style CORE fill:#e8f5e9,stroke:#2e7d32
style STORE fill:#fff8e1,stroke:#f9a825
style EXT fill:#ffebee,stroke:#c62828
这里需要明确三个角色:
本地消息表是可靠来源 :只要事务提交成功,任务就不会凭空消失;
Spring Event 是快速通道 :正常情况下无需等待定时任务,提交后立即尝试通知;
定时任务是补偿通道 :处理事件丢失、进程崩溃、网络异常和外部服务故障。
Spring Event 本身不是消息中间件。它只能在当前进程内传递事件,进程突然退出时,尚未处理的内存事件会丢失。真正兜底的是数据库里的任务记录。
四、工程分层设计 组件不是一个只包含几个静态方法的工具包,而是一个具备领域行为、存储能力、事件机制和调度能力的小型业务内核。可以按照 DDD 风格划分为 config、domain、infrastructure 和 trigger 四层。
flowchart TB
CONFIG["config<br/>自动配置、属性绑定、AOP"]
TRIGGER["trigger<br/>事件监听、定时任务、外部入口"]
DOMAIN["domain<br/>实体、领域服务、策略、端口接口"]
INFRA["infrastructure<br/>JDBC、Spring Event、HTTP、MQ"]
SPRING["Spring 容器"]
DB[("MySQL")]
OUT["HTTP / MQ"]
CONFIG --> SPRING
TRIGGER --> DOMAIN
DOMAIN --> INFRA
INFRA --> DB
INFRA --> OUT
SPRING --> CONFIG
SPRING --> TRIGGER
SPRING --> DOMAIN
SPRING --> INFRA
style CONFIG fill:#607d8b,color:#fff
style TRIGGER fill:#2196f3,color:#fff
style DOMAIN fill:#f2cf00,color:#111
style INFRA fill:#d0006f,color:#fff
推荐的包结构如下:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 local-task-message ├── config │ ├── LocalTaskMessageAutoConfig.java │ ├── LocalTaskMessageProperties.java │ └── aop │ └── LocalTaskMessageAop.java ├── message │ └── LocalTaskMessage.java ├── domain │ ├── model │ │ └── entity │ │ └── TaskMessageEntityCommand.java │ ├── service │ │ ├── ILocalTaskMessageHandleService.java │ │ ├── ILocalTaskMessageNotifyService.java │ │ └── strategy │ │ ├── INotifyStrategy.java │ │ ├── HTTPNotifyStrategy.java │ │ └── RabbitMQNotifyStrategy.java │ └── adapter │ ├── repository │ │ └── ILocalTaskMessageRepository.java │ ├── event │ │ └── ILocalTaskMessageEvent.java │ └── port │ └── ILocalTaskMessagePort.java ├── infrastructure │ ├── dao │ ├── repository │ ├── event │ ├── gateway │ └── mq └── trigger ├── listener │ └── TaskMessageEventListener.java └── job └── TaskMessageEventJob.java
分层的核心价值不是“目录看起来很高级”,而是让领域逻辑只依赖抽象接口。将来把 Retrofit 换成 OkHttp、RabbitMQ 换成 Kafka、JDBC 换成其他存储实现时,不需要修改核心业务流程。
五、本地任务消息表设计 5.1 表结构 erDiagram
LOCAL_TASK_MESSAGE {
BIGINT id PK "自增主键"
VARCHAR task_id UK "业务任务唯一标识"
VARCHAR task_name "任务名称"
VARCHAR notify_type "http / rabbit_mq"
TEXT notify_config "通知配置 JSON"
INT status "0待处理 1处理中 2完成 3失败"
JSON parameter_json "外部通知参数"
INT house_number "门牌号/分片号"
DATETIME create_time "创建时间"
DATETIME update_time "更新时间"
}
基础建表语句可以设计为:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 CREATE TABLE local_task_message ( id BIGINT UNSIGNED NOT NULL AUTO_INCREMENT COMMENT '自增主键' , task_id VARCHAR (128 ) NOT NULL COMMENT '任务唯一标识' , task_name VARCHAR (128 ) NOT NULL COMMENT '任务名称' , notify_type VARCHAR (32 ) NOT NULL COMMENT '通知类型:http、rabbit_mq' , notify_config TEXT NOT NULL COMMENT '通知配置 JSON' , status TINYINT NOT NULL DEFAULT 0 COMMENT '0待处理 1处理中 2已完成 3失败' , parameter_json JSON NOT NULL COMMENT '通知参数 JSON' , house_number TINYINT NOT NULL COMMENT '门牌号/分片号' , create_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP COMMENT '创建时间' , update_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP COMMENT '更新时间' , PRIMARY KEY (id), UNIQUE KEY uk_local_task_message_task_id (task_id), KEY idx_local_task_message_scan (status, house_number, id) ) ENGINE = InnoDB DEFAULT CHARSET = utf8mb4 COMMENT = '本地任务消息表' ;
task_id 应当具备业务唯一性,可以使用订单号、结算单号、券号或稳定生成的业务 UUID。唯一索引既能防止同一业务重复创建任务,也可以作为外部幂等键。
5.2 状态机 stateDiagram-v2
[*] --> WAITING: 创建任务 status=0
WAITING --> PROCESSING: 抢占执行 status=1
FAILED --> PROCESSING: 补偿重试
PROCESSING --> SUCCESS: 通知成功 status=2
PROCESSING --> FAILED: 通知失败 status=3
SUCCESS --> [*]
状态含义:
状态
名称
说明
0
待处理
已经可靠写入数据库,尚未开始通知
1
处理中
已被某个执行器抢占,防止多个节点重复处理
2
已完成
外部通知成功
3
失败
本次执行失败,等待下一次补偿
生产环境建议继续增加以下字段:
1 2 3 4 5 6 7 8 ALTER TABLE local_task_message ADD COLUMN retry_count INT NOT NULL DEFAULT 0 COMMENT '已重试次数' , ADD COLUMN next_retry_time DATETIME NULL COMMENT '下次重试时间' , ADD COLUMN last_error VARCHAR (2000 ) NULL COMMENT '最近一次异常摘要' , ADD COLUMN version INT NOT NULL DEFAULT 0 COMMENT '乐观锁版本' ;CREATE INDEX idx_local_task_message_retry ON local_task_message (status, house_number, next_retry_time, id);
这些字段可以支持指数退避、最大重试次数、死信处理、异常排查和并发抢占。
5.3 门牌号分片 如果所有 Job 都执行同一条扫描 SQL,那么增加 Job 数量并不会提升吞吐量,只会制造重复竞争。因此,可以根据 task_id 计算一个稳定的门牌号:
1 2 int hash = taskId.hashCode() & Integer.MAX_VALUE;int houseNumber = hash % 10 ;
例如:
group01 负责门牌号 [0, 1, 2, 3];
group02 负责门牌号 [4, 5, 6, 7, 8, 9]。
这样每个任务组只处理自己的数据分片,可以水平扩展扫描能力。
六、任务命令模型 任务命令需要描述“做什么、怎么通知、通知什么内容”。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 @Data @NoArgsConstructor @AllArgsConstructor public class TaskMessageEntityCommand { private Long id; private String taskId; private String taskName; private String notifyType; private NotifyConfig notifyConfig; private Integer status; private String parameterJson; private Integer houseNumber; @Data @Builder @NoArgsConstructor @AllArgsConstructor public static class NotifyConfig { private MQ mq; private HTTP http; @Data @Builder @NoArgsConstructor @AllArgsConstructor public static class MQ { private String topic; private String exchange; } @Data @Builder @NoArgsConstructor @AllArgsConstructor public static class HTTP { private String url; private String method; private String contentType; private String authorization; } } }
通知类型枚举同时保存对应的 Spring 策略 Bean 名称:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 @Getter @AllArgsConstructor public enum TaskNotifyEnum { HTTP("http" , "httpNotifyStrategy" , "HTTP 通知" ), RABBIT_MQ("rabbit_mq" , "rabbitMQNotifyStrategy" , "RabbitMQ 通知" ); private final String type; private final String strategy; private final String description; public static TaskNotifyEnum of (String type) { return Arrays.stream(values()) .filter(item -> item.type.equalsIgnoreCase(type)) .findFirst() .orElseThrow(() -> new IllegalArgumentException ( "Unsupported notify type: " + type)); } }
需要在组件入口处校验 notifyType 与 notifyConfig 是否匹配。例如,notifyType=http 时必须配置 notifyConfig.http,避免任务进入数据库后才发现参数不完整。
七、两种接入方式 flowchart LR
BIZ["业务方法"] --> MODE{接入方式}
MODE -->|注解方式| ANN["@LocalTaskMessage"]
MODE -->|主动调用| CALL["handleService.acceptTaskMessage"]
ANN --> HANDLE["任务受理服务"]
CALL --> HANDLE
HANDLE --> SAVE["保存本地任务"]
SAVE --> EVENT["发布 Spring Event"]
7.1 主动调用组件服务 这种方式最直接,调用关系清晰,适合需要显式控制任务构造过程的场景。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 @Service @RequiredArgsConstructor public class OrderApplicationService { private final OrderRepository orderRepository; private final ILocalTaskMessageHandleService handleService; @Transactional public void createOrder (CreateOrderCommand request) { orderRepository.insert(request.toOrder()); TaskMessageEntityCommand command = new TaskMessageEntityCommand (); command.setTaskId("ORDER_CREATED:" + request.getOrderId()); command.setTaskName("订单创建通知" ); command.setNotifyType(TaskNotifyEnum.RABBIT_MQ.getType()); command.setStatus(0 ); command.setNotifyConfig( TaskMessageEntityCommand.NotifyConfig.builder() .mq(TaskMessageEntityCommand.NotifyConfig.MQ.builder() .exchange("order.exchange" ) .topic("order.created" ) .build()) .build()); command.setParameterJson(JSON.toJSONString(request)); handleService.acceptTaskMessage(command); } }
7.2 注解方式 自定义注解用于从方法参数中提取 TaskMessageEntityCommand,切面负责判断事务并调用组件服务。
1 2 3 4 5 6 7 8 9 10 @Retention(RetentionPolicy.RUNTIME) @Target(ElementType.METHOD) @Documented public @interface LocalTaskMessage { String entityAttributeName () default "" ; }
使用示例:
1 2 3 4 5 @Transactional @LocalTaskMessage(entityAttributeName = "command") public void createOrder (TaskMessageEntityCommand command) { orderRepository.insert(buildOrder(command)); }
嵌套参数示例:
1 2 3 4 5 @Transactional @LocalTaskMessage(entityAttributeName = "request.command") public void createOrder (CreateOrderRequest request) { orderRepository.insert(request.toOrder()); }
切面的核心逻辑是:
执行业务方法;
从方法参数中解析任务命令;
若当前已经存在事务,则加入当前事务;
若当前没有事务,则通过 TransactionTemplate 新建事务;
保存本地任务并发布事件;
任意一步失败都必须抛出异常,让业务数据和任务记录一起回滚。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 @Aspect @Component @RequiredArgsConstructor public class LocalTaskMessageAop { private final ILocalTaskMessageHandleService handleService; private final TransactionTemplate transactionTemplate; @Around("@annotation(localTaskMessage)") public Object around ( ProceedingJoinPoint joinPoint, LocalTaskMessage localTaskMessage) throws Throwable { if (TransactionSynchronizationManager.isActualTransactionActive()) { Object result = joinPoint.proceed(); TaskMessageEntityCommand command = resolveCommand( joinPoint, localTaskMessage.entityAttributeName()); handleService.acceptTaskMessage(command); return result; } try { return transactionTemplate.execute(status -> { try { Object result = joinPoint.proceed(); TaskMessageEntityCommand command = resolveCommand( joinPoint, localTaskMessage.entityAttributeName()); handleService.acceptTaskMessage(command); return result; } catch (Throwable throwable) { status.setRollbackOnly(); throw new LocalTaskMessageException (throwable); } }); } catch (LocalTaskMessageException exception) { throw exception.getCause(); } } }
这里有三个容易被忽略的工程细节:
解析 Java 参数名通常需要编译时开启 -parameters;
AOP 与 @Transactional 都是代理增强,需要明确切面顺序;
handleService 内部绝不能只记录异常然后吞掉,否则业务事务仍会提交,消息记录却可能没有写入。
对于公共组件,长期更推荐使用 Spring ParameterNameDiscoverer 或 SpEL 解析属性路径,而不是自己维护大量反射代码。
八、同一事务内保存业务数据与本地任务 主流程如下:
sequenceDiagram
autonumber
actor User as 调用方
participant Biz as 业务方法
participant Aop as 注解切面/直接入口
participant Handle as 任务受理服务
participant Repo as 任务仓储
participant DB as MySQL
participant Publisher as 事件发布器
participant Listener as 提交后监听器
User->>Biz: 执行业务操作
Biz->>DB: 写入业务数据
Biz->>Aop: 提交任务命令
Aop->>Handle: acceptTaskMessage(command)
Handle->>Repo: saveTaskMessage(command)
Repo->>DB: INSERT local_task_message
DB-->>Repo: success
Handle->>Publisher: publishEvent(command)
Publisher-->>Handle: return
Handle-->>Biz: return
Biz-->>User: 本地事务提交
Publisher-->>Listener: AFTER_COMMIT 触发异步通知
8.1 使用 Spring JDBC,而不是裸连接 组件为了减少与上游 MyBatis 版本的冲突,可以只依赖 Spring JDBC。需要注意,若直接使用:
1 Connection connection = dataSource.getConnection();
可能绕开 Spring 当前事务绑定的连接。生产实现更适合使用 JdbcTemplate,或者显式使用 DataSourceUtils.getConnection(dataSource)。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 @Repository @RequiredArgsConstructor public class TaskMessageDao { private final JdbcTemplate jdbcTemplate; public int insert (TaskMessagePO po) { String sql = "INSERT INTO local_task_message " + "(task_id, task_name, notify_type, notify_config, " + "status, parameter_json, house_number, create_time, update_time) " + "VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)" ; return jdbcTemplate.update( sql, po.getTaskId(), po.getTaskName(), po.getNotifyType(), po.getNotifyConfig(), po.getStatus(), po.getParameterJson(), po.getHouseNumber(), po.getCreateTime(), po.getUpdateTime()); } }
领域仓储负责完成领域对象到持久化对象的转换:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 @Repository @RequiredArgsConstructor public class LocalTaskMessageRepository implements ILocalTaskMessageRepository { private final TaskMessageDao taskMessageDao; @Override public void saveTaskMessage (TaskMessageEntityCommand command) { TaskMessagePO po = new TaskMessagePO (); po.setTaskId(command.getTaskId()); po.setTaskName(command.getTaskName()); po.setNotifyType(command.getNotifyType()); po.setNotifyConfig(JSON.toJSONString(command.getNotifyConfig())); po.setStatus(0 ); po.setParameterJson(command.getParameterJson()); po.setHouseNumber(calculateHouseNumber(command.getTaskId())); po.setCreateTime(LocalDateTime.now()); po.setUpdateTime(LocalDateTime.now()); int affectedRows = taskMessageDao.insert(po); if (affectedRows != 1 ) { throw new IllegalStateException ( "Save task message failed, taskId=" + command.getTaskId()); } } private int calculateHouseNumber (String taskId) { return (taskId.hashCode() & Integer.MAX_VALUE) % 10 ; } }
任务受理服务保持简单,并让异常继续向外传播:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 @Service @RequiredArgsConstructor public class LocalTaskMessageHandleService implements ILocalTaskMessageHandleService { private final ILocalTaskMessageRepository repository; private final ILocalTaskMessageEvent event; @Override public void acceptTaskMessage (TaskMessageEntityCommand command) { validate(command); repository.saveTaskMessage(command); event.publishEvent(command); } }
要真正加入同一个事务,业务表与 local_task_message 必须使用:
同一个 DataSource;
同一个 PlatformTransactionManager;
同一条事务传播链路。
如果业务数据和任务消息分别写入两个数据库,这个方案就不再具备本地事务原子性。
九、Spring Event:快速触发,而不是可靠存储 9.1 发布事件 1 2 3 public interface ILocalTaskMessageEvent { void publishEvent (TaskMessageEntityCommand command) ; }
1 2 3 4 5 6 7 8 9 @Getter public class SpringTaskMessageEvent { private final TaskMessageEntityCommand command; public SpringTaskMessageEvent (TaskMessageEntityCommand command) { this .command = command; } }
1 2 3 4 5 6 7 8 9 10 11 @Component @RequiredArgsConstructor public class LocalTaskMessageEvent implements ILocalTaskMessageEvent { private final ApplicationEventPublisher eventPublisher; @Override public void publishEvent (TaskMessageEntityCommand command) { eventPublisher.publishEvent(new SpringTaskMessageEvent (command)); } }
9.2 在事务提交后处理 基线实现可以使用 @EventListener + @Async,但异步线程可能在事务真正提交前开始执行。更稳妥的方式是使用:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 @Slf4j @Component @RequiredArgsConstructor public class TaskMessageEventListener { private final ILocalTaskMessageNotifyService notifyService; @Async("localTaskMessageExecutor") @TransactionalEventListener(phase = TransactionPhase.AFTER_COMMIT) public void handle (SpringTaskMessageEvent event) { TaskMessageEntityCommand command = event.getCommand(); try { notifyService.notify(command); } catch (Exception exception) { log.error("Task notify failed, taskId={}" , command.getTaskId(), exception); } } }
这样可以保证:
本地事务回滚时,不执行外部通知;
本地事务提交成功后,立即异步通知;
通知失败不会反向影响已经提交的业务事务;
即使事件尚未处理应用就崩溃,本地任务仍会被 Job 扫描。
不要把通用组件的异步任务扔进默认线程池。应当提供独立线程池,并暴露核心参数:
1 2 3 4 5 6 7 8 9 10 11 12 @Bean("localTaskMessageExecutor") public Executor localTaskMessageExecutor () { ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor (); executor.setCorePoolSize(4 ); executor.setMaxPoolSize(16 ); executor.setQueueCapacity(1000 ); executor.setThreadNamePrefix("local-task-message-" ); executor.setRejectedExecutionHandler( new ThreadPoolExecutor .CallerRunsPolicy()); executor.initialize(); return executor; }
线程池拒绝并不意味着任务永久丢失,因为数据库记录还在;但拒绝次数必须监控,否则补偿任务会慢慢变成“真正的主流程”。
十、通知策略:HTTP 与 MQ 解耦 不同通知方式的配置和执行逻辑差异较大,不适合塞进一个巨大的 if...else。策略模式可以把变化隔离到独立实现中。
classDiagram
direction LR
class INotifyStrategy {
<<interface>>
+notify(command) String
}
class HTTPNotifyStrategy {
-ILocalTaskMessagePort port
-ILocalTaskMessageRepository repository
+notify(command) String
}
class RabbitMQNotifyStrategy {
-ILocalTaskMessagePort port
-ILocalTaskMessageRepository repository
+notify(command) String
}
class LocalTaskMessageNotifyService {
-Map~String, INotifyStrategy~ strategies
+notify(command) String
}
class ILocalTaskMessagePort {
<<interface>>
+notify2http(command) String
+notify2rabbitmq(command) String
}
INotifyStrategy <|.. HTTPNotifyStrategy
INotifyStrategy <|.. RabbitMQNotifyStrategy
LocalTaskMessageNotifyService --> INotifyStrategy
HTTPNotifyStrategy --> ILocalTaskMessagePort
RabbitMQNotifyStrategy --> ILocalTaskMessagePort
策略接口:
1 2 3 public interface INotifyStrategy { String notify (TaskMessageEntityCommand command) throws Exception; }
领域通知服务通过 Spring 自动注入所有策略 Bean:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 @Service @RequiredArgsConstructor public class LocalTaskMessageNotifyService implements ILocalTaskMessageNotifyService { private final Map<String, INotifyStrategy> strategies; @Override public String notify (TaskMessageEntityCommand command) throws Exception { TaskNotifyEnum notifyType = TaskNotifyEnum.of(command.getNotifyType()); INotifyStrategy strategy = strategies.get(notifyType.getStrategy()); if (strategy == null ) { throw new IllegalStateException ( "Notify strategy not found: " + notifyType.getStrategy()); } return strategy.notify(command); } }
10.1 HTTP 策略 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 @Slf4j @Component("httpNotifyStrategy") @RequiredArgsConstructor public class HTTPNotifyStrategy implements INotifyStrategy { private final ILocalTaskMessagePort port; private final ILocalTaskMessageRepository repository; @Override public String notify (TaskMessageEntityCommand command) throws Exception { if (!repository.tryMarkProcessing(command.getTaskId())) { return "task already claimed" ; } try { String result = port.notify2http(command); repository.markSuccess(command.getTaskId()); return result; } catch (Exception exception) { repository.markFailed(command.getTaskId(), exception.getMessage()); throw exception; } } }
HTTP 调用应设置:
连接超时、读取超时和整体调用超时;
最大响应体限制;
允许的 URL、协议和域名白名单,避免 SSRF;
Idempotency-Key 或稳定业务键;
敏感 Header 加密存储或通过配置引用,不要明文长期落库。
10.2 RabbitMQ 策略 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 @Slf4j @Component("rabbitMQNotifyStrategy") @RequiredArgsConstructor public class RabbitMQNotifyStrategy implements INotifyStrategy { private final ILocalTaskMessagePort port; private final ILocalTaskMessageRepository repository; @Override public String notify (TaskMessageEntityCommand command) throws Exception { if (!repository.tryMarkProcessing(command.getTaskId())) { return "task already claimed" ; } try { String result = port.notify2rabbitmq(command); repository.markSuccess(command.getTaskId()); return result; } catch (Exception exception) { repository.markFailed(command.getTaskId(), exception.getMessage()); throw exception; } } }
RabbitMQ 适配器可以使用可选依赖,避免没有配置 RabbitMQ 的业务系统启动失败:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 @Slf4j @Component public class RabbitMQEvent { @Autowired(required = false) private RabbitTemplate rabbitTemplate; public void publish (String exchange, String routingKey, String message) { if (rabbitTemplate == null ) { throw new IllegalStateException ("RabbitTemplate is not configured" ); } rabbitTemplate.convertAndSend(exchange, routingKey, message, msg -> { msg.getMessageProperties() .setDeliveryMode(MessageDeliveryMode.PERSISTENT); return msg; }); } }
仅仅把消息设置为持久化还不等于发送可靠。生产环境还应配置 Publisher Confirm、Return Callback、交换机与队列持久化,并根据确认结果决定是否将本地任务标记为成功。
十一、定时补偿与门牌号扫描 完整补偿流程如下:
sequenceDiagram
autonumber
participant Job as TaskMessageEventJob
participant Repo as 任务仓储
participant DB as local_task_message
participant Notify as 通知领域服务
participant HTTP as HTTP API
participant MQ as MQ Broker
Job->>Repo: 按任务组和门牌号查询可重试任务
Repo->>DB: status in (0,3) and next_retry_time <= now
DB-->>Repo: 返回 limit 条任务
Repo-->>Job: TaskMessageEntityCommand 列表
loop 每条任务
Job->>Repo: CAS 抢占 status 0/3 -> 1
alt 抢占成功
Job->>Notify: notify(command)
alt HTTP
Notify->>HTTP: 发起请求,携带幂等键
HTTP-->>Notify: response
else MQ
Notify->>MQ: 发送持久化消息
MQ-->>Notify: confirm
end
Notify->>Repo: 成功改为 2,失败改为 3
else 已被其他节点处理
Job-->>Job: 跳过
end
end
11.1 任务组配置 1 2 3 4 5 6 7 8 9 10 11 12 13 14 xfg: wrench: task: config: groups: - groupId: group01 houseNumbers: [0 , 1 , 2 , 3 ] cron: "0/10 * * * * ?" limit: 100 - groupId: group02 houseNumbers: [4 , 5 , 6 , 7 , 8 , 9 ] fixedDelayMs: 5000 limit: 50
属性对象:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 @Data @ConfigurationProperties(prefix = "xfg.wrench.task.config") public class LocalTaskMessageProperties { private List<TaskGroupConfig> groups = new ArrayList <>(); @Data public static class TaskGroupConfig { private String groupId = "default" ; private List<Integer> houseNumbers = new ArrayList <>(); private String cron; private Long fixedDelayMs; private Integer limit = 100 ; } }
11.2 扫描 SQL 推荐让“任务状态”成为扫描条件,而不是只依赖内存中的 lastId:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 SELECT id, task_id, task_name, notify_type, notify_config, status, parameter_json, house_number, create_time, update_timeFROM local_task_messageWHERE house_number IN (?, ?, ?) AND status IN (0 , 3 ) AND (next_retry_time IS NULL OR next_retry_time <= NOW())ORDER BY id LIMIT ?;
然后使用条件更新抢占任务:
1 2 3 4 5 6 7 UPDATE local_task_messageSET status = 1 , update_time = NOW(), version = version + 1 WHERE id = ? AND status IN (0 , 3 ) AND version = ?;
只有影响行数为 1 的执行器才获得处理权。
原始方案中常见的“找到最小 ID,再使用 id > lastId LIMIT N 增量扫描”可以降低查询范围,但存在一个陷阱:如果某条较小 ID 的任务失败,而游标已经向后移动,这条失败任务可能被长期越过。因此:
lastId 不能替代任务状态条件;
失败任务必须能够重新进入扫描集合;
多节点部署时必须有抢占或租约机制;
处理中任务需要超时恢复,防止节点宕机后永久停留在状态 1。
11.3 重试退避 不要每 5 秒无限轰炸一个已经宕机的外部系统。可以使用指数退避:
1 nextRetryTime = now + min(baseDelay * 2^retryCount, maxDelay)
例如:
1 10 秒 -> 30 秒 -> 1 分钟 -> 5 分钟 -> 30 分钟 -> 2 小时
超过最大重试次数后,可以转为“人工处理”或“死信状态”,并触发告警。
十二、端到端流程 flowchart TD
START([业务请求]) --> TX[开启或加入本地事务]
TX --> WRITE_BIZ[写入业务表]
WRITE_BIZ --> BUILD[构造 TaskMessageEntityCommand]
BUILD --> WRITE_TASK[写入 local_task_message]
WRITE_TASK --> PUBLISH[发布 Spring Event]
PUBLISH --> COMMIT{事务是否提交成功}
COMMIT -->|否| ROLLBACK[业务数据与任务记录一起回滚]
COMMIT -->|是| AFTER[AFTER_COMMIT 异步监听]
AFTER --> CLAIM[抢占任务 status -> 1]
CLAIM --> NOTIFY{通知类型}
NOTIFY -->|HTTP| CALL_HTTP[调用外部 HTTP]
NOTIFY -->|MQ| SEND_MQ[发送 MQ 并等待确认]
CALL_HTTP --> RESULT{是否成功}
SEND_MQ --> RESULT
RESULT -->|是| SUCCESS[status -> 2]
RESULT -->|否| FAILED[status -> 3<br/>记录异常和下次重试时间]
FAILED --> JOB[定时任务再次扫描]
JOB --> CLAIM
SUCCESS --> END([流程完成])
ROLLBACK --> END
这套设计的关键不是“保证第一次一定成功”,而是把失败变成一种可观察、可恢复、可重试的正常状态。
十三、自动配置与组件发布 13.1 Spring Boot 2 Spring Boot 2 可以通过 META-INF/spring.factories 注册自动配置:
1 2 org.springframework.boot.autoconfigure.EnableAutoConfiguration =\ cn.example.localtaskmessage.config.LocalTaskMessageAutoConfig
13.2 Spring Boot 3 Spring Boot 3 推荐使用:
1 META-INF/spring/org.springframework.boot.autoconfigure.AutoConfiguration.imports
文件内容:
1 cn.example.localtaskmessage.config.LocalTaskMessageAutoConfig
自动配置类:
1 2 3 4 5 6 7 8 9 10 11 12 @AutoConfiguration @EnableAsync @EnableScheduling @EnableConfigurationProperties(LocalTaskMessageProperties.class) @ComponentScan(basePackages = { "cn.example.localtaskmessage.config.aop", "cn.example.localtaskmessage.domain", "cn.example.localtaskmessage.infrastructure", "cn.example.localtaskmessage.trigger" }) public class LocalTaskMessageAutoConfig { }
更成熟的 Starter 可以进一步增加:
@ConditionalOnProperty:允许关闭组件;
@ConditionalOnClass:仅在 RabbitTemplate 或 HTTP 客户端存在时装配对应适配器;
@ConditionalOnMissingBean:允许业务系统覆盖默认实现;
配置元数据:为 IDE 提供 YAML 自动提示;
独立线程池和调度器;
健康检查与 Actuator 指标。
十四、测试方案 仅验证“日志里收到一条事件”远远不够。至少需要覆盖以下测试矩阵:
测试场景
预期结果
业务写库成功,任务写库成功
两张表同时提交
业务写库异常
业务表与任务表同时回滚
任务写库异常
业务事务回滚
事务回滚
AFTER_COMMIT 监听器不执行
HTTP 调用成功
状态变为 2
HTTP 调用超时
状态变为 3,并生成下次重试时间
MQ Confirm 成功
状态变为 2
MQ 发送结果不确定
保留可重试状态,下游依赖幂等去重
应用提交事务后立即宕机
重启后 Job 可以扫描并补偿
两个节点同时扫描同一任务
只有一个节点抢占成功
同一任务重复创建
唯一索引拦截或返回已有任务
同一消息重复投递
下游根据 taskId 或业务键幂等处理
可以在集成测试中主动模拟异常:
1 2 3 4 5 6 @Transactional public void createOrderAndRollback (TaskMessageEntityCommand command) { orderRepository.insert(buildOrder(command)); handleService.acceptTaskMessage(command); throw new IllegalStateException ("mock rollback" ); }
断言业务表和 local_task_message 都没有新增数据。
十五、生产环境必须补齐的能力 15.1 幂等 本地消息表通常只能提供“至少一次”语义。以下场景都可能导致重复:
HTTP 服务已经处理成功,但客户端读取响应超时;
MQ Broker 已接收消息,但生产者没有收到 Confirm;
消息处理完成后,更新本地状态失败;
执行器在完成外部调用后宕机。
因此:
HTTP 使用 taskId 作为 Idempotency-Key;
MQ 消息携带稳定 messageId/taskId;
消费方建立业务唯一索引或幂等记录表;
不要依赖“我这里看起来只调用了一次”。分布式系统最擅长把“一次”变成“你猜几次”。
15.2 多节点并发 单机内存 lastId 只适合演示,不适合作为唯一协调机制。集群部署至少要选择一种方式:
数据库条件更新抢占;
SELECT ... FOR UPDATE SKIP LOCKED;
数据库租约字段;
Redis 分布式锁或调度平台分片。
对本地消息表而言,数据库条件更新通常最简单,也最容易和状态机保持一致。
15.3 可观测性 建议输出以下指标:
待处理任务数;
失败任务数;
最老待处理任务的等待时间;
每分钟成功量、失败量和重试量;
HTTP/MQ 各通道耗时;
每个门牌号的任务积压量;
线程池活跃线程、队列长度和拒绝次数。
日志至少包含:
1 task_id、notify_type、house_number、status、retry_count、cost_ms、error_code、trace_id
15.4 数据清理 成功任务不能永久留在业务库。可以按月归档,或者定时删除超过保留期的状态 2 数据。失败和人工处理数据需要更长保留时间。
15.5 安全 notify_config 可能包含 URL、Token、Authorization 等敏感信息:
不要在日志中完整输出;
Token 尽量只保存配置引用,不保存明文;
HTTP URL 应配置白名单;
请求体和响应体需要设置大小限制;
对外调用应配置超时、熔断和并发隔离。
十六、方案边界与演进方向 本地任务消息表适合中等规模、业务数据库与外部副作用之间的最终一致性场景。随着业务增长,可以逐步演进:
将任务执行器与业务服务拆分,但仍读取同一任务表;
使用 CDC/Canal/Debezium 订阅任务表变更,替代高频轮询;
使用 Outbox Pattern,将任务表变成标准事件 Outbox;
将通知配置改为类型安全的版本化协议;
增加管理后台,支持失败查询、人工重试、暂停和重放;
增加租约、死信、告警和全链路追踪;
对 Kafka、Pulsar、Webhook、RPC 泛化调用提供插件式适配器。
无论如何演进,都应该保留一个核心原则:
先可靠记录意图,再执行外部副作用;快速通道负责效率,持久化补偿负责可靠性。
十七、总结 本地任务消息组件把分散在各业务系统中的重复代码抽象为一个通用内核:
业务数据与任务消息在同一个数据库事务中提交;
Spring Event 在事务提交后立即触发通知;
HTTP、MQ 使用策略模式扩展;
本地消息表记录完整状态;
定时任务按门牌号分组扫描并补偿;
外部系统通过幂等设计承受重复投递;
集群环境通过状态抢占、租约或锁避免并发重复处理。
它不是魔法,也不会让分布式系统突然变得听话,但它能把“偶发的数据不一致事故”变成“有状态、有日志、可重试、可追踪的工程流程”。这已经是可靠性设计中非常重要的一步。
本文根据“本地任务消息组件”系列学习材料重新整理,并结合实际工程中的事务连接、事务提交后事件、并发抢占、失败补偿和幂等要求进行了补充。