如何在Spring Boot中使用观察者模式优化转账监控服务的架构?
嘿,刚好我正在开发一套转账监控系统,结合Spring Boot用观察者模式来优化架构的思路,完全适配你提到的动态构建SQL、异步发送外部服务,还有未来不断新增监控条件的需求,我给你详细说说:
核心思路背景
我这套系统的核心逻辑是:根据不同的监控条件动态构建SQL查询语句,找到匹配的转账记录后,异步发送给外部的Markaz服务。而且未来会不断新增监控规则,每加一个新规则就对应一个独立的Observer类,把规则逻辑封装起来,不影响现有代码。
具体落地步骤(Spring Boot环境下)
第一步:定义观察者接口
先抽象出一个观察者的统一接口,比如TransferMonitorObserver,里面定义两个核心方法:一个用来生成对应的SQL条件片段,另一个可以用来处理匹配后的转账(如果需要的话)。示例代码大概是这样:public interface TransferMonitorObserver { // 生成当前规则对应的SQL条件片段 String buildQueryCriteria(); // 可选:处理匹配到的转账记录(如果每个规则有特殊处理逻辑) void handleMatchedTransfers(List<Transfer> transfers); }第二步:实现具体的观察者类
每一个监控规则就对应一个具体的观察者实现,比如大额转账监控、跨境转账监控,每个类里只关注自己的规则逻辑,完全独立。比如大额转账的观察者:@Component public class HighAmountTransferObserver implements TransferMonitorObserver { @Override public String buildQueryCriteria() { // 封装大额转账的SQL条件 return "amount > 10000 AND currency = 'USD'"; } @Override public void handleMatchedTransfers(List<Transfer> transfers) { // 如果这个规则有特殊处理,比如打标,就写在这里 transfers.forEach(transfer -> transfer.setTag("HIGH_AMOUNT")); } }用
@Component把这些类注册到Spring容器,后面主题类可以自动收集所有观察者。第三步:定义主题类(统一协调者)
创建一个主题类,比如TransferMonitorSubject,它负责维护所有观察者,触发SQL构建、执行查询,以及后续的异步发送逻辑。在Spring里把它做成单例组件:@Component public class TransferMonitorSubject { // Spring自动注入所有实现了TransferMonitorObserver的Bean @Autowired private List<TransferMonitorObserver> observers; // 注入转账DAO用来执行SQL查询 @Autowired private TransferDao transferDao; // 注入异步发送服务 @Autowired private MarkazAsyncSender markazAsyncSender; public void triggerMonitor() { // 1. 收集所有观察者的SQL条件,拼接成完整的WHERE子句 StringBuilder criteriaBuilder = new StringBuilder(); for (TransferMonitorObserver observer : observers) { if (criteriaBuilder.length() > 0) { criteriaBuilder.append(" OR "); } criteriaBuilder.append("(").append(observer.buildQueryCriteria()).append(")"); } String whereClause = criteriaBuilder.length() > 0 ? criteriaBuilder.toString() : "1=1"; // 2. 执行SQL查询,获取匹配的转账记录 List<Transfer> matchedTransfers = transferDao.queryTransfers(whereClause); // 3. 通知所有观察者处理记录(如果有需要) for (TransferMonitorObserver observer : observers) { observer.handleMatchedTransfers(matchedTransfers); } // 4. 异步发送到Markaz服务 markazAsyncSender.sendToMarkaz(matchedTransfers); } }第四步:实现异步发送逻辑
用Spring的异步注解来实现非阻塞发送,首先在启动类加@EnableAsync,然后定义一个异步发送的服务:@Service public class MarkazAsyncSender { // 自定义线程池,避免用默认线程池 @Autowired private ThreadPoolTaskExecutor monitorTaskExecutor; @Async("monitorTaskExecutor") public void sendToMarkaz(List<Transfer> transfers) { // 这里写调用Markaz服务的逻辑,比如HTTP请求 System.out.println("异步发送" + transfers.size() + "条转账记录到Markaz"); } }同时记得配置自定义线程池的Bean,比如在配置类里:
@Configuration public class AsyncConfig { @Bean("monitorTaskExecutor") public ThreadPoolTaskExecutor monitorTaskExecutor() { ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); executor.setCorePoolSize(5); executor.setMaxPoolSize(10); executor.setQueueCapacity(20); executor.setThreadNamePrefix("Monitor-Async-"); executor.initialize(); return executor; } }
为什么这么做的优势?
- 完美适配未来扩展:新增监控规则只需要新建一个
TransferMonitorObserver的实现类,不用修改任何现有代码,完全符合开闭原则。 - 职责单一易维护:每个规则的逻辑都封装在自己的观察者类里,代码清晰,出问题也好定位,测试也方便。
- 解耦性强:主题类和观察者之间只依赖抽象接口,不管观察者怎么变,主题类的核心逻辑不用改,系统稳定性更高。
备注:内容来源于stack exchange,提问作者John Williams

