跳到正文
Joeplover
后端开发·2025-10-28·约 2 分钟阅读

电商实时分析平台: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(毫秒级延迟)
  • 加入实时大屏展示