共计 2919 个字符,预计需要花费 8 分钟才能阅读完成。
背景痛点
在电商秒杀、游戏道具获取等场景中,capture skill 服务需要处理瞬时高并发请求。传统同步处理模式(如直接操作数据库)面临三大典型问题:

- 请求超时:当并发量超过数据库连接池上限时,新请求会堆积在 TCP 队列,最终触发网关超时(504 错误)
- 重复执行:因网络抖动导致客户端重试时,缺乏幂等控制会造成资源超额分配
- 状态不一致:跨服务调用出现部分成功时,没有补偿机制会导致业务数据对账异常
某电商平台黑色星期五的实战数据显示:同步模式在 5000QPS 下,订单创建失败率高达 12%,而超时订单中有 37% 实际已扣减库存。
架构演进对比
| 方案 | TPS (万级 QPS) | 平均延迟 | 数据一致性 | 复杂度 |
|---|---|---|---|---|
| 纯数据库事务 | 0.5-1.2 | 200-500ms | 强一致 | 低 |
| 消息队列(如 RabbitMQ) | 3-5 | 50-100ms | 最终一致(秒级) | 中 |
| 事件溯源 +CQRS | 8-15 | 20-50ms | 最终一致(毫秒级) | 高 |
关键结论 :事件溯源方案通过将状态变更记录为不可变事件流,结合读写分离(CQRS) 实现:
– 写模型专注业务规则校验
– 读模型通过物化视图快速响应查询
核心实现
事件驱动基础架构
// 事件发布示例
@RestController
public class CaptureController {
@Autowired
private KafkaTemplate<String, Object> kafkaTemplate;
@PostMapping("/capture")
public ResponseEntity<String> triggerCapture(@Valid @RequestBody CaptureRequest request) {
// 生成事件 ID 保障幂等
String eventId = UUID.randomUUID().toString();
kafkaTemplate.send("capture-events",
eventId,
new CaptureEvent(eventId, request.getUserId(), request.getSkillId()));
return ResponseEntity.accepted().body(eventId);
}
}
// 消费者配置
@KafkaListener(topics = "${kafka.topic.capture}",
groupId = "${spring.kafka.consumer.group-id}")
public void handleCaptureEvent(
@Payload CaptureEvent event,
@Header(KafkaHeaders.RECEIVED_KEY) String key) {
try {skillService.processCapture(event);
} catch (JsonProcessingException e) {
// 死信队列处理
deadLetterPublish(event, e);
}
}
Saga 事务实现
@startuml
state "初始化" as init
state "扣减库存" : deduct
state "发放道具" : grant
state "完成" : done
init --> deduct : 扣减命令
deduct --> grant : 库存充足
grant --> done : 发放成功
deduct --> cancel : 库存不足
grant --> compensate : 发放失败
@enduml
@Saga
public class CaptureSaga {
@StartSaga
@SagaEventHandler(associationProperty = "captureId")
public void handle(CaptureStartedEvent event) {
// 启动库存预占
commandGateway.send(new ReserveStockCommand(...));
}
@EndSaga
@SagaEventHandler(associationProperty = "captureId")
public void handle(CaptureCompletedEvent event) {// 清理上下文}
@SagaEventHandler(associationProperty = "captureId")
public void handle(StockReservationFailedEvent event) {
// 触发补偿流程
commandGateway.send(new CancelCaptureCommand(...));
}
}
生产环境关键配置
- Kafka 消费者优化
- 设置
max.poll.interval.ms= 5 分钟避免频繁 rebalance -
使用
partition.assignment.strategy=StickyAssignor保持分区分配稳定 -
JVM 调参建议
-XX:+UseG1GC -XX:MaxGCPauseMillis=200 -XX:InitiatingHeapOccupancyPercent=35 -
分布式追踪实现
// 在事件头中传递 TraceID @Bean public ProducerFactory<String, Object> producerFactory() {Map<String, Object> config = new HashMap<>(); config.put(ProducerConfig.INTERCEPTOR_CLASSES_CONFIG, TracingProducerInterceptor.class.getName()); return new DefaultKafkaProducerFactory<>(config); }
血泪教训:三大线上事故
- 事件重复消费
- 现象:用户夜间收到重复道具
- 根因:消费者重启后未处理已提交的 offset
-
修复:实现
IdempotentEventProcessor装饰器 -
版本冲突丢失更新
- 现象:高并发时技能数量统计不准
- 根因:聚合根未实现乐观锁
-
修复:添加
@Version字段并捕获 OptimisticLockingFailureException -
死信队列堆积
- 现象:凌晨服务雪崩
- 根因:异常事件未设置 TTL
- 修复:配置
kafka.dead-letter.topic.retention.ms=86400000
本地测试环境
# docker-compose.yml
version: '3'
services:
zookeeper:
image: confluentinc/cp-zookeeper:7.0.1
ports:
- "2181:2181"
kafka:
image: confluentinc/cp-kafka:7.0.1
depends_on:
- zookeeper
ports:
- "9092:9092"
environment:
KAFKA_AUTO_CREATE_TOPICS_ENABLE: 'true'
启动命令:docker-compose up -d
通过这套架构,某 MOBA 游戏赛季更新时的 capture 请求处理能力从 800QPS 提升至 1.2 万 QPS,且 99 分位延迟稳定在 80ms 以内。关键在于:将同步阻塞操作转化为异步事件流,通过 Saga 模式保障业务正确性。
正文完
发表至: 未分类
近两天内
