必须用kafka前置缓冲:秒杀请求经echo handler校验后立即发kafka并返回202,库存扣减与订单生成全异步;consumer需单goroutine+本地channel+按sku分组串行处理,并通过redis setnx实现幂等防超卖。

秒杀请求直接打到 Echo handler 会崩,必须用 Kafka 拦在中间
秒杀场景下,瞬时流量远超数据库和库存服务承载能力,Echo 的 HTTP handler 处理不过来,连接池打满、goroutine 爆炸、超时堆积。Kafka 不是“锦上添花”,而是必须前置的缓冲层——所有秒杀请求进 POST /seckill 后,只做基础校验(用户 ID、商品 ID、签名),立刻序列化为消息发往 Kafka Topic,HTTP 层立即返回 202 Accepted。真正的库存扣减、订单生成全部异步消费。
常见错误是把 Kafka 生产逻辑塞进 handler 里同步调用,还加了重试和事务包装,结果反而拖慢响应。正确做法是:handler 内只调用 producer.SendMessage(),且必须设 RequiredAcks: kafka.NoAck 或 kafka.WaitForLocal,避免网络抖动导致 HTTP 延迟飙升。
实操建议:
- 用
sync.Pool复用*sarama.ProducerMessage和 JSON 序列化 buffer,避免高频秒杀下的 GC 压力 - Topic 分区数至少等于消费者实例数,推荐设为 12 或 24,避免单分区成为吞吐瓶颈
- 不要在 handler 里解析完整订单结构体,只提取必要字段(
userID,skuID,timestamp)构造成轻量SeckillEvent结构体序列化
Consumer 必须用单 goroutine 消费 + 本地内存队列做二级缓冲
Kafka Consumer 直接调用库存服务,依然会压垮下游。哪怕用了 sarama.SyncProducer,只要并发消费,Redis Lua 扣减或 DB UPDATE 就可能超卖。真正安全的做法是:每个消费者实例启动一个单 goroutine 拉取消息,写入本地 chan SeckillEvent(容量设为 1000),再由另一个 goroutine 从该 channel 取出事件,按 skuID 做分组排队(用 map[string]*list.List + sync.RWMutex),同一商品的所有请求串行处理。
关键点在于“分组串行”不是靠 Kafka 分区保证的——分区只保证顺序,不保证同 sku 落同一分区。必须在 consumer 内存中做二次路由。
实操建议:
- 本地 channel 容量不能过大,否则 OOM;也不能过小,否则
sarama.Consumer的Messages()channel 积压导致 Kafka rebalance - 每个 sku 队列用独立的
time.Timer控制超时(比如 5 秒未处理完就丢弃),防止某个热 sku 卡死全局 - 扣库存前先查 Redis 缓存的
stock:{skuID},为 0 直接 skip,避免穿透 DB
如何避免 Kafka 消息重复导致超卖
Kafka 本身提供 at-least-once 语义,Consumer crash 或 rebalance 时可能重复拉取同一批 offset。如果只是简单地“扣库存 → 写订单”,重复消费必然超卖。
必须引入幂等性控制:在消息体中带上唯一 requestID(由 Echo handler 生成,如 uuid.New().String()),Consumer 扣库存前先用 SETNX redis:seckill:idempotent:{requestID} 1 EX 3600 占坑。失败则跳过;成功才继续后续流程,并在订单落库后异步删除该 key(或依赖 TTL 自动过期)。
注意:不能用 MySQL 唯一索引代替,因为订单表写入在扣库存之后,中间存在时间窗口。
实操建议:
-
requestID必须由 Echo handler 生成并透传,不能由 Consumer 自己生成,否则无法关联前端请求 - Redis key 的 TTL 要大于整个秒杀生命周期(含补偿任务),建议设为 2 小时
- 幂等检查必须在任何状态变更操作之前,且需捕获
redis.Nil错误作为“已存在”信号
监控和降级开关必须嵌入 Kafka 生产/消费链路
没有监控的削峰等于裸奔。要实时知道:每秒进 Kafka 的请求数(seckill_produce_rate)、Consumer 消费延迟(kafka_lag{topic="seckill"} )、本地队列积压数(local_queue_length{sku="1001"})、幂等 key 冲突率(idempotent_conflict_rate)。
一旦 local_queue_length 持续 > 500 或 kafka_lag > 10000,就要触发降级:Echo handler 改为直接返回 429 Too Many Requests,不再发消息到 Kafka。
实操建议:
- 用
prometheus.Counter和Gauge在 producer/consumer 关键路径埋点,避免用 log 打点(性能损耗大) - 降级开关存在 Redis 中,key 为
feature:seckill:enabled,Echo handler 每次请求都GET一次,缓存 1 秒(用time.Now().Unix()%10 == 0控制刷新频率) - Consumer 启动时检查 Kafka topic 是否存在,不存在则 panic,避免静默失败
最易被忽略的是本地内存队列的 GC 友好性——别用 map[requestID]struct{} 存已处理 ID,而要用 LRU cache 或带 TTL 的 map,否则内存只增不减。还有就是 Kafka 的 MaxRetry 别设成 -1,重试无限循环会卡死 goroutine。
golang免费学习笔记(深入):立即使用
在学习笔记中,你将探索golang的核心概念和高级技巧!











