SpringBoot中观察者模式、事件驱动架构、防御编程、短路求值思想示例
·
SpringBoot中观察者模式、事件驱动架构、防御编程、短路求值思想示例
一、观察者模式(Observer Pattern)
概念
观察者模式定义了对象间一对多的依赖关系。当被观察对象(Subject)状态发生变化时,所有依赖它的观察者(Observer)会自动收到通知并执行相应逻辑。
核心角色:
- Subject(被观察者):维护一组观察者,状态变更时通知所有观察者
- Observer(观察者):定义更新接口,接收 Subject 的通知
- ConcreteObserver(具体观察者):实现具体的响应逻辑
适用场景
- 当一个对象的变化需要同时通知多个其他对象时
- 发布-订阅系统(消息中间件、事件总线)
- GUI 事件监听(按钮点击、输入变化)
- 数据库变更通知(binlog 监听、CDC)
- 配置中心的配置变更推送
优缺点
| 优点 | 缺点 |
|---|---|
| 松耦合:Subject 不需要知道 Observer 的具体实现 | 通知顺序不可控:多个 Observer 的执行顺序依赖注册顺序 |
| 支持广播通信:一次状态变更通知所有观察者 | 可能引发级联更新:Observer 中又触发其他变更 |
| 符合开闭原则:新增 Observer 无需修改 Subject | 内存泄漏风险:忘记注销 Observer 导致引用不释放 |
注:
博客:
https://blog.csdn.net/badao_liumang_qizhi
代码示例
import java.util.ArrayList;
import java.util.List;
/**
* 被观察者接口.
*/
interface Subject {
void registerObserver(Observer observer);
void removeObserver(Observer observer);
void notifyObservers(String event, Object data);
}
/**
* 观察者接口.
*/
interface Observer {
void update(String event, Object data);
}
/**
* 订单状态变更 - 被观察者.
*/
class OrderStatusSubject implements Subject {
private List<Observer> observers = new ArrayList<>();
private String orderStatus;
@Override
public void registerObserver(Observer observer) {
observers.add(observer);
}
@Override
public void removeObserver(Observer observer) {
observers.remove(observer);
}
@Override
public void notifyObservers(String event, Object data) {
for (Observer observer : observers) {
observer.update(event, data);
}
}
public void changeStatus(String newStatus) {
this.orderStatus = newStatus;
notifyObservers("ORDER_STATUS_CHANGED", newStatus);
}
}
/**
* 短信通知观察者.
*/
class SmsNotifyObserver implements Observer {
@Override
public void update(String event, Object data) {
System.out.println("发送短信通知:订单状态变更为 " + data);
}
}
/**
* 库存释放观察者.
*/
class StockReleaseObserver implements Observer {
@Override
public void update(String event, Object data) {
if ("CANCELLED".equals(data)) {
System.out.println("订单取消,释放库存");
}
}
}
/**
* 积分变更观察者.
*/
class PointsObserver implements Observer {
@Override
public void update(String event, Object data) {
if ("COMPLETED".equals(data)) {
System.out.println("订单完成,发放积分");
}
}
}
// 使用示例
public class ObserverDemo {
public static void main(String[] args) {
OrderStatusSubject orderSubject = new OrderStatusSubject();
// 注册观察者
orderSubject.registerObserver(new SmsNotifyObserver());
orderSubject.registerObserver(new StockReleaseObserver());
orderSubject.registerObserver(new PointsObserver());
// 订单状态变更,自动通知所有观察者
orderSubject.changeStatus("COMPLETED");
// 输出:
// 发送短信通知:订单状态变更为 COMPLETED
// 订单完成,发放积分
}
}
Spring 中的观察者模式
import org.springframework.context.ApplicationEvent;
import org.springframework.context.ApplicationListener;
import org.springframework.context.event.EventListener;
import org.springframework.stereotype.Component;
/**
* 自定义事件.
*/
class UserRegisteredEvent extends ApplicationEvent {
private final String username;
public UserRegisteredEvent(Object source, String username) {
super(source);
this.username = username;
}
public String getUsername() {
return username;
}
}
/**
* 观察者1:发送欢迎邮件.
*/
@Component
class WelcomeEmailListener {
@EventListener
public void onUserRegistered(UserRegisteredEvent event) {
System.out.println("发送欢迎邮件给 " + event.getUsername());
}
}
/**
* 观察者2:初始化用户配置.
*/
@Component
class UserConfigInitListener {
@EventListener
public void onUserRegistered(UserRegisteredEvent event) {
System.out.println("初始化用户默认配置:" + event.getUsername());
}
}
二、事件驱动架构(Event-Driven Architecture)
概念
事件驱动架构(EDA)是一种以事件的产生、检测、消费为核心的系统架构风格。系统中的各组件通过异步事件进行通信,而不是直接调用。
核心组成:
- 事件生产者(Producer):产生事件并发布到事件通道
- 事件通道(Channel):传输事件的中间件(Kafka、RabbitMQ、EventBus)
- 事件消费者(Consumer):订阅并处理特定类型的事件
- 事件(Event):携带业务语义的不可变数据包
与观察者模式的区别
| 维度 | 观察者模式 | 事件驱动架构 |
|---|---|---|
| 范围 | 单进程内的对象通信 | 跨进程/跨服务的系统通信 |
| 耦合度 | Subject 持有 Observer 引用(弱耦合) | 生产者和消费者完全解耦 |
| 通信方式 | 同步调用(通常) | 异步消息传递 |
| 中间件 | 不需要 | 通常需要消息中间件 |
| 可靠性 | 依赖进程存活 | 消息持久化,支持重试 |
适用场景
- 微服务间异步通信
- 数据同步与变更通知
- 审计日志记录
- 异步工作流编排
- 系统解耦(订单→库存→物流→通知)
优缺点
| 优点 | 缺点 |
|---|---|
| 高度解耦:生产者无需知道消费者存在 | 最终一致性:不保证强一致,调试复杂 |
| 弹性伸缩:消费者可独立扩容 | 事件溯源困难:分布式链路追踪成本高 |
| 削峰填谷:消息队列缓冲流量尖峰 | 消息丢失/重复:需要幂等设计 |
| 易扩展:新增消费者不影响现有逻辑 | 事件风暴:高频事件可能压垮下游 |
代码示例
import java.util.*;
import java.util.concurrent.*;
import java.util.function.Consumer;
/**
* 轻量级进程内事件总线.
*/
class EventBus {
private final Map<String, List<Consumer<Object>>> subscribers = new ConcurrentHashMap<>();
private final ExecutorService executor = Executors.newFixedThreadPool(4);
/**
* 订阅事件.
*/
public void subscribe(String eventType, Consumer<Object> handler) {
subscribers.computeIfAbsent(eventType, k -> new CopyOnWriteArrayList<>()).add(handler);
}
/**
* 同步发布事件.
*/
public void publish(String eventType, Object eventData) {
List<Consumer<Object>> handlers = subscribers.get(eventType);
if (handlers != null) {
for (Consumer<Object> handler : handlers) {
handler.accept(eventData);
}
}
}
/**
* 异步发布事件.
*/
public void publishAsync(String eventType, Object eventData) {
List<Consumer<Object>> handlers = subscribers.get(eventType);
if (handlers != null) {
for (Consumer<Object> handler : handlers) {
executor.submit(() -> {
try {
handler.accept(eventData);
} catch (Exception e) {
System.err.println("事件处理异常: " + e.getMessage());
}
});
}
}
}
}
/**
* 支付完成事件.
*/
class PaymentCompletedEvent {
private final String orderId;
private final double amount;
private final long timestamp;
public PaymentCompletedEvent(String orderId, double amount) {
this.orderId = orderId;
this.amount = amount;
this.timestamp = System.currentTimeMillis();
}
public String getOrderId() { return orderId; }
public double getAmount() { return amount; }
public long getTimestamp() { return timestamp; }
}
// 使用示例
public class EventDrivenDemo {
public static void main(String[] args) {
EventBus eventBus = new EventBus();
// 注册消费者:更新订单状态
eventBus.subscribe("PAYMENT_COMPLETED", event -> {
PaymentCompletedEvent e = (PaymentCompletedEvent) event;
System.out.println("更新订单状态为已支付: " + e.getOrderId());
});
// 注册消费者:发送支付成功通知
eventBus.subscribe("PAYMENT_COMPLETED", event -> {
PaymentCompletedEvent e = (PaymentCompletedEvent) event;
System.out.println("推送支付成功通知: " + e.getOrderId() + ", 金额: " + e.getAmount());
});
// 注册消费者:触发发货流程
eventBus.subscribe("PAYMENT_COMPLETED", event -> {
PaymentCompletedEvent e = (PaymentCompletedEvent) event;
System.out.println("触发自动发货流程: " + e.getOrderId());
});
// 支付完成,发布事件
eventBus.publishAsync("PAYMENT_COMPLETED",
new PaymentCompletedEvent("ORD-20250210-001", 299.00));
}
}
Kafka 事件驱动示例(Spring Boot)
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.stereotype.Service;
/**
* 事件生产者.
*/
@Service
class OrderEventProducer {
private final KafkaTemplate<String, String> kafkaTemplate;
public OrderEventProducer(KafkaTemplate<String, String> kafkaTemplate) {
this.kafkaTemplate = kafkaTemplate;
}
/**
* 发布订单创建事件.
*/
public void publishOrderCreated(String orderId, String orderJson) {
kafkaTemplate.send("topic-order-events", orderId, orderJson)
.addCallback(
result -> System.out.println("事件发布成功: " + orderId),
ex -> System.err.println("事件发布失败: " + ex.getMessage())
);
}
}
/**
* 事件消费者 - 库存服务.
*/
@Service
class InventoryEventConsumer {
@KafkaListener(topics = "topic-order-events", groupId = "inventory-service")
public void onOrderCreated(String orderJson) {
System.out.println("库存服务收到订单事件,执行库存预占: " + orderJson);
}
}
/**
* 事件消费者 - 通知服务.
*/
@Service
class NotificationEventConsumer {
@KafkaListener(topics = "topic-order-events", groupId = "notification-service")
public void onOrderCreated(String orderJson) {
System.out.println("通知服务收到订单事件,发送确认短信: " + orderJson);
}
}
三、防御性编程(Defensive Programming)
概念
防御性编程是一种编码实践,核心思想是:永远不要信任外部输入(包括参数、返回值、配置、网络数据),在代码中主动防御各种异常情况,确保系统在异常条件下依然能安全运行或优雅降级。
核心原则:
- 不信任外部输入:所有输入都需要校验
- 不信任外部调用:第三方接口可能超时、返回异常数据
- 快速失败,优雅降级:发现问题立即处理,不让错误扩散
- 异常隔离:某个模块的异常不应影响其他模块
适用场景
- 接收外部系统数据(API 参数、MQ 消息、文件导入)
- 调用第三方接口(HTTP 调用、RPC 调用)
- 处理不可控数据源(binlog、日志解析)
- 非核心逻辑不影响主链路(推送通知、异步同步)
- 并发场景下的状态保护
防御性编程的层次
| 层次 | 手段 | 举例 |
|---|---|---|
| 输入层 | 参数校验、边界检查 | null 检查、范围校验、格式校验 |
| 逻辑层 | 异常捕获、默认值、兜底 | try-catch、Optional、默认返回 |
| 输出层 | 结果校验、超时控制 | 响应码校验、重试机制 |
| 系统层 | 熔断、限流、降级 | Sentinel、Hystrix |
代码示例
import java.util.Collections;
import java.util.List;
import java.util.Optional;
/**
* 防御性编程示例:用户服务.
*/
class UserService {
private final UserRepository userRepository;
private final ThirdPartyApiClient apiClient;
public UserService(UserRepository userRepository, ThirdPartyApiClient apiClient) {
this.userRepository = userRepository;
this.apiClient = apiClient;
}
/**
* 获取用户信息 - 防御性编程示例.
* 1. 入参校验
* 2. 数据库查询结果校验
* 3. 外部调用异常隔离
*/
public UserDto getUserInfo(Integer userId) {
// 防御1:入参空值检查
if (userId == null || userId <= 0) {
throw new IllegalArgumentException("用户ID无效: " + userId);
}
// 防御2:数据库查询结果可能为null
User user = userRepository.findById(userId).orElse(null);
if (user == null) {
throw new UserNotFoundException("用户不存在: " + userId);
}
UserDto dto = new UserDto();
dto.setUserId(user.getId());
dto.setUsername(user.getUsername());
// 防御3:外部接口调用隔离,失败不影响主流程
dto.setVipLevel(getVipLevelSafely(userId));
dto.setRecommendations(getRecommendationsSafely(userId));
return dto;
}
/**
* 安全获取VIP等级 - 外部接口调用失败时返回默认值.
*/
private Integer getVipLevelSafely(Integer userId) {
try {
Integer level = apiClient.queryVipLevel(userId);
return level != null ? level : 0;
} catch (Exception e) {
// 外部接口异常不影响主流程,返回默认值
System.err.println("查询VIP等级异常, userId: " + userId + ", error: " + e.getMessage());
return 0;
}
}
/**
* 安全获取推荐列表 - 外部接口调用失败时返回空列表.
*/
private List<String> getRecommendationsSafely(Integer userId) {
try {
List<String> list = apiClient.queryRecommendations(userId);
return list != null ? list : Collections.emptyList();
} catch (Exception e) {
System.err.println("查询推荐列表异常, userId: " + userId + ", error: " + e.getMessage());
return Collections.emptyList();
}
}
}
/**
* 防御性编程示例:批量数据处理.
* 单条数据异常不影响整批数据处理.
*/
class BatchProcessor {
/**
* 批量处理消息 - 单条异常不中断整批.
*/
public void processBatch(List<Message> messages) {
int successCount = 0;
int failCount = 0;
for (Message message : messages) {
try {
// 防御:逐条 try-catch,避免一条失败导致整批回滚
processMessage(message);
successCount++;
} catch (Exception e) {
failCount++;
System.err.println("消息处理失败, id: " + message.getId()
+ ", error: " + e.getMessage());
}
}
System.out.println("批量处理完成, 成功: " + successCount + ", 失败: " + failCount);
}
private void processMessage(Message message) {
// 防御:空值检查
if (message == null || message.getContent() == null) {
throw new IllegalArgumentException("消息内容为空");
}
// 业务处理逻辑...
}
}
/**
* 防御性编程示例:使用 Optional 避免 NPE.
*/
class OrderQueryService {
/**
* 链式调用中的空值防御.
*/
public String getOrderWarehouseName(Order order) {
return Optional.ofNullable(order)
.map(Order::getWarehouse)
.map(Warehouse::getAddress)
.map(Address::getName)
.orElse("未知仓库");
}
}
四、短路求值 / 快速失败(Short-Circuit Evaluation)
概念
短路求值是一种优化策略:在判断多个条件时,按照从「代价最小」到「代价最大」的顺序排列条件,一旦某个条件不满足就立即返回,跳过后续开销更大的判断。
核心思想:
- 尽早返回:条件不满足时立即结束,不做无用功
- 按代价排序:便宜的判断放前面,昂贵的判断放后面
- 减少资源消耗:避免不必要的 DB 查询、HTTP 调用、计算
代价参考排序
| 操作 | 典型耗时 | 排序优先级 |
|---|---|---|
| 内存变量判断(null、flag) | < 1ns | 最先 |
| 本地缓存读取 | < 1ms | 次之 |
| 数据库主键查询 | 1-10ms | 中等 |
| 数据库复杂查询 | 10-100ms | 较后 |
| HTTP/RPC 调用 | 50-500ms | 最后 |
适用场景
- 多条件判断的业务规则校验
- 权限校验链(角色→资源→操作)
- 数据同步的触发条件判断
- 搜索过滤的条件筛选
- 流水线式的数据处理(ETL)
代码示例
/**
* 短路求值示例:优惠券领取资格判断.
* 条件按开销从小到大排序.
*/
class CouponEligibilityChecker {
private final UserCache userCache;
private final OrderRepository orderRepository;
private final RiskControlClient riskControlClient;
/**
* 判断用户是否有资格领取优惠券.
* 判断顺序:内存判断 → 缓存查询 → DB查询 → 远程调用
*/
public boolean isEligible(Integer userId, String couponId) {
// 第1层:内存判断(几乎无开销)
if (userId == null || couponId == null) {
return false;
}
// 第2层:本地缓存判断(< 1ms)
// 检查优惠券是否已过期或已领完
CouponInfo couponInfo = userCache.getCouponInfo(couponId);
if (couponInfo == null || couponInfo.isExpired() || couponInfo.getStock() <= 0) {
return false;
}
// 第3层:检查用户是否已领取过(缓存 or 简单DB查询,1-5ms)
if (userCache.hasUserClaimed(userId, couponId)) {
return false;
}
// 第4层:DB查询 - 检查用户订单金额是否达标(10-50ms)
Double totalAmount = orderRepository.sumUserOrderAmount(userId);
if (totalAmount == null || totalAmount < couponInfo.getMinOrderAmount()) {
return false;
}
// 第5层:远程调用 - 风控系统校验(100-500ms)
// 只有前面所有条件都通过了,才做这个最昂贵的调用
boolean riskPass = riskControlClient.checkUser(userId);
if (!riskPass) {
return false;
}
return true;
}
}
/**
* 短路求值示例:Guard Clause(卫语句)风格.
* 通过提前返回减少嵌套层级,提升可读性.
*/
class PaymentService {
/**
* 处理退款请求.
* 使用卫语句逐步校验,不满足条件立即返回.
*/
public RefundResult processRefund(RefundRequest request) {
// 卫语句1:基础参数校验
if (request == null) {
return RefundResult.fail("请求参数为空");
}
if (request.getOrderId() == null || request.getAmount() == null) {
return RefundResult.fail("订单号或金额不能为空");
}
if (request.getAmount().compareTo(BigDecimal.ZERO) <= 0) {
return RefundResult.fail("退款金额必须大于0");
}
// 卫语句2:业务状态校验(DB查询)
Order order = orderRepository.findById(request.getOrderId());
if (order == null) {
return RefundResult.fail("订单不存在");
}
if (!"PAID".equals(order.getStatus())) {
return RefundResult.fail("订单状态不支持退款");
}
if (request.getAmount().compareTo(order.getPaidAmount()) > 0) {
return RefundResult.fail("退款金额超过实付金额");
}
// 卫语句3:幂等校验(DB查询)
if (refundRepository.existsByOrderIdAndStatus(request.getOrderId(), "SUCCESS")) {
return RefundResult.fail("该订单已退款成功,请勿重复操作");
}
// 所有校验通过,执行退款
return doRefund(order, request);
}
}
/**
* 短路求值示例:集合过滤的流水线.
* 利用Stream的惰性求值特性实现短路.
*/
class ProductFilterService {
/**
* 查找第一个满足条件的商品.
* Stream的filter是惰性的,findFirst找到后立即停止.
*/
public Optional<Product> findFirstEligibleProduct(List<Product> products) {
return products.stream()
.filter(p -> p.getStock() > 0) // 第1层:内存字段判断
.filter(p -> p.getStatus() == 1) // 第2层:内存字段判断
.filter(p -> !isInBlacklist(p.getId())) // 第3层:缓存查询
.filter(p -> checkInventory(p.getId())) // 第4层:远程调用
.findFirst(); // 找到第一个就停止
}
}
反模式对比
/**
* 反模式:没有短路,所有判断无论是否必要都执行.
*/
class BadExample {
public boolean shouldSync(Integer memberId, Integer warehouseId) {
// 问题:即使 memberId 无效,仍然会执行昂贵的 DB 查询和 HTTP 调用
boolean isValidMember = memberId != null && memberId > 0;
boolean isTargetWarehouse = warehouseRepository.isTypeThree(warehouseId); // DB查询
boolean isInCustomerList = apiClient.getCustomerList().contains(memberId); // HTTP调用
return isValidMember && isTargetWarehouse && isInCustomerList;
}
}
/**
* 正确模式:短路求值,逐步过滤.
*/
class GoodExample {
public boolean shouldSync(Integer memberId, Integer warehouseId) {
// 第1层:内存判断,0开销
if (memberId == null || memberId <= 0) {
return false;
}
// 第2层:DB主键查询,开销小
if (!warehouseRepository.isTypeThree(warehouseId)) {
return false;
}
// 第3层:HTTP调用,开销大,只有前面都通过才执行
List<Integer> customerList = apiClient.getCustomerList();
return customerList != null && customerList.contains(memberId);
}
}
总结
| 模式/思想 | 核心价值 | 一句话概括 |
|---|---|---|
| 观察者模式 | 解耦变更通知 | 状态变了,所有关心它的人自动知道 |
| 事件驱动架构 | 系统级异步解耦 | 我只管发事件,谁消费我不关心 |
| 防御性编程 | 提升系统健壮性 | 永远不信任输入,永远为最坏情况做准备 |
| 短路求值 | 优化执行效率 | 便宜的先判断,不行就别往下走了 |
这四者经常组合使用:用观察者模式或事件驱动实现解耦通知,用防御性编程保证每个观察者/消费者不会因异常影响其他组件,用短路求值在触发逻辑中尽早过滤掉不需要处理的事件。
更多推荐
所有评论(0)