学习目标
- 理解阅读、点赞、收藏三类互动业务的差异,并能用「biz + bizId」的通用模型设计可复用的互动服务
- 掌握用 Redis Hash + Lua 脚本实现高频计数类业务(点赞数、阅读数、收藏数)的方案
- 理解 Upsert 语义、软删除策略在互动业务中的应用,以及为什么缓存一致性问题「不需要解决」
- 掌握 Kafka 的核心概念(topic / partition / consumer group / acks / ISR)和 Sarama 客户端的基本用法
- 学会用「领域事件 + 批量消费」改造同步写库的高频接口,理解异步、解耦、削峰三个消息队列的核心价值
前置知识(先看这部分,避免中途卡壳):
- 已学完第01章(Go 语法)和第02章(Gin + GORM + 分层架构)。
- 知道 GORM 怎么定义模型、怎么写 CRUD,知道 Redis 是什么。
- 本地能跑 MySQL(前面用 Docker Compose 起过)。
- 不用先会 Kafka,本章从零讲。
本章你会动手做的事:
- 用
biz + bizId思路,给"评论"也设计一套可复用的点赞能力(不写新代码,只换 biz)。 - 起一个 Redis,用 Lua 脚本体验"检查 + 自增"的原子性。
- 用 Sarama 发一条
read_article消息,再起一个消费者组把它消费掉。
一、需求分析:阅读、点赞、收藏有什么不一样
生活类比:阅读、点赞、收藏,就像一家书店的三种"互动"。阅读是"每有人进门就记一笔客流",只增不减;点赞是"在书上贴便签,贴了能撕",需要记录"谁贴过"才能撕;收藏是"把书放进我的某个书架",一个人可以有好几个书架。三者形态不同,但本质都是"某个用户对某个资源的某种计数动作"——所以能抽象成同一个服务。
把这些互动抽象成一张"通用互动模型"图:
flowchart LR
U[用户] -->|对某个资源做动作| A["资源 = biz + bizId
article/123, comment/456"]
A --> I[InteractiveService
统一处理 阅读/点赞/收藏]
I --> R[(Redis 计数缓存)]
I --> DB[(MySQL 计数表)]下面马上用工程语言把这三种互动的差异拆开讲。
内容生产平台几乎都有这三个功能,看起来相似,但形态上有差别:
- 阅读数:每打开一次 +1,不可减。是最频繁的写操作。
- 点赞:可加可减(点赞 / 取消点赞),需要记录「谁点过赞」用于判重和取消。
- 收藏:可加可减,且通常带有「收藏夹」概念(一个用户可以有多个收藏夹,每篇文章可放入不同夹子)。
1.1 要不要做成「通用」功能
仔细想一下:评论可以被点赞,视频可以被点赞,动态也可以被点赞。如果每个业务自己实现一套点赞,代码重复严重。
所以应该抽象出一个通用互动服务,用 biz + bizId 标识「哪个业务的哪条记录」:
biz = "article",bizId = 123表示文章 id=123biz = "comment",bizId = 456表示评论 id=456
类似的设计在工业界很常见,也叫 resourceType + resourceId 或 requestType + requestId。
1.2 拆分还是合并
阅读 / 点赞 / 收藏是合成一个服务,还是拆三个?
| 维度 | 合并方案 | 拆分方案 |
|---|---|---|
| 性能 | 数据集中,难以分散压力 | 各自独立数据库 / Redis 集群,分散压力 |
| 研发效率 | 一套代码搞定 | 三套相似代码,重复劳动 |
| 适用场景 | 中小公司、流量不大 | 大公司、互动 QPS 极高 |
本课程采用合并方案,这也是大多数中小公司的选择。
工程视角:合并不是"偷懒",而是"先活下来"。中小公司互动 QPS 远没到需要拆分 Redis 集群的程度,一套代码维护三种相似互动反而省人力。等某天某个互动(比如点赞)量真爆了,再单独拆出去——因为用了统一的
biz + bizId模型,拆分时业务代码几乎不用改,只是挪到独立服务 + 独立存储。这就是通用模型带来的"演进弹性"。
二、阅读计数:高频计数业务的模板
2.1 设计思路
新建一个 InteractiveService,对外暴露 IncrReadCnt(biz, bizId) 方法。在文章详情接口的 Web 层聚合 ArticleService 和 InteractiveService。
为什么不在 GetPublishedById 内部直接 +1?因为「阅读」是一个独立的领域,未来还会有点赞、收藏等扩展,独立服务更清晰。
2.2 数据库设计要点
-- interactive 表:互动计数表
CREATE TABLE `interactive` (
`id` bigint UNSIGNED NOT NULL AUTO_INCREMENT,
`biz` varchar(32) NOT NULL COMMENT '业务标识,如 article',
`biz_id` bigint NOT NULL COMMENT '业务对象 ID',
`read_cnt` bigint NOT NULL DEFAULT 0,
`like_cnt` bigint NOT NULL DEFAULT 0,
`collect_cnt` bigint NOT NULL DEFAULT 0,
`ctime` bigint NOT NULL,
`utime` bigint NOT NULL,
PRIMARY KEY (`id`),
-- 关键:联合唯一索引,保证同一 biz+bizId 只有一条记录
UNIQUE KEY `biz_biz_id` (`biz`, `biz_id`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;
一张图看清"阅读计数"的数据链路,以及为什么自增必须走数据库原子操作:
flowchart TD
Req[文章详情请求] --> S[InteractiveService.IncrReadCnt]
S --> C[Redis Hash 缓存 +1
扛高频读]
S --> DB[MySQL Upsert
read_cnt = read_cnt + 1]
DB -->|原子自增| Safe[行锁保证并发安全]关键点 1:Upsert 语义。 不知道是首次插入还是已存在,用 ON DUPLICATE KEY UPDATE:
// dao/interactive.go
func (d *InteractiveDAO) InsertOrUpdate(ctx context.Context, biz string, bizId int64) error {
// Upsert:插入如果冲突则更新 read_cnt = read_cnt + 1
// 不能 "先 SELECT 再 +1 再 UPDATE",会有并发问题
// 直接利用 MySQL 的原子自增
return d.db.WithContext(ctx).Clauses(clause.OnConflict{
Columns: []clause.Column{{Name: "biz"}, {Name: "biz_id"}},
DoUpdates: clause.Assignments(map[string]interface{}{
"read_cnt": gorm.Expr("read_cnt + ?", 1),
}),
}).Create(&Interactive{
Biz: biz, BizId: bizId, ReadCnt: 1,
Ctime: time.Now().UnixMilli(), Utime: time.Now().UnixMilli(),
}).Error
}
⚠️ 新手必踩的坑:计数千万别"先查再改"。
read_cnt = read_cnt + 1必须让 MySQL 在一条语句里原子完成。如果你先SELECT读出 10,在内存里 +1 成 11,再UPDATE写回,两个 goroutine 同时读到 10,都会写 11,实际应该是 12——少记了一次。MySQL 的read_cnt = read_cnt + 1由行锁保证原子,并发安全。
关键点 2:自增必须用 read_cnt = read_cnt + 1。 千万不要先查再改,否则会有并发问题(两个 goroutine 同时读到 10,都写 11,实际应该是 12)。MySQL 的原子自增由行锁保证安全。
2.3 Redis 缓存设计
互动数据是高频访问数据,不做缓存数据库会被打爆。用 Redis Hash 结构:
key: interactive:article:123
field: read_cnt -> 1024
field: like_cnt -> 56
field: collect_cnt -> 12
为什么用 Hash 而不是 String?一篇文章的阅读/点赞/收藏在同一个 key 下,方便统一管理和批量读取。
2.4 用 Lua 脚本保证「检查 + 自增」原子性
「检查 key 是否存在 → 存在则 HIncrBy」是典型的 check-then-act 场景,并发下会出问题。Redis 单线程执行 Lua 脚本能保证原子性。
// lua 脚本:key 存在时把指定 field 自增 delta
const luaIncrIfPresent = `
if redis.call("EXISTS", KEYS[1]) == 1 then
return redis.call("HINCRBY", KEYS[1], ARGV[1], ARGV[2])
else
return 0
end
`
// IncrReadCnt 的缓存操作
func (r *InteractiveRedisCache) IncrReadCntIfPresent(ctx context.Context, biz string, bizId int64) error {
key := r.key(biz, bizId)
// ARGV[1] = field 名,ARGV[2] = 增量
_, err := r.client.Eval(ctx, luaIncrIfPresent, []string{key}, "read_cnt", 1).Result()
return err
}
⚠️ 新手必踩的坑:check-then-act 不是原子的。“先看 key 在不在、在的话再自增"这种"检查再操作"逻辑,并发下会出 race——两个请求同时看到 key 不在,都走"不操作"分支,计数就丢了。Redis 单线程执行 Lua 脚本能保证整段脚本原子执行,所以"检查 + 自增"必须塞进同一个 Lua 脚本里。
注意: HIncrBy 本身在 field 不存在时会自动置 0 再自增,所以脚本只判断 key 是否存在即可。
2.5 缓存一致性问题:为什么不解决
时序问题:缓存没有数据时,读接口会回源数据库并回写缓存,同时写接口也在更新数据库。两者交错可能出现缓存和数据库不一致。
老师的结论:不需要解决。 理由:
- 业务上阅读数少一点点不会有问题,用户感知不到。
- 只有高并发文章才会出现一致性问题,而高并发文章本身阅读量就大,少几个无所谓。
- 低并发文章几乎不会遇到这个并发场景。
数据一致性虽然重要,但并不是所有的数据一致性问题都需要解决。需要彻底解决一致性问题的场景,本来就不应该用缓存。
工程视角:判断"要不要解决一致性”,先看业务对"绝对准确"的要求有多高。阅读数、点赞数这类"展示型计数",差几个用户根本无感,花大力气保证强一致是过度设计。真正需要强一致的(如账户余额、库存),一开始就不该用"缓存 + 异步回写"结构,而是走事务或分布式事务。可观测性里这也成立:监控数据允许"最终一致",不必为它引入复杂一致性协议。
三、点赞:去重 + 计数
点赞比阅读数多了一件事:要记录「某个用户是否点过赞」,否则无法判重和取消。
3.1 数据库设计
-- user_like biz 表:记录哪个用户对哪个资源点过赞
CREATE TABLE `user_like_biz` (
`id` bigint UNSIGNED NOT NULL AUTO_INCREMENT,
`biz` varchar(32) NOT NULL,
`biz_id` bigint NOT NULL,
`user_id` bigint NOT NULL,
`status` tinyint NOT NULL DEFAULT 1 COMMENT '1 点赞 0 取消',
`ctime` bigint NOT NULL,
`utime` bigint NOT NULL,
PRIMARY KEY (`id`),
UNIQUE KEY `biz_biz_id_user_id` (`biz`, `biz_id`, `user_id`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;
3.2 硬删除 vs 软删除
考虑「点赞 → 取消 → 再点赞」的场景:
- 硬删除:取消时直接 DELETE,再点赞时 INSERT。简单,但删除+插入比 UPDATE 慢。
- 软删除:用
status字段标记 1/0,再点赞时只 UPDATE status=1。性能好,能保留数据。
⚠️ 新手必踩的坑:点赞/取消用硬删除会慢且丢数据。硬删除在"取消→再点赞"时要先
DELETE再INSERT,比UPDATE慢;而且删掉后你就失去了"这个人历史点过赞"的数据。实践上偏好软删除(用status字段标记 1/0),再点赞只UPDATE status=1,性能好还保留了行为数据。
实践建议:偏好软删除,因为性能好且保留数据,很多公司都这么干。
3.3 点赞 Service 实现
// biz/interactive.go
func (s *interactiveService) Like(ctx context.Context, biz string, id int64, uid int64) error {
// 步骤 1:给资源的总点赞数 +1
err := s.repo.IncrLikeCnt(ctx, biz, id, 1)
if err != nil {
return err
}
return s.repo.MarkLiked(ctx, biz, id, uid, true)
}
func (s *interactiveService) CancelLike(ctx context.Context, biz string, id int64, uid int64) error {
err := s.repo.IncrLikeCnt(ctx, biz, id, -1)
if err != nil {
return err
}
return s.repo.MarkLiked(ctx, biz, id, uid, false)
}
DAO 层用 Upsert,因为「点赞-取消-再点赞」时记录已存在,需要更新 status:
func (d *InteractiveDAO) UpsertLike(ctx context.Context, biz string, bizId, uid int64, status int8) error {
var st int8
if status { st = 1 } else { st = 0 }
return d.db.WithContext(ctx).Clauses(clause.OnConflict{
Columns: []clause.Column{{Name: "biz"}, {Name: "biz_id"}, {Name: "user_id"}},
DoUpdates: clause.Assignments(map[string]interface{}{
"status": st,
"utime": time.Now().UnixMilli(),
}),
}).Create(&UserLikeBiz{
Biz: biz, BizId: bizId, UserId: uid, Status: st,
Ctime: time.Now().UnixMilli(), Utime: time.Now().UnixMilli(),
}).Error
}
四、收藏:1:N 的收藏夹
收藏比点赞再复杂一点:一个用户有多个收藏夹(业务上叫「文件夹」),文章入哪个夹要记录。
4.1 表设计
-- 收藏夹本体
CREATE TABLE `collection` (
`id` bigint UNSIGNED NOT NULL AUTO_INCREMENT,
`name` varchar(128) NOT NULL,
`user_id` bigint NOT NULL,
`ctime` bigint NOT NULL,
`utime` bigint NOT NULL,
PRIMARY KEY (`id`),
KEY `idx_user_id` (`user_id`)
);
-- 收藏夹与资源的关联
CREATE TABLE `collection_item` (
`id` bigint UNSIGNED NOT NULL AUTO_INCREMENT,
`collection_id` bigint NOT NULL,
`biz` varchar(32) NOT NULL,
`biz_id` bigint NOT NULL,
`user_id` bigint NOT NULL,
`ctime` bigint NOT NULL,
`utime` bigint NOT NULL,
PRIMARY KEY (`id`),
UNIQUE KEY `biz_biz_id_user_id` (`biz`, `biz_id`, `user_id`)
);
收藏夹和收藏项是 1:N 关系。
4.2 事务保证原子性
收藏时要同时做两件事:插入收藏项 + 更新总收藏数。必须用事务保证 ACID:
func (d *InteractiveDAO) AddCollectionItem(ctx context.Context, biz string, bizId, uid, colId int64) error {
return d.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
// 1. 插入收藏项(Upsert 防重)
err := tx.Clauses(clause.OnConflict{
Columns: []clause.Column{{Name: "biz"}, {Name: "biz_id"}, {Name: "user_id"}},
DoUpdates: clause.Assignments(map[string]interface{}{
"collection_id": colId,
"utime": time.Now().UnixMilli(),
}),
}).Create(&CollectionItem{
CollectionId: colId, Biz: biz, BizId: bizId, UserId: uid,
Ctime: time.Now().UnixMilli(), Utime: time.Now().UnixMilli(),
}).Error
if err != nil {
return err
}
// 2. 更新 collect_cnt(Upsert)
return tx.Clauses(clause.OnConflict{
Columns: []clause.Column{{Name: "biz"}, {Name: "biz_id"}},
DoUpdates: clause.Assignments(map[string]interface{}{
"collect_cnt": gorm.Expr("collect_cnt + ?", 1),
}),
}).Create(&Interactive{
Biz: biz, BizId: bizId, CollectCnt: 1,
Ctime: time.Now().UnixMilli(), Utime: time.Now().UnixMilli(),
}).Error
})
}
五、查询接口:聚合 + 并发
文章详情页要展示「阅读/点赞/收藏总数 + 当前用户是否点赞/收藏」。Web 层用 errgroup 并发拉文章内容和互动信息:
func (h *ArticleHandler) PubDetail(c *gin.Context) {
var (
art domain.Article
intr domain.Interactive
)
g, ctx := errgroup.WithContext(c.Request.Context())
// 拉文章
g.Go(func() error {
var err error
art, err = h.articleSvc.PublicById(ctx, id)
return err
})
// 拉互动信息
g.Go(func() error {
var err error
intr, err = h.intrSvc.Get(ctx, "article", id, uid)
return err
})
if err := g.Wait(); err != nil { /* ... */ }
// 聚合返回
}
只缓存总数,不缓存「个人是否点赞/收藏」。 因为用户很少重复访问同一篇文章,缓存命中率低,不划算。是否缓存要基于业务理解判断。
六、Kafka 入门
6.1 消息队列三大价值
- 异步:写库慢的操作丢到队列异步处理,主流程更快返回。
- 解耦:生产者不用关心有谁消费,新增消费者对生产者透明。
- 削峰:突发流量写入队列,消费者按自己的节奏慢慢处理。
生活类比:Kafka 就像一个"快递代收点"。商家(生产者)不用盯着你在家没在家,把包裹(消息)丢进对应格口(topic 的某个分区)就走;快递员(消费者组)按自己的节奏去取,双十一爆仓时包裹先在代收点堆着(削峰),不会把商家门口挤爆。
flowchart LR
P[Producer 生产者] --> T[Topic 主题]
T --> PA1[Partition 0]
T --> PA2[Partition 1]
T --> PA3[Partition 2]
PA1 --> CG[Consumer Group 消费者组]
PA2 --> CG
PA3 --> CG
CG --> APP[业务处理]6.2 Kafka 核心概念
| 概念 | 说明 |
|---|---|
| Broker | 消息服务器进程,一台机器通常一个 broker |
| Topic | 业务主题,一个业务一个 topic |
| Partition | 分区,topic 的物理分片,是并发的基本单位 |
| Producer | 生产者,往 topic 发消息 |
| Consumer Group | 消费者组,组内消费者分摊分区 |
| ISR | In Sync Replicas,跟得上主分区的从分区集合 |
6.3 分区与消费者关系(核心面试题)
一张图说明"一分区一消费者"和消息积压的关系:
flowchart TD
T[Topic 有 3 个分区] --> P0[P0]
T --> P1[P1]
T --> P2[P2]
P0 --> C1[消费者 A]
P1 --> C2[消费者 B]
P2 --> C3[消费者 C]
D[消费者 D 分到 0 个分区] -. 多余消费者闲置 .-> P0关键规则:一个分区在同一时刻只能被同一个消费者组内的一个消费者消费。
由此衍生:
- topic 有 N 个分区 → 消费者组最多 N 个消费者有效,多出来的会闲置
- 想加快消费速率不能无限加消费者,要先加分区
这就是「消息积压」问题的根因。
6.4 消息有序性
Kafka 的有序性以分区为单位保证。一个分区内消息按写入顺序被消费。
工程视角:理解"有序性以分区为单位"是理解 Kafka 的钥匙。它意味着"全局有序"和"高并发"不可兼得——你想绝对有序就只能单分区(吞吐上不去),想高并发就得多分区(但跨分区不保序)。工程上几乎都选"业务有序"(同 key 进同分区),既保住需要的顺序,又保留并发。这正是 Kafka 在"顺序"和"扩展"之间做的务实取舍。
- 全局有序:只用一个分区(牺牲并发)
- 业务有序:用 Hash Partitioner,把同一业务 key 的消息路由到同一分区
工程视角:绝大多数业务只要"业务有序"就够了——同一篇文章的阅读事件按发生顺序累计,但不同文章之间谁先谁后无所谓。用文章
id做 key 哈希,既保证了单篇文章内的顺序,又保留了多分区的并发能力。只有极少数场景(如全局严格的事件溯源)才需要牺牲并发换全局有序。
config := sarama.NewConfig()
config.Producer.Partitioner = sarama.NewHashPartitioner // 按 key 哈希
// 发送时设置 key,比如 biz_id
msg := &sarama.ProducerMessage{
Topic: "read_article",
Key: sarama.StringEncoder(strconv.FormatInt(bizId, 10)),
Value: sarama.ByteEncoder(payload),
}
6.5 acks 参数(可靠性核心)
生产者发送消息时 acks 有三个取值:
| acks | 含义 | 性能 | 可靠性 |
|---|---|---|---|
| 0 | 不等任何确认 | 最高 | 最低,可能丢消息 |
| 1 | 主分区写入即可 | 中 | 中,主分区宕机会丢 |
| -1 / all | 所有 ISR 都确认 | 最低 | 最高,不丢消息 |
核心问「谁 ack 才算数」:
- 0:TCP 层 ack 即可
- 1:主分区 ack
- -1:所有 ISR ack
工程视角:
acks本质是在"吞吐"和"不丢消息"之间划线。金融交易、订单类消息选-1(宁可慢点也不能丢);而阅读计数、埋点日志这种"丢几条无感知"的场景,1甚至0都能接受,换来更高吞吐。选错了要么丢关键数据,要么白白牺牲性能——必须按业务定,不能无脑-1。
6.6 Docker 启动 Kafka
新版 Kafka 已不需要 ZooKeeper:
# docker-compose.yaml
services:
kafka:
image: bitnami/kafka:latest
ports:
- "9092:9092"
environment:
- KAFKA_CFG_NODE_ID=0
- KAFKA_CFG_PROCESS_ROLES=controller,broker
- KAFKA_CFG_LISTENERS=PLAINTEXT://:9092
- KAFKA_CFG_ADVERTISED_LISTENERS=PLAINTEXT://localhost:9092
- KAFKA_CFG_CONTROLLER_QUORUM_VOTERS=0@kafka:9093
- KAFKA_CFG_LISTENER_SECURITY_PROTOCOL_MAP=CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT
- KAFKA_CFG_CONTROLLER_LISTENER_NAMES=CONTROLLER
6.7 Go 客户端选型
| 库 | 优点 | 缺点 |
|---|---|---|
| IBM/sarama(原 Shopify) | 用户最多,遇到问题能搜到 | 偶有 bug |
| segmentio/kafka-go | API 简洁 | 生态小 |
| confluent-kafka-go | 性能好 | 需要 cgo,跨平台差 |
课程选 Sarama,因为出去工作遇到它的概率最大。
七、Sarama 实战
7.1 同步生产者
func NewSyncProducer(addrs []string) (sarama.SyncProducer, error) {
config := sarama.NewConfig()
// acks 选择:-1 最安全,1 平衡,0 最快
config.Producer.RequiredAcks = sarama.WaitForAll // 等价 acks=-1
config.Producer.Return.Successes = true // 同步发送必须开启
return sarama.NewSyncProducer(addrs, config)
}
func produce(p sarama.SyncProducer, topic, key, value string) error {
msg := &sarama.ProducerMessage{
Topic: topic,
Key: sarama.StringEncoder(key), // 同 key 进同分区,保证业务有序
Value: sarama.ByteEncoder(value),
}
_, _, err := p.SendMessage(msg)
return err
}
7.2 异步生产者
func NewAsyncProducer(addrs []string) (sarama.AsyncProducer, error) {
config := sarama.NewConfig()
config.Producer.Return.Successes = true // 想拿到成功通知就开启
config.Producer.Return.Errors = true // 想拿到错误通知就开启
return sarama.NewAsyncProducer(addrs, config)
}
func runAsyncProducer(p sarama.AsyncProducer) {
// 发送消息:直接丢进 Input channel
p.Input() <- &sarama.ProducerMessage{Topic: "test", Value: sarama.StringEncoder("hi")}
// 用 select 同时监听成功和失败
for {
select {
case msg := <-p.Successes():
log.Printf("发送成功 partition=%d offset=%d", msg.Partition, msg.Offset)
case err := <-p.Errors():
log.Printf("发送失败 %v", err)
}
}
}
7.3 消费者组(ConsumerGroup)
Sarama 的消费者组稍微复杂,需要实现 ConsumerGroupHandler 接口:
type readArticleHandler struct {
// 注入业务 service
}
// Setup 在消费者组 rebalance 之后、消费之前调用
func (h *readArticleHandler) Setup(sarama.ConsumerGroupSession) error { return nil }
// Cleanup 在消费结束、rebalance 之前调用
func (h *readArticleHandler) Cleanup(sarama.ConsumerGroupSession) error { return nil }
// ConsumeClaim 核心消费逻辑
func (h *readArticleHandler) ConsumeClaim(sess sarama.ConsumerGroupSession, claim sarama.ConsumerGroupClaim) error {
for msg := range claim.Messages() {
// 1. 业务处理
err := h.svc.OnReadEvent(context.Background(), msg.Value)
if err != nil {
// 失败处理:重试 / 死信队列 / 记录日志后跳过
log.Printf("消费失败 topic=%s partition=%d offset=%d err=%v",
msg.Topic, msg.Partition, msg.Offset, err)
}
// 2. 业务成功后标记消费完成,提交偏移量
sess.MarkMessage(msg, "")
}
return nil
}
工程视角:消费失败怎么处理,取决于消息能不能丢。阅读计数这种"丢了无感知"的,记日志跳过就行;订单、扣款这类"一条都不能丢"的,必须重试,重试几次还失败就进死信队列(DLQ)人工兜底,绝不能
continue默默跳过。Kafka 的"提交偏移量"节奏也跟着变:能丢的消息可以先提交再处理,不能丢的必须"处理成功再提交"——这正好呼应前面"忘记 MarkMessage 会重复消费"那条坑。
启动消费者,用 context 控制退出:
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
go func() {
for {
// Consume 阻塞消费,出错或 ctx 取消会返回
err := cg.Consume(ctx, []string{"read_article"}, &readArticleHandler{})
if err != nil { log.Println(err) }
if ctx.Err() != nil { return } // 主动退出
}
}()
八、用 Kafka 改造阅读计数
8.1 原方案的痛点
原方案在 GetPublishedById 里同步调用 IncrReadCnt 写库。读文章是高频操作,原本接口只有数据库读,现在多了一次写,DB 压力陡增。
工程视角:这个痛点的本质,是把"读路径"和"写路径"耦合了。读文章本应越快越好(用户等着看),却因为要顺手记一笔阅读数而被拖慢,还把数据库写压力引到了最高频的接口上。解法就是"读归读、写异步"——用领域事件把阅读计数甩到后台慢慢处理,读接口只管返回文章。这也是后面"事件驱动 / CQRS"思想最朴素的一次实践。
8.2 领域事件 + 异步消费
生活类比:原来你每看一篇文章,店员当场拿笔记本记一笔"阅读+1"(同步写库,拖慢你看文章)。改造后,店员只给你递文章,同时往"统计箱"里投一张小纸条(领域事件),后台有专门的统计员慢慢把纸条汇总(异步消费)。你看文章更快了,统计也没漏。
flowchart LR
Req[用户读文章] --> H[ArticleService]
H -->|当场返回文章| U[用户]
H -->|异步投纸条| K[(Kafka
read_article)]
K --> C[消费者组]
C --> S[InteractiveService.IncrReadCnt]
S --> DB[(MySQL)]引入 DDD 中的「领域事件」概念:用户读了文章,发出一个 ReadArticleEvent,对原接口透明。
// 领域事件定义
type ReadArticleEvent struct {
Uid int64
Aid int64
}
// ArticleService 内部发送消息
func (s *articleService) PublicById(ctx context.Context, id, uid int64) (domain.Article, error) {
art, err := s.repo.PublicById(ctx, id)
if err != nil { return art, err }
// 异步发送领域事件,不阻塞主流程
evt := ReadArticleEvent{Uid: uid, Aid: id}
payload, _ := json.Marshal(evt)
_ = s.producer.Produce(ctx, "read_article", strconv.FormatInt(id, 10), payload)
return art, nil
}
消费者侧订阅 read_article topic,调用 InteractiveService.IncrReadCnt:
func (h *readArticleHandler) ConsumeClaim(sess sarama.ConsumerGroupSession, claim sarama.ConsumerGroupClaim) error {
for msg := range claim.Messages() {
var evt ReadArticleEvent
_ = json.Unmarshal(msg.Value, &evt)
_ = h.intrSvc.IncrReadCnt(context.Background(), "article", evt.Aid)
// 步骤:业务成功后提交偏移量,否则重启会重复消费
sess.MarkMessage(msg, "")
}
return nil
}
> ⚠️ **新手必踩的坑:忘记 `MarkMessage` 会重复消费**。只有调用 `MarkMessage` 提交偏移量,Kafka 才认为这条消息"消费过了"。如果处理成功却没提交,消费者重启后会从旧偏移量重新拉,这条消息被再处理一次。所以"先确保业务成功,再 MarkMessage"是铁律。
后续要加新功能(比如「阅读记录」给用户看历史),只要再起一个消费者组订阅同一个 topic 即可,对原接口完全无侵入。
8.3 批量消费提升吞吐
一张图看"单条消费 vs 批量消费"的差别:
flowchart LR
subgraph 单条["单条消费"]
M1[每条消息] --> W1[一次事务写库]
end
subgraph 批量["批量消费"]
B1[攒 100 条] --> W2[一次事务批量写库]
end读文章这种场景特别适合批量处理:生产者一次发一条,消费者攒一批一起处理。
type batchReadHandler struct {
svc InteractiveService
batch int // 一批大小,比如 100
wait time.Duration // 凑批超时
}
func (h *batchReadHandler) ConsumeClaim(sess sarama.ConsumerGroupSession, claim sarama.ConsumerGroupClaim) error {
batch := make([]int64, 0, h.batch)
timer := time.NewTimer(h.wait)
defer timer.Stop()
for {
select {
case msg, ok := <-claim.Messages():
if !ok { return nil }
var evt ReadArticleEvent
_ = json.Unmarshal(msg.Value(), &evt)
batch = append(batch, evt.Aid)
sess.MarkMessage(msg, "")
if len(batch) >= h.batch {
h.flush(batch)
batch = batch[:0]
timer.Reset(h.wait)
}
case <-timer.C:
if len(batch) > 0 {
h.flush(batch)
batch = batch[:0]
}
timer.Reset(h.wait)
}
}
}
func (h *batchReadHandler) flush(aids []int64) {
// 调用批量接口 BatchIncrReadCnt,一次事务更新
_ = h.svc.BatchIncrReadCnt(context.Background(), "article", aids)
}
DAO 层用一条事务批量更新:
func (d *InteractiveDAO) BatchIncrReadCnt(ctx context.Context, biz string, ids []int64) error {
return d.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
for _, id := range ids {
err := tx.Clauses(clause.OnConflict{
Columns: []clause.Column{{Name: "biz"}, {Name: "biz_id"}},
DoUpdates: clause.Assignments(map[string]interface{}{
"read_cnt": gorm.Expr("read_cnt + ?", 1),
}),
}).Create(&Interactive{Biz: biz, BizId: id, ReadCnt: 1,
Ctime: time.Now().UnixMilli(), Utime: time.Now().UnixMilli()}).Error
if err != nil { return err }
}
return nil
})
}
插入 100 条数据一次性提交 vs 每条提交一次,性能差一个数量级。
工程视角:批量写为什么快?每一笔
INSERT都有"网络往返 + 事务提交 + 刷盘"的固定开销,100 笔就是 100 次开销;攒成 1 笔INSERT ... VALUES (...),(...)或一条事务包 100 条,固定开销只付一次。Kafka 消费侧用批量,本质也是同一个道理——把"零散小操作"聚成"一次大操作",用吞吐换延迟。这也是"批量接口是性能利器"这句经验话的底层逻辑。
九、Go channel 速查(异步 Producer 必备)
// 声明与创建
var ch chan int // 声明
ch = make(chan int) // 无缓冲
ch = make(chan int, 10) // 有缓冲,容量 10
// 发送与接收
ch <- 42 // 发送
val := <-ch // 接收
val, ok := <-ch // 接收并判断是否关闭
// 关闭
close(ch)
// 关闭后:写入 panic;读取得到零值,ok=false;重复 close panic
// for range 读取,close 后自动退出
for v := range ch {
fmt.Println(v)
}
// select 多路复用
select {
case ch <- 1:
case v := <-ch2:
fmt.Println(v)
default:
// 所有 case 阻塞时执行
}
原则:谁创建谁关闭。 channel 用不当会导致 goroutine 泄漏:发送者等不到接收者、接收者等不到发送者都会阻塞。
工程视角:channel 是 Go 并发的"管道",但它不是银弹。最典型的坑是"向已关闭的 channel 发送会 panic"“从已关闭的 channel 接收会立刻拿到零值”——所以约定"只让发送方关闭 channel",接收方永远不关。Kafka 的异步生产者内部就是 channel:你往
Input() <- msg投消息,背后 goroutine 消费它发出去,理解 channel 才能看懂 Sarama 的异步 API。
十、工程实践要点
- 计数类业务优先用 Redis Hash + Lua,不要 SELECT COUNT 查数据库。
- 自增用 SQL 原子操作(
cnt = cnt + 1),不要先读后写。 - Upsert 几乎是必需的:你永远不知道数据库里有没有这条记录。
- 缓存一致性未必都要解决:高并发数据差一点点不影响业务,低并发数据压根不会有一致性问题。
- Kafka 选 acks 要看业务:金融类用 -1,日志/计数可以用 1 甚至 0。
- 消费者记得
MarkMessage:不提交偏移量下次还会消费,导致重复处理。 - 批量接口是性能利器:消费侧攒批 + 生产侧聚合都能极大减轻 broker 压力。
- 消费者组要用 context 控制退出,避免 goroutine 泄漏。
十一、Kafka 高频面试题速记
- 为什么要用消息队列? 异步、解耦、削峰。
- ISR 是什么? 跟上主分区节奏的从分区集合,acks=-1 时全部 ISR 确认才算写入成功。
- 一个分区可以有多个消费者吗? 同一消费者组内不行;不同消费者组可以。
- 消息积压怎么办? 三招:加分区(现实基本不可能)、异步消费(goroutine 并发 + 合并提交)、批量消费。
- 怎么保证消息有序? 全局有序只用一个分区;业务有序用 Hash Partitioner 把同一 key 路由到同一分区。
- 有序消息积压怎么办? 异步消费时也要保证哈希路由,让同一 key 的消息进同一个 goroutine。
十二、自测题与动手练习
自测题(合上书能答出来,才算懂):
- 阅读、点赞、收藏为什么能抽象成同一个 Interactive 服务?
biz + bizId这个模型解决了什么问题? - 阅读计数为什么"先
SELECT读出现值再UPDATE“会少记?正确写法是什么,靠什么保证并发安全? - “检查 key 是否存在 → 存在则自增"为什么必须用 Lua 脚本?不用会出什么问题?
- Kafka 里"一个分区在同一时刻只能被同一个消费者组内的一个消费者消费"这条规则,会导致什么后果?想加快消费速率为什么不能只加消费者?
- 用 Kafka 改造阅读计数后,如果消费者忘记
MarkMessage会怎样?用什么办法给"阅读记录"这种新功能零侵入地加能力?
动手练习(建议真做一遍):
- 复用互动服务:用
biz="comment"、bizId=456的思路,给"评论"也接上点赞能力,体会不用新写一套代码的复用。 - 体验 Lua 原子性:在 Redis 里用
EVAL跑一遍luaIncrIfPresent脚本,观察 key 不存在 vs 存在时返回值的差别。 - 跑通领域事件:用 Sarama 向
read_article发一条消息,再起一个消费者组消费它,验证计数真的 +1;然后故意在消费里panic看消息会不会被重复投递。
十三、本章小结
- 阅读、点赞、收藏用
biz + bizId通用模型抽象为一个 Interactive 服务,是中小公司的高性价比方案。 - Redis Hash + Lua 脚本是高频计数业务的标配,但缓存一致性问题「不需要解决」是工程取舍的关键认知。
- Kafka 通过 topic / partition / consumer group 实现解耦和并发,理解「一分区一消费者」是搞懂消息积压和有序性的钥匙。
- 用领域事件改造阅读计数:发消息解耦、批量消费提性能、新增消费方零侵入扩展。