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);
    }
}

总结

模式/思想核心价值一句话概括
观察者模式解耦变更通知状态变了,所有关心它的人自动知道
事件驱动架构系统级异步解耦我只管发事件,谁消费我不关心
防御性编程提升系统健壮性永远不信任输入,永远为最坏情况做准备
短路求值优化执行效率便宜的先判断,不行就别往下走了

这四者经常组合使用:用观察者模式或事件驱动实现解耦通知,用防御性编程保证每个观察者/消费者不会因异常影响其他组件,用短路求值在触发逻辑中尽早过滤掉不需要处理的事件。

Logo

北京人形旗下天工造物具身智能开源社区,聚焦具身天工与慧思开物两大平台

更多推荐