点赞收藏与 Kafka 入门

2023-01-13T14:11:02+08:00 | 19分钟阅读 | 更新于 2026-01-13T14:11:02+08:00

@

学习目标

  • 理解阅读、点赞、收藏三类互动业务的差异,并能用「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=123
  • biz = "comment"bizId = 456 表示评论 id=456

类似的设计在工业界很常见,也叫 resourceType + resourceIdrequestType + requestId

1.2 拆分还是合并

阅读 / 点赞 / 收藏是合成一个服务,还是拆三个?

维度合并方案拆分方案
性能数据集中,难以分散压力各自独立数据库 / Redis 集群,分散压力
研发效率一套代码搞定三套相似代码,重复劳动
适用场景中小公司、流量不大大公司、互动 QPS 极高

本课程采用合并方案,这也是大多数中小公司的选择。

工程视角:合并不是"偷懒",而是"先活下来"。中小公司互动 QPS 远没到需要拆分 Redis 集群的程度,一套代码维护三种相似互动反而省人力。等某天某个互动(比如点赞)量真爆了,再单独拆出去——因为用了统一的 biz + bizId 模型,拆分时业务代码几乎不用改,只是挪到独立服务 + 独立存储。这就是通用模型带来的"演进弹性"。


二、阅读计数:高频计数业务的模板

2.1 设计思路

新建一个 InteractiveService,对外暴露 IncrReadCnt(biz, bizId) 方法。在文章详情接口的 Web 层聚合 ArticleServiceInteractiveService

为什么不在 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 缓存一致性问题:为什么不解决

时序问题:缓存没有数据时,读接口会回源数据库并回写缓存,同时写接口也在更新数据库。两者交错可能出现缓存和数据库不一致。

老师的结论:不需要解决。 理由:

  1. 业务上阅读数少一点点不会有问题,用户感知不到。
  2. 只有高并发文章才会出现一致性问题,而高并发文章本身阅读量就大,少几个无所谓。
  3. 低并发文章几乎不会遇到这个并发场景。

数据一致性虽然重要,但并不是所有的数据一致性问题都需要解决。需要彻底解决一致性问题的场景,本来就不应该用缓存。

工程视角:判断"要不要解决一致性”,先看业务对"绝对准确"的要求有多高。阅读数、点赞数这类"展示型计数",差几个用户根本无感,花大力气保证强一致是过度设计。真正需要强一致的(如账户余额、库存),一开始就不该用"缓存 + 异步回写"结构,而是走事务或分布式事务。可观测性里这也成立:监控数据允许"最终一致",不必为它引入复杂一致性协议。


三、点赞:去重 + 计数

点赞比阅读数多了一件事:要记录「某个用户是否点过赞」,否则无法判重和取消。

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。性能好,能保留数据。

⚠️ 新手必踩的坑:点赞/取消用硬删除会慢且丢数据。硬删除在"取消→再点赞"时要先 DELETEINSERT,比 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消费者组,组内消费者分摊分区
ISRIn 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-goAPI 简洁生态小
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。


十、工程实践要点

  1. 计数类业务优先用 Redis Hash + Lua,不要 SELECT COUNT 查数据库。
  2. 自增用 SQL 原子操作cnt = cnt + 1),不要先读后写。
  3. Upsert 几乎是必需的:你永远不知道数据库里有没有这条记录。
  4. 缓存一致性未必都要解决:高并发数据差一点点不影响业务,低并发数据压根不会有一致性问题。
  5. Kafka 选 acks 要看业务:金融类用 -1,日志/计数可以用 1 甚至 0。
  6. 消费者记得 MarkMessage:不提交偏移量下次还会消费,导致重复处理。
  7. 批量接口是性能利器:消费侧攒批 + 生产侧聚合都能极大减轻 broker 压力。
  8. 消费者组要用 context 控制退出,避免 goroutine 泄漏。

十一、Kafka 高频面试题速记

  • 为什么要用消息队列? 异步、解耦、削峰。
  • ISR 是什么? 跟上主分区节奏的从分区集合,acks=-1 时全部 ISR 确认才算写入成功。
  • 一个分区可以有多个消费者吗? 同一消费者组内不行;不同消费者组可以。
  • 消息积压怎么办? 三招:加分区(现实基本不可能)、异步消费(goroutine 并发 + 合并提交)、批量消费。
  • 怎么保证消息有序? 全局有序只用一个分区;业务有序用 Hash Partitioner 把同一 key 路由到同一分区。
  • 有序消息积压怎么办? 异步消费时也要保证哈希路由,让同一 key 的消息进同一个 goroutine。

十二、自测题与动手练习

自测题(合上书能答出来,才算懂)

  1. 阅读、点赞、收藏为什么能抽象成同一个 Interactive 服务?biz + bizId 这个模型解决了什么问题?
  2. 阅读计数为什么"先 SELECT 读出现值再 UPDATE“会少记?正确写法是什么,靠什么保证并发安全?
  3. “检查 key 是否存在 → 存在则自增"为什么必须用 Lua 脚本?不用会出什么问题?
  4. Kafka 里"一个分区在同一时刻只能被同一个消费者组内的一个消费者消费"这条规则,会导致什么后果?想加快消费速率为什么不能只加消费者?
  5. 用 Kafka 改造阅读计数后,如果消费者忘记 MarkMessage 会怎样?用什么办法给"阅读记录"这种新功能零侵入地加能力?

动手练习(建议真做一遍)

  1. 复用互动服务:用 biz="comment"bizId=456 的思路,给"评论"也接上点赞能力,体会不用新写一套代码的复用。
  2. 体验 Lua 原子性:在 Redis 里用 EVAL 跑一遍 luaIncrIfPresent 脚本,观察 key 不存在 vs 存在时返回值的差别。
  3. 跑通领域事件:用 Sarama 向 read_article 发一条消息,再起一个消费者组消费它,验证计数真的 +1;然后故意在消费里 panic 看消息会不会被重复投递。

十三、本章小结

  • 阅读、点赞、收藏用 biz + bizId 通用模型抽象为一个 Interactive 服务,是中小公司的高性价比方案。
  • Redis Hash + Lua 脚本是高频计数业务的标配,但缓存一致性问题「不需要解决」是工程取舍的关键认知。
  • Kafka 通过 topic / partition / consumer group 实现解耦和并发,理解「一分区一消费者」是搞懂消息积压和有序性的钥匙。
  • 用领域事件改造阅读计数:发消息解耦、批量消费提性能、新增消费方零侵入扩展。
About Me

没什么想介绍的,一个很大众的码农…

喜欢代码,车,马,真的是 🐎

讨厌别人让我给自己的代码写注释 最厌烦别人的程序没有写注释

目标

学AI,加油!加油!