设计模式详解-中介者模式
设计模式详解:中介者模式
一、模式概述
中介者模式(Mediator Pattern)是行为型设计模式中最具系统解耦价值的模式,其核心意图在于用一个中介对象来封装一系列的对象交互。中介者使各对象不需要显式地相互引用,从而使其耦合松散,而且可以独立地改变它们之间的交互。这一模式将多对多的复杂通信网络,简化为以中介者为中心的一对多星型结构。
中介者模式的命名直接揭示了其功能——"中介"即居间协调、促成交易的第三方。在房地产领域,房产中介连接买卖双方,处理看房、议价、签约、过户等复杂流程;在航空领域,空中交通管制员协调多架飞机的起降顺序,防止碰撞冲突;在企业组织架构中,项目经理协调开发、测试、产品、运维等多团队协作。这些现实隐喻均指向同一核心特征:将原本分散的直接通信,收敛为通过中介者的间接协调。
中介者模式的深层价值在于控制复杂性的指数增长。当系统中对象数量增加时,对象间的直接通信链路数量呈平方级增长(n个对象有n(n-1)/条潜在连接)。中介者模式将这一复杂度降为线性(n个对象各连接1个中介者),使系统从"网状 chaos"转变为"星型秩序"。在微服务架构、前端组件通信、游戏引擎、工作流引擎等场景中,中介者是控制系统复杂度的关键架构模式。
二、模式结构
中介者模式包含两个核心角色,形成清晰的星型拓扑:
抽象中介者(Mediator):定义与各个同事对象通信的接口,通常包含事件注册、消息分发、状态同步等方法。
具体中介者(Concrete Mediator):实现抽象中介者接口,协调各同事对象的行为。它了解并维护所有同事对象的引用,负责实现具体的协调逻辑。
抽象同事类(Colleague):定义与中介者通信的接口,维护一个对中介者的引用。
具体同事类(Concrete Colleague):实现同事类接口,通过中介者与其他同事通信,而非直接引用。
中介者模式的关键特征在于同事对象之间的所有通信必须通过中介者。同事对象不知道彼此的存在,它们只知道中介者。这种彻底解耦使单个同事的修改不会影响其他同事,新增同事也只需向中介者注册即可。
三、深度案例:企业级微服务编排平台
以下展示一个真实场景下的中介者模式应用——金融领域的分布式事务协调与微服务编排平台,处理跨服务调用、 Saga事务、事件驱动流程等复杂交互场景。
3.1 问题域分析:微服务通信的混沌
java
// 反模式:服务间的直接网状调用
@Service
public class NaiveOrderService {
@Autowired
private InventoryServiceClient inventoryService;
@Autowired
private PaymentServiceClient paymentService;
@Autowired
private LogisticsServiceClient logisticsService;
@Autowired
private NotificationServiceClient notificationService;
@Autowired
private PointsServiceClient pointsService;
@Autowired
private RiskControlServiceClient riskControlService;
@Autowired
private AuditServiceClient auditService;
public OrderResult createOrder(CreateOrderRequest request) {
// 直接调用库存服务
InventoryReservation reservation = inventoryService.reserve(request.getItems());
// 直接调用风控服务
RiskAssessment risk = riskControlService.assess(request);
if (risk.isBlocked()) {
inventoryService.release(reservation.getId()); // 补偿库存
return OrderResult.blocked(risk.getReason());
}
// 直接调用支付服务
PaymentResult payment = paymentService.charge(request.getPayment());
if (!payment.isSuccess()) {
inventoryService.release(reservation.getId()); // 补偿库存
return OrderResult.paymentFailed(payment.getError());
}
// 直接调用物流服务
Shipment shipment = logisticsService.createShipment(request.getShipping());
// 直接调用积分服务
pointsService.earn(request.getUserId(), payment.getAmount());
// 直接调用通知服务
notificationService.sendOrderConfirmation(request.getUserId(), orderId);
// 直接调用审计服务
auditService.recordTransaction(payment);
// 问题:每新增一个服务,修改此代码;任一服务故障,影响整体;补偿逻辑散落在各处
return OrderResult.success();
}
}
上述代码中,订单服务与七个下游服务直接耦合,补偿逻辑散落在各处,新增服务需修改核心代码。中介者模式通过引入编排引擎,将服务间的直接调用转化为事件驱动的间接协调。
3.2 抽象层:中介者与同事契约
java
/**
* 抽象中介者:流程编排引擎
*/
public interface OrchestrationEngine {
/**
* 注册服务参与者
*/
void registerParticipant(ServiceParticipant participant);
/**
* 注销服务参与者
*/
void unregisterParticipant(String serviceId);
/**
* 启动流程实例
*/
ProcessInstance startProcess(String processType, ProcessContext context);
/**
* 发送事件到引擎
*/
void publishEvent(ProcessEvent event);
/**
* 查询流程状态
*/
ProcessStatus queryStatus(String processInstanceId);
/**
* 触发补偿(Saga回滚)
*/
void triggerCompensation(String processInstanceId, String failedStepId);
}
/**
* 抽象同事:服务参与者
*/
public interface ServiceParticipant {
/**
* 服务标识
*/
String getServiceId();
/**
* 感兴趣的事件类型
*/
Set<String> getInterestedEventTypes();
/**
* 处理引擎分派的事件
*/
EventResult handleEvent(ProcessEvent event, ProcessContext context);
/**
* 执行补偿操作
*/
CompensationResult compensate(CompensationCommand command, ProcessContext context);
/**
* 查询服务能力(用于路由决策)
*/
ServiceCapability getCapability();
}
/**
* 流程上下文:贯穿整个编排过程的状态容器
*/
public class ProcessContext {
private final String processInstanceId;
private final String processType;
private final Map<String, Object> variables;
private final List<ProcessStep> completedSteps;
private final List<ProcessStep> pendingCompensations;
private final ProcessStatus status;
private final Instant startTime;
// 构建方法...
public <T> T getVariable(String key) {
return (T) variables.get(key);
}
public ProcessContext setVariable(String key, Object value) {
this.variables.put(key, value);
return this;
}
public ProcessContext recordStep(ProcessStep step) {
this.completedSteps.add(step);
return this;
}
}
3.3 具体中介者:分布式Saga编排引擎
java
/**
* 具体中介者:Saga编排引擎
* 协调分布式事务的各参与服务,管理执行顺序与补偿
*/
@Component
public class SagaOrchestrationEngine implements OrchestrationEngine {
private final Map<String, ServiceParticipant> participants = new ConcurrentHashMap<>();
private final Map<String, ProcessInstance> runningInstances = new ConcurrentHashMap<>();
private final SagaDefinitionRepository definitionRepository;
private final EventBus eventBus;
private final StateMachineFactory stateMachineFactory;
private final CompensationExecutor compensationExecutor;
private final MeterRegistry meterRegistry;
@Autowired
public SagaOrchestrationEngine(SagaDefinitionRepository definitionRepository,
EventBus eventBus,
StateMachineFactory stateMachineFactory,
CompensationExecutor compensationExecutor,
MeterRegistry meterRegistry) {
this.definitionRepository = definitionRepository;
this.eventBus = eventBus;
this.stateMachineFactory = stateMachineFactory;
this.compensationExecutor = compensationExecutor;
this.meterRegistry = meterRegistry;
}
@Override
public void registerParticipant(ServiceParticipant participant) {
participants.put(participant.getServiceId(), participant);
meterRegistry.counter("orchestration.participant.registered",
"service", participant.getServiceId()).increment();
}
@Override
public void unregisterParticipant(String serviceId) {
participants.remove(serviceId);
}
@Override
public ProcessInstance startProcess(String processType, ProcessContext context) {
// 加载流程定义
SagaDefinition definition = definitionRepository.findByType(processType)
.orElseThrow(() -> new ProcessDefinitionNotFoundException(processType));
// 创建状态机
StateMachine stateMachine = stateMachineFactory.create(definition);
// 创建流程实例
ProcessInstance instance = ProcessInstance.builder()
.instanceId(generateInstanceId())
.processType(processType)
.context(context)
.stateMachine(stateMachine)
.status(ProcessStatus.RUNNING)
.startTime(Instant.now())
.build();
runningInstances.put(instance.getInstanceId(), instance);
// 启动执行
executeNextStep(instance);
meterRegistry.counter("orchestration.process.started",
"type", processType).increment();
return instance;
}
@Override
public void publishEvent(ProcessEvent event) {
// 路由到对应的流程实例
ProcessInstance instance = runningInstances.get(event.getProcessInstanceId());
if (instance == null) {
log.warn("事件对应的流程实例不存在: {}", event.getProcessInstanceId());
return;
}
// 更新状态机
StateMachine stateMachine = instance.getStateMachine();
stateMachine.fire(event.getEventType(), event.getPayload());
// 根据新状态决定下一步
if (stateMachine.getCurrentState() == ProcessState.COMPLETED) {
completeProcess(instance);
} else if (stateMachine.getCurrentState() == ProcessState.FAILED) {
triggerCompensation(instance.getInstanceId(),
stateMachine.getFailedStepId());
} else {
executeNextStep(instance);
}
}
/**
* 执行流程的下一步
*/
private void executeNextStep(ProcessInstance instance) {
SagaDefinition definition = instance.getStateMachine().getDefinition();
String currentState = instance.getStateMachine().getCurrentState().name();
// 查找当前状态可触发的步骤
List<StepDefinition> nextSteps = definition.getStepsForState(currentState);
for (StepDefinition step : nextSteps) {
ServiceParticipant participant = participants.get(step.getServiceId());
if (participant == null) {
handleMissingParticipant(instance, step);
continue;
}
// 构建事件
ProcessEvent event = ProcessEvent.builder()
.processInstanceId(instance.getInstanceId())
.eventType(step.getActionType())
.payload(buildStepPayload(instance, step))
.timestamp(Instant.now())
.build();
// 异步分派到服务参与者
dispatchAsync(participant, event, instance.getContext())
.thenAccept(result -> {
if (result.isSuccess()) {
handleStepSuccess(instance, step, result);
} else {
handleStepFailure(instance, step, result);
}
})
.exceptionally(ex -> {
handleStepException(instance, step, ex);
return null;
});
}
}
/**
* 异步分派事件到服务参与者
*/
private CompletableFuture<EventResult> dispatchAsync(ServiceParticipant participant,
ProcessEvent event,
ProcessContext context) {
return CompletableFuture.supplyAsync(() -> {
long start = System.currentTimeMillis();
try {
EventResult result = participant.handleEvent(event, context);
meterRegistry.timer("orchestration.dispatch",
"service", participant.getServiceId(),
"status", result.isSuccess() ? "success" : "failure")
.record(System.currentTimeMillis() - start, TimeUnit.MILLISECONDS);
return result;
} catch (Exception e) {
meterRegistry.counter("orchestration.dispatch.errors",
"service", participant.getServiceId()).increment();
throw new ParticipantException(participant.getServiceId(), e);
}
}, participantExecutor(participant));
}
@Override
public void triggerCompensation(String processInstanceId, String failedStepId) {
ProcessInstance instance = runningInstances.get(processInstanceId);
if (instance == null) return;
instance.setStatus(ProcessStatus.COMPENSATING);
// 获取需要补偿的步骤(逆序)
List<ProcessStep> stepsToCompensate = instance.getContext()
.getCompletedSteps().stream()
.filter(step -> !step.isCompensated())
.sorted(Comparator.comparing(ProcessStep::getSequence).reversed())
.collect(Collectors.toList());
// 顺序执行补偿
CompensationChain chain = new CompensationChain();
for (ProcessStep step : stepsToCompensate) {
ServiceParticipant participant = participants.get(step.getServiceId());
chain.addCompensation(() -> {
CompensationCommand command = CompensationCommand.builder()
.originalStepId(step.getStepId())
.originalPayload(step.getRequestPayload())
.compensationPayload(step.getCompensationState())
.reason("Upstream failure: " + failedStepId)
.build();
return participant.compensate(command, instance.getContext());
});
}
chain.execute()
.thenAccept(result -> {
if (result.isAllCompensated()) {
instance.setStatus(ProcessStatus.COMPENSATED);
publishProcessEvent(instance, ProcessEventType.COMPENSATED);
} else {
instance.setStatus(ProcessStatus.COMPENSATION_FAILED);
publishProcessEvent(instance, ProcessEventType.COMPENSATION_FAILED);
// 进入人工干预队列
escalationQueue.enqueue(instance);
}
});
}
/**
* 处理步骤成功
*/
private void handleStepSuccess(ProcessInstance instance,
StepDefinition step,
EventResult result) {
instance.getContext()
.recordStep(ProcessStep.builder()
.stepId(step.getStepId())
.serviceId(step.getServiceId())
.sequence(step.getSequence())
.requestPayload(step.getPayload())
.responsePayload(result.getPayload())
.compensationState(result.getCompensationState())
.completedAt(Instant.now())
.build());
// 发布内部事件,驱动状态机
publishEvent(ProcessEvent.builder()
.processInstanceId(instance.getInstanceId())
.eventType(step.getSuccessEventType())
.payload(result.getPayload())
.build());
}
/**
* 处理步骤失败
*/
private void handleStepFailure(ProcessInstance instance,
StepDefinition step,
EventResult result) {
instance.getContext().setVariable("failureStep", step.getStepId());
instance.getContext().setVariable("failureReason", result.getErrorMessage());
publishEvent(ProcessEvent.builder()
.processInstanceId(instance.getInstanceId())
.eventType(step.getFailureEventType())
.payload(Map.of("error", result.getErrorMessage()))
.build());
}
private void completeProcess(ProcessInstance instance) {
instance.setStatus(ProcessStatus.COMPLETED);
instance.setEndTime(Instant.now());
runningInstances.remove(instance.getInstanceId());
meterRegistry.timer("orchestration.process.duration",
"type", instance.getProcessType(),
"status", "completed")
.record(Duration.between(instance.getStartTime(), instance.getEndTime()));
}
private Executor participantExecutor(ServiceParticipant participant) {
// 根据服务能力分配线程池
return Executors.newFixedThreadPool(
participant.getCapability().getMaxConcurrency());
}
}
3.4 具体同事:服务参与者实现
java
/**
* 具体同事:库存服务参与者
*/
@Component
public class InventoryParticipant implements ServiceParticipant {
private final InventoryService inventoryService;
private final EventBus eventBus;
@Override
public String getServiceId() {
return "inventory-service";
}
@Override
public Set<String> getInterestedEventTypes() {
return Set.of("RESERVE_INVENTORY", "RELEASE_INVENTORY", "CONFIRM_DEDUCTION");
}
@Override
public EventResult handleEvent(ProcessEvent event, ProcessContext context) {
switch (event.getEventType()) {
case "RESERVE_INVENTORY":
return handleReserve(event, context);
case "RELEASE_INVENTORY":
return handleRelease(event, context);
case "CONFIRM_DEDUCTION":
return handleConfirm(event, context);
default:
return EventResult.ignored("未知事件类型");
}
}
private EventResult handleReserve(ProcessEvent event, ProcessContext context) {
ReserveRequest request = event.getPayloadAs(ReserveRequest.class);
try {
InventoryReservation reservation = inventoryService.reserve(request);
return EventResult.success()
.withPayload(Map.of(
"reservationId", reservation.getId(),
"reservedItems", reservation.getItems(),
"expireAt", reservation.getExpireAt()
))
.withCompensationState(Map.of(
"reservationId", reservation.getId(),
"action", "RELEASE"
))
.build();
} catch (InsufficientInventoryException e) {
return EventResult.failure("INSUFFICIENT_INVENTORY", e.getMessage());
}
}
private EventResult handleRelease(ProcessEvent event, ProcessContext context) {
String reservationId = event.getPayloadAsString("reservationId");
inventoryService.release(reservationId);
return EventResult.success().build();
}
private EventResult handleConfirm(ProcessEvent event, ProcessContext context) {
String reservationId = event.getPayloadAsString("reservationId");
inventoryService.confirmDeduction(reservationId);
return EventResult.success().build();
}
@Override
public CompensationResult compensate(CompensationCommand command,
ProcessContext context) {
String reservationId = command.getCompensationPayload()
.get("reservationId").toString();
try {
inventoryService.release(reservationId);
return CompensationResult.success(reservationId);
} catch (Exception e) {
return CompensationResult.failure(reservationId, e.getMessage());
}
}
@Override
public ServiceCapability getCapability() {
return ServiceCapability.builder()
.maxConcurrency(50)
.supportsCompensation(true)
.compensationTimeout(Duration.ofSeconds(30))
.idempotent(true)
.build();
}
}
/**
* 具体同事:支付服务参与者
*/
@Component
public class PaymentParticipant implements ServiceParticipant {
private final PaymentService paymentService;
private final IdempotencyKeyGenerator idempotencyGenerator;
@Override
public String getServiceId() {
return "payment-service";
}
@Override
public Set<String> getInterestedEventTypes() {
return Set.of("CHARGE_PAYMENT", "REFUND_PAYMENT");
}
@Override
public EventResult handleEvent(ProcessEvent event, ProcessContext context) {
switch (event.getEventType()) {
case "CHARGE_PAYMENT":
return handleCharge(event, context);
case "REFUND_PAYMENT":
return handleRefund(event, context);
default:
return EventResult.ignored("未知事件类型");
}
}
private EventResult handleCharge(ProcessEvent event, ProcessContext context) {
ChargeRequest request = event.getPayloadAs(ChargeRequest.class);
// 生成幂等键
String idempotencyKey = idempotencyGenerator.generate(
context.getProcessInstanceId(), "CHARGE");
try {
ChargeResult result = paymentService.charge(request.withIdempotencyKey(idempotencyKey));
return EventResult.success()
.withPayload(Map.of(
"transactionId", result.getTransactionId(),
"chargedAmount", result.getActualAmount(),
"status", result.getStatus()
))
.withCompensationState(Map.of(
"transactionId", result.getTransactionId(),
"refundAmount", result.getActualAmount()
))
.build();
} catch (PaymentException e) {
return EventResult.failure(e.getErrorCode(), e.getMessage());
}
}
@Override
public CompensationResult compensate(CompensationCommand command,
ProcessContext context) {
String transactionId = command.getCompensationPayload()
.get("transactionId").toString();
BigDecimal refundAmount = new BigDecimal(
command.getCompensationPayload().get("refundAmount").toString());
try {
RefundResult result = paymentService.refund(RefundRequest.builder()
.originalTransactionId(transactionId)
.refundAmount(refundAmount)
.reason(command.getReason())
.build());
return CompensationResult.success(transactionId);
} catch (Exception e) {
// 退款失败进入人工处理队列
return CompensationResult.requiresManualIntervention(transactionId, e.getMessage());
}
}
@Override
public ServiceCapability getCapability() {
return ServiceCapability.builder()
.maxConcurrency(100)
.supportsCompensation(true)
.compensationTimeout(Duration.ofMinutes(5))
.idempotent(true)
.build();
}
}
/**
* 具体同事:通知服务参与者
*/
@Component
public class NotificationParticipant implements ServiceParticipant {
private final NotificationDispatcher dispatcher;
@Override
public String getServiceId() {
return "notification-service";
}
@Override
public Set<String> getInterestedEventTypes() {
return Set.of("SEND_NOTIFICATION");
}
@Override
public EventResult handleEvent(ProcessEvent event, ProcessContext context) {
NotificationRequest request = event.getPayloadAs(NotificationRequest.class);
// 通知服务不参与Saga,失败不影响主流程
dispatcher.dispatchAsync(request);
return EventResult.success()
.withPayload(Map.of("notificationQueued", true))
.build();
}
@Override
public CompensationResult compensate(CompensationCommand command,
ProcessContext context) {
// 通知服务无需补偿,返回成功
return CompensationResult.success(command.getOriginalStepId());
}
@Override
public ServiceCapability getCapability() {
return ServiceCapability.builder()
.maxConcurrency(200)
.supportsCompensation(false) // 通知不补偿
.bestEffort(true)
.build();
}
}
3.5 前端中介者:组件通信协调
java
/**
* 具体中介者:前端状态管理器(Redux/MobX风格的Java实现)
*/
public class FrontendStateMediator {
private final Map<String, Component> components = new HashMap<>();
private final Map<String, List<String>> eventSubscriptions = new HashMap<>();
private final BlockingQueue<UIEvent> eventQueue = new LinkedBlockingQueue<>();
private final ExecutorService executor = Executors.newSingleThreadExecutor();
/**
* 注册UI组件
*/
public void registerComponent(Component component) {
components.put(component.getComponentId(), component);
component.setMediator(this);
// 注册组件感兴趣的事件
for (String eventType : component.getInterestedEvents()) {
eventSubscriptions.computeIfAbsent(eventType, k -> new ArrayList<>())
.add(component.getComponentId());
}
}
/**
* 发送UI事件(异步,避免阻塞)
*/
public void sendEvent(UIEvent event) {
eventQueue.offer(event);
}
/**
* 启动事件循环
*/
public void startEventLoop() {
executor.submit(() -> {
while (!Thread.interrupted()) {
try {
UIEvent event = eventQueue.take();
dispatchEvent(event);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
break;
}
}
});
}
/**
* 分派事件到感兴趣的组件
*/
private void dispatchEvent(UIEvent event) {
List<String> subscribers = eventSubscriptions.get(event.getEventType());
if (subscribers == null) return;
for (String componentId : subscribers) {
Component component = components.get(componentId);
if (component != null) {
component.handleEvent(event);
}
}
}
/**
* 协调复杂交互:表单提交触发多组件更新
*/
public void coordinateFormSubmission(String formId, FormData data) {
// 1. 验证表单
ValidationResult validation = validateForm(formId, data);
if (!validation.isValid()) {
sendEvent(UIEvent.formValidationFailed(formId, validation.getErrors()));
return;
}
// 2. 提交数据
sendEvent(UIEvent.formSubmitting(formId));
// 3. 异步提交后更新多个组件
submitAsync(formId, data)
.thenAccept(result -> {
if (result.isSuccess()) {
// 更新列表组件
sendEvent(UIEvent.listDataRefresh(formId));
// 更新统计组件
sendEvent(UIEvent.statisticsUpdate(formId));
// 关闭模态框
sendEvent(UIEvent.modalClose(formId));
// 显示成功提示
sendEvent(UIEvent.toastSuccess("提交成功"));
} else {
sendEvent(UIEvent.formSubmissionFailed(formId, result.getError()));
}
});
}
}
/**
* 抽象同事:UI组件
*/
public abstract class Component {
protected FrontendStateMediator mediator;
protected String componentId;
public void setMediator(FrontendStateMediator mediator) {
this.mediator = mediator;
}
/**
* 发送事件到中介者
*/
protected void emit(UIEvent event) {
mediator.sendEvent(event);
}
/**
* 处理来自中介者的事件
*/
public abstract void handleEvent(UIEvent event);
/**
* 感兴趣的事件类型
*/
public abstract Set<String> getInterestedEvents();
}
/**
* 具体同事:商品列表组件
*/
public class ProductListComponent extends Component {
private List<Product> products = new ArrayList<>();
private ProductFilter currentFilter;
@Override
public void handleEvent(UIEvent event) {
switch (event.getEventType()) {
case "FILTER_CHANGED":
currentFilter = event.getPayloadAs(ProductFilter.class);
refreshProducts();
break;
case "PRODUCT_CREATED":
case "PRODUCT_UPDATED":
case "PRODUCT_DELETED":
refreshProducts(); // 数据变更时自动刷新
break;
case "SORT_CHANGED":
SortOption sort = event.getPayloadAs(SortOption.class);
sortProducts(sort);
break;
}
}
@Override
public Set<String> getInterestedEvents() {
return Set.of("FILTER_CHANGED", "PRODUCT_CREATED", "PRODUCT_UPDATED",
"PRODUCT_DELETED", "SORT_CHANGED");
}
private void refreshProducts() {
// 通过中介者协调加载
mediator.sendEvent(UIEvent.dataLoading(componentId));
loadProductsAsync(currentFilter)
.thenAccept(products -> {
this.products = products;
emit(UIEvent.componentUpdated(componentId, products));
});
}
}
/**
* 具体同事:购物车组件
*/
public class ShoppingCartComponent extends Component {
private Cart cart = new Cart();
@Override
public void handleEvent(UIEvent event) {
switch (event.getEventType()) {
case "PRODUCT_ADDED_TO_CART":
Product product = event.getPayloadAs(Product.class);
cart.addItem(product);
emit(UIEvent.cartUpdated(cart));
break;
case "PRODUCT_REMOVED_FROM_CART":
String productId = event.getPayloadAsString("productId");
cart.removeItem(productId);
emit(UIEvent.cartUpdated(cart));
break;
case "CHECKOUT_INITIATED":
// 通知中介者协调结算流程
mediator.coordinateCheckout(cart);
break;
}
}
@Override
public Set<String> getInterestedEvents() {
return Set.of("PRODUCT_ADDED_TO_CART", "PRODUCT_REMOVED_FROM_CART",
"CHECKOUT_INITIATED");
}
}
四、中介者模式的高级主题
4.1 多级中介者:分层协调
java
/**
* 域中介者:处理单个业务域内的协调
*/
public class DomainMediator implements Mediator {
private List<Colleague> colleagues = new ArrayList<>();
private GlobalMediator parent;
@Override
public void notify(Colleague sender, String event) {
// 先处理域内协调
if (canHandleLocally(event)) {
handleLocally(sender, event);
} else {
// 升级到全局中介者
parent.notify(sender, event);
}
}
}
/**
* 全局中介者:跨域协调
*/
public class GlobalMediator implements Mediator {
private Map<String, DomainMediator> domains = new HashMap<>();
@Override
public void notify(Colleague sender, String event) {
// 路由到目标域
String targetDomain = resolveDomain(event);
DomainMediator mediator = domains.get(targetDomain);
if (mediator != null) {
mediator.notify(sender, event);
}
}
}
五、中介者模式与相关模式的辨析
中介者 vs 外观:外观简化子系统接口,中介者协调同事交互。外观是单向的简化封装,中介者是双向的通信协调。
中介者 vs 观察者:观察者模式中观察者与主题直接通信;中介者模式中所有通信通过中介者。中介者可视为观察者的集中化变体。
中介者 vs 责任链:责任链沿链传递请求直到处理;中介者将请求分派给特定同事。责任链是"逐个尝试",中介者是"精准路由"。
六、设计陷阱与规避策略
陷阱一:中介者成为上帝对象
所有逻辑集中在中介者,导致其过于庞大。解决方案:提取子中介者、使用状态机分解、或将部分协调逻辑下沉到同事。
陷阱二:单点故障
中介者故障导致整个系统瘫痪。解决方案:中介者集群化、状态持久化、故障转移机制。
陷阱三:性能瓶颈
所有通信经过中介者,可能成为瓶颈。解决方案:异步消息队列、分片路由、缓存优化。
七、结语
中介者模式是控制系统复杂度的核心架构工具,它将多对多的网状通信简化为星型拓扑,使系统从混沌走向秩序。在微服务编排、前端状态管理、游戏引擎、工作流系统等场景中,中介者模式是不可或缺的基础设施。理解其星型解耦原理、掌握Saga事务协调、多级分层等高级形态,警惕上帝对象、单点故障、性能瓶颈等工程陷阱,是构建可扩展分布式系统的能力核心。中介者模式的精髓在于承认直接通信的复杂性爆炸,以间接协调换取系统的可理解性与可维护性。