共计 1173 个字符,预计需要花费 3 分钟才能阅读完成。
背景痛点
数据挖掘竞赛平台的实时评分系统面临两大核心挑战:
- 突发流量 :竞赛提交高峰期的 QPS 可达日常的 50 倍以上,传统数据库直连方案易引发连接池耗尽
- 强一致性要求 :排名变动需在 200ms 内全局可见,避免出现选手端显示分数与排行榜不一致
典型场景中,单次评分涉及特征提取、规则计算、结果持久化等步骤,传统同步处理模式在流量峰值时会造成请求堆积。
架构设计
方案对比
- 纯队列方案 (Kafka+Spark Streaming)
- 优点:吞吐量高
-
缺点:端到端延迟通常 >500ms,状态管理复杂
-
混合方案 (Redis Stream+ 微服务批处理)
- 流处理层:Redis Stream 处理实时写入(8000 QPS/ 节点)
- 批处理层:Spring Batch 聚合 10ms 窗口内的请求
- 存储层:MySQL 分片存储最终结果,Redis 缓存排行榜

核心代码实现
消息分区策略
// 按用户 ID 哈希分片,保证同一用户请求有序
public String determineStreamKey(Long userId) {int shard = userId.hashCode() % STREAM_SHARDS;
return "score_stream_" + shard;
}
批量消费实现
@Scheduled(fixedDelay = 10)
public void batchProcess() {RLock lock = redisson.getLock("score_lock");
try {if (lock.tryLock(5, TimeUnit.MILLISECONDS)) {
List<MapRecord<String, String, String>> records =
redisTemplate.opsForStream().read(Consumer.from("group1", "instance1"),
StreamReadOptions.empty().count(100),
StreamOffset.create(streamKey, ReadOffset.lastConsumed())
);
// 批量处理逻辑
}
} finally {lock.unlock();
}
}
性能测试
使用 JMeter 模拟 10000 并发用户持续 5 分钟:
| 指标 | 数值 |
|---|---|
| 平均延迟 | 128ms |
| P99 延迟 | 193ms |
| 错误率 | 0.002% |
避坑指南
- Redis 内存控制 :
- 设置 stream 的 maxlen 参数(建议 10000)
-
监控 used_memory_peak 指标
-
实例数计算公式 :
实例数 = 峰值 QPS / (单实例处理能力 * 0.7) -
幂等性保障 :
- 消息携带唯一 ID
- 处理前查询去重表
延伸思考
动态评分规则可通过以下方式实现:
- 规则引擎层:将评分公式存储在 Drools 规则库
- 版本管理:每个消息携带规则版本号
- 热更新:通过 Spring Cloud Config 推送新规则
实际部署时建议采用蓝绿发布验证规则变更,避免大规模回溯历史数据。
正文完
