电商实时分析平台:Spring Boot + Kafka + MySQL
电商用户行为实时分析系统,Kafka 消费用户行为数据,Spring Boot 提供 REST API 查询实时热销商品。
电商实时分析平台:Spring Boot + Kafka + MySQL
业务需求
电商平台需要实时(秒级)了解运营数据:
- 今日实时订单数
- 实时销售额
- 热门商品排名
- 分区销量对比
传统方案是后台跑 SQL 查询,但数据量大时 MySQL 扛不住高并发查询。本方案的思路:用 Kafka 做缓冲削峰,异步写入分析结果。
架构
订单系统 → Kafka(缓冲) → 消费者(实时聚合) → MySQL/ClickHouse(结果存储) → 前端查询
Spring Boot 收集数据
@RestController
@RequestMapping("/api/events")
public class EventCollector {
@Autowired
private KafkaTemplate<String, String> kafkaTemplate;
@PostMapping("/order")
public Result collectOrder(@RequestBody @Valid OrderEvent event) {
// 不做业务处理,直接发到 Kafka
kafkaTemplate.send("order-events", event.getOrderId(), json.toJson(event));
return Result.success();
}
}
Kafka 缓冲削峰
Kafka 在这里承担关键角色:流量暴增时数据积在 Kafka,不会压垮后端 MySQL。
假设双十一促销时订单量暴增 10 倍:
- 没有 Kafka:MySQL 直接被打满,订单丢失
- 有 Kafka:Kafka 积压消息,消费者慢慢处理,MySQL 压力稳定
实时聚合消费者
@Component
public class OrderAggregator {
@Autowired
private JdbcTemplate jdbcTemplate;
@KafkaListener(topics = "order-events", groupId = "realtime-analysis")
public void consume(String message) {
OrderEvent event = json.parse(message, OrderEvent.class);
// 更新实时统计数据
jdbcTemplate.update(
"INSERT INTO realtime_stats (stat_date, stat_hour, order_count, total_amount) " +
"VALUES (CURRENT_DATE, HOUR(NOW()), 1, ?) " +
"ON DUPLICATE KEY UPDATE order_count = order_count + 1, total_amount = total_amount + ?",
event.getAmount(), event.getAmount()
);
}
}
从 MySQL 到 ClickHouse
当数据量达到千万级别时,MySQL 的聚合查询会变慢:
-- 百万级数据时,这条查询可能超过 10 秒
SELECT product_id, SUM(amount) FROM orders
WHERE created_at > NOW() - INTERVAL 7 DAY
GROUP BY product_id ORDER BY SUM(amount) DESC LIMIT 20;
ClickHouse 是列式存储数据库,对此类分析查询快 10-100 倍:
-- ClickHouse:MB 级扫描速度
SELECT product_id, sum(amount) FROM orders
WHERE created_at > now() - INTERVAL 7 DAY
GROUP BY product_id ORDER BY sum(amount) DESC LIMIT 20;
实时告警
在消费者中增加阈值检测:
if (event.getAmount() > 10000) {
// 大额订单通知
alertService.sendAlert("大额订单: " + event.getAmount());
}
if (errorRate > 0.05) {
// 错误率超过 5% 告警
alertService.sendAlert("错误率异常: " + errorRate);
}
总结
该方案的核心价值:Kafka 缓冲层将突发流量与处理能力解耦。后续升级可以:
- 替换 MySQL 为 ClickHouse(查询性能)
- 替换 Spring Boot 消费者为 Flink(毫秒级延迟)
- 加入实时大屏展示