Kratos 云盘项目:Kafka 事件驱动与可靠消费

2025-01-15T10:30:00+08:00 | 29分钟阅读 | 更新于 2025-01-15T10:30:00+08:00

@

本篇基于真实项目 kratos-cloud-diskinternal/eventinternal/biz/event.gointernal/data/kafka.go 源码整理,所有常量、函数名、重试次数、退避时间、TTL、DLQ 命名均来自代码,可直接对照阅读。

学习目标

阅读完本文,你应该能够:

  1. 说清为什么用事件驱动:用云盘"删文件 / 移入回收站后异步取消相关分享"的真实场景,解释解耦、削峰、最终一致三件事。
  2. 画出分层与依赖方向:biz 只依赖 EventPublisher 接口,event 包用 eventPublisherAdapter 把事件类型映射到 topic,避免循环依赖。
  3. 讲透可靠消费三板斧:消费端幂等(Redis SetNX / CheckAndSet,24h TTL)、失败重试(3 次、退避 200*(i+1)ms)、死信队列(topic+".dlq")。
  4. 解释分区与有序PartitionKey="user:%d" 按用户维度保证同用户事件同分区有序,不同用户可并行。
  5. 回答降级与可靠性边界:Kafka 未配置时返回 noop 实现;本项目是 at-least-once 语义,必须靠消费端幂等兜底。

前置知识

  • 了解 Kafka 的基本概念:Broker、Topic、Partition、Consumer Group、offset。
  • 了解 Go 接口与依赖注入(go-wire)的基本用法。
  • 了解 Redis 的 SETNX 原子语义。

动手三件事

  1. 找到 internal/biz/event.go,把 7 种事件常量与对应 Payload 抄一遍,想清楚每类事件"谁生产、谁消费"。
  2. 在本地起一个 Docker Kafka,把 conf.Data_Kafka.Addrs 配好,打一条 file.deleted 事件,观察 EventHandlerService 是否取消了分享。
  3. conf.Data_Kafka.Addrs 清空,启动服务,确认业务流程不因为"没有 Kafka"而崩溃(验证 noop 降级)。

一、为什么是事件驱动:从"删文件"说起

Q1. 为什么云盘项目要引入 Kafka 事件驱动?直接同步调用不行吗?

答: 类比一下。假设你在朋友圈发一张照片,系统要同步去做三件事:给所有好友发通知、更新你的相册统计、给图片打水印。如果这三件事都在你点"发布"的请求里同步做完,那么:

  • 任何一件慢了(比如通知服务卡住),你的发布请求就被拖死;
  • 通知服务挂了,你的发布就失败,哪怕照片本身存好了;
  • 以后想再加"发积分"“同步搜索引擎”,每次都要改发布接口,耦合爆炸。

云盘项目里最典型的真实场景就是:用户删除了一个文件(或把文件移入回收站),系统需要取消掉所有指向这个文件的分享链接

为什么这件事适合异步?因为"取消分享"不是用户删除文件的"主干"。用户的核心诉求是"文件删掉了",至于分享链接在 100ms 内还是 2s 内被取消,用户感知不强;但如果在删除请求里直接去查所有分享、逐个删除,一旦分享表很大、或数据库慢,删除接口就变慢甚至超时。更糟的是,如果分享服务暂时不可用,用户连文件都删不掉——这显然不合理。

所以项目把这类"主流程之外、但最终必须发生"的动作,抽象成事件(Event),丢给 Kafka,由消费者异步处理。源码里一共定义了 7 种事件常量,见 internal/biz/event.go

// internal/biz/event.go
const (
	EventFileDeleted = "file.deleted"
	EventFileMoved   = "file.moved"
	EventFileTrashed = "file.trashed"

	EventShareCreated   = "share.created"
	EventShareCancelled = "share.cancelled"
	EventShareAccessed  = "share.accessed"

	EventRecycleClean = "recycle.clean"

	EventUploadCompleted = "upload.completed"
)

这 7 种事件覆盖了文件生命周期(删除 / 移动 / 移入回收站 / 上传完成)、分享生命周期(创建 / 取消 / 访问)和回收站清理。它们的共同特点是:生产方只管"发生了什么",不关心"谁来处理、怎么处理"

⚠️ 坑:事件驱动不等于"随便异步"。只有满足"失败可重试、允许短暂延迟、不需要即时返回结果给调用方"的场景才适合事件化。像"扣款"这种需要强一致、立即返回成功/失败的动作,不适合用本项目的 at-least-once 语义直接异步。

引入事件驱动带来三个工程收益,也是面试高频点:

  • 解耦:文件用例层(FileUC)不再直接依赖分享仓储。它只调用 EventPublisher.Publish,由 event 包去取消分享。删文件的主流程完全不感知下游。
  • 削峰:高峰期大量删除/上传事件先堆在 Kafka,消费者按自己节奏处理,不会因为瞬时流量把分享库打挂。
  • 最终一致:分享最终会被取消,但不要求与删除"同一时刻"完成。只要消费者最终处理成功,系统达到一致状态。

下面这张图展示整体的事件驱动架构,注意依赖方向是单向的(业务层 → 事件层 → MQ)。

graph LR
    A[文件用例层
FileUC] -->|发布 文件事件| B[事件发布器
EventPublisher] C[分享用例层
ShareUC] -->|发布 分享事件| B D[上传用例层
UploadUC] -->|发布 上传事件| B B --> E[Kafka 集群
多 Topic 分区] E --> F[消费者组
Consumer Group] F --> G[事件处理器服务
EventHandlerService] G --> H[(MySQL
取消分享)] G --> I[(Redis
幂等去重)]

二、分层与依赖:biz 为什么只认接口

Q2. 业务层(biz)怎么做到不依赖具体 MQ 实现?EventPublisher 接口与 adapter 怎么分工?

答: 类比:你家公司要发快递,业务部门(biz)只关心"把包裹交给前台(接口)",至于前台用顺丰还是邮政(Kafka 还是 RabbitMQ),业务部门不关心。这样哪天换物流公司,业务部门一行代码都不用改。

项目里 biz 层定义了一个极简的接口 EventPublisher,位于 internal/biz/event.go

// internal/biz/event.go
// EventPublisher 定义业务层发布异步事件的接口,由 event 包实现。
// 通过接口注入避免 biz 与 event 包之间的循环依赖。
type EventPublisher interface {
	Publish(ctx context.Context, eventType string, payload interface{}) error
}

注意接口只有 Publish 一个方法,入参是"事件类型字符串 + 任意 payload"。biz 完全不知道 Kafka、topic、分区键这些概念。

真正把事件发到 Kafka 的是 event 包里的 eventPublisherAdapter(见 internal/event/biz_publisher.go)。它做了两件关键的事:

  1. 事件类型 → topic 的映射:业务层只说"file.deleted",adapter 负责把它路由到具体 Kafka topic(比如 cloud_disk_file_events)。
  2. 设置分区键 PartitionKey:按用户维度,保证同一用户的事件同分区有序。
// internal/event/biz_publisher.go
func (a *eventPublisherAdapter) Publish(ctx context.Context, eventType string, payload interface{}) error {
	topic, ok := a.topics[eventType]
	if !ok {
		return nil
	}
	evt := NewEvent(eventType, SourceSystem, payload)
	// 以用户维度作为分区键,保证同一用户的事件落在同一分区(分区内有序)
	switch p := payload.(type) {
	case *biz.FileChangedPayload:
		evt.PartitionKey = fmt.Sprintf("user:%d", p.UserID)
	case *biz.SharePayload:
		evt.PartitionKey = fmt.Sprintf("user:%d", p.UserID)
	case *biz.UploadCompletedPayload:
		evt.PartitionKey = fmt.Sprintf("user:%d", p.UserID)
	}
	return a.producer.Publish(ctx, topic, evt)
}

topic 映射在构造时建立(NewEventPublisherAdapter):

// internal/event/biz_publisher.go
topics := map[string]string{
	biz.EventFileDeleted:     prefix + cfg.Topics.FileEvents,
	biz.EventFileMoved:       prefix + cfg.Topics.FileEvents,
	biz.EventFileTrashed:     prefix + cfg.Topics.FileEvents,
	biz.EventShareCreated:    prefix + cfg.Topics.ShareEvents,
	biz.EventShareCancelled:  prefix + cfg.Topics.ShareEvents,
	biz.EventRecycleClean:    prefix + cfg.Topics.RecycleEvents,
	biz.EventUploadCompleted: prefix + cfg.Topics.UploadEvents,
}

这里有个设计要点:topic 前缀 cfg.TopicPrefix 默认是 cloud_disk_,再拼上 cfg.Topics.FileEvents 等子配置。业务层根本不知道 topic 叫什么,改 topic 名字只动配置,不动代码。

⚠️ 坑:为什么要把映射放在 event 包而不是 biz 包?因为 biz → event 会形成循环依赖(event 的 handler 又依赖 biz 的 ShareRepo)。把"按事件类型选 topic"这层放在 event 包,biz 只暴露接口,依赖方向是单向的 biz ← event,编译通过。

另外注意 biz 的生产调用通常是不阻塞主流程的。在用例层里你会看到类似 _ = publisher.Publish(ctx, ...) 的写法:发布失败只记日志,不影响"删文件"返回成功。这是"可选依赖"思想的体现——消息总线是增强能力,不是主干必选项。


三、生产者:用 segmentio/kafka-go 发消息

Q3. 生产者怎么用 segmentio/kafka-go 发送消息?kafka.Writer 配置与 Publish 流程是什么?

答: 类比:Kafka 的 kafka.Writer 就像"快递公司的揽收台"。你把写好的包裹(kafka.Message)交给它,它负责按地址(Broker)发出去。segmentio/kafka-go 是 Go 里比较主流的 Kafka 客户端,相比 sarama 它 API 更简洁,且原生支持 context 与消费组。

先看底层 Writer 的配置,在 internal/data/kafka.go

// internal/data/kafka.go
writer := &kafka.Writer{
	Addr:                   kafka.TCP(addrs...),
	Balancer:               &kafka.RoundRobin{},
	RequiredAcks:           kafka.RequireOne,
	AllowAutoTopicCreation: true,
	BatchTimeout:           10 * time.Millisecond,
	BatchSize:              1,
}

几个关键配置,面试常问:

  • Addr:Broker 地址列表,从 c.Addrs 来。
  • Balancer:分区策略。RoundRobin 轮询分配到各分区。注意:本项目生产消息时显式带了 Key(分区键),此时 Balancer 的轮询不生效,Kafka 会用 hash(Key) % 分区数 决定分区,从而保证同 Key 同分区。
  • RequiredAcks:这里是 kafka.RequireOne(leader 确认即可)。如果要更强的不丢消息保证,面试里要答"Leader 和所有 ISR 都确认 = RequireAll,对应 acks=all",本项目为了吞吐没开到最高。
  • AllowAutoTopicCreation:topic 不存在自动创建,开发期方便,生产环境一般关掉、改为运维预建。
  • BatchSize: 1 + BatchTimeout: 10ms:每条都尽快发,牺牲批量吞吐换低延迟,适合事件这类"小且不频繁"的消息。

⚠️ 坑:RequiredAcks: kafka.RequireOne 是"至少一次"语义的来源之一。leader 写完就可能 ack,如果 leader 写完还没同步给 follower 就宕机,消息可能丢。要保证不丢,需要 RequireAll + 生产者重试 + 多副本。本项目靠消费端幂等兜"重复",但"丢失"只能靠 acks 配置缓解。

生产者接口 Producer 也定义在 event 包,方便替换实现:

// internal/event/producer.go
type Producer interface {
	Publish(ctx context.Context, topic string, event *Event) error
	Close() error
}

kafkaProducer.Publish 的核心流程(真实代码):

// internal/event/producer.go(核心节选)
event.ID = generateEventID()
if event.IdempotencyKey == "" {
	event.IdempotencyKey = fmt.Sprintf("%s:%s", event.Type, event.ID)
}
if event.Timestamp == "" {
	event.Timestamp = time.Now().UTC().Format(time.RFC3339)
}

data, err := json.Marshal(event)
// ...

partitionKey := event.PartitionKey
if partitionKey == "" {
	partitionKey = event.IdempotencyKey // 缺省回退到幂等键
}

msg := kafka.Message{
	Topic: topic,
	Key:   []byte(partitionKey),
	Value: data,
	Headers: []kafka.Header{
		{Key: "event_type", Value: []byte(event.Type)},
		{Key: "event_source", Value: []byte(event.Source)},
	},
}

var lastErr error
for i := 0; i < 3; i++ {
	select {
	case <-ctx.Done():
		return ctx.Err()
	default:
	}
	if err := p.writer.WriteMessages(ctx, msg); err != nil {
		lastErr = err
		time.Sleep(time.Duration(100*(i+1)) * time.Millisecond) // 生产者侧退避
		continue
	}
	return nil
}
return fmt.Errorf("event: failed to publish after 3 retries: %w", lastErr)

注意两点:

  • 分区键优先用 PartitionKey,缺失时回退到 IdempotencyKey:这保证了"即使没设用户维度,也不会乱序到不可控"。
  • 消息里带了 Headerevent_typeevent_source):消费者即使不反序列化也能从 Header 拿到类型,方便做路由和监控。

统一事件结构 Eventinternal/event/event.go)是所有事件的信封:

// internal/event/event.go
type Event struct {
	ID             string      `json:"id"`
	IdempotencyKey string      `json:"idempotency_key"`
	PartitionKey   string      `json:"partition_key,omitempty"`
	Type           string      `json:"type"`
	Source         string      `json:"source"`
	Timestamp      string      `json:"timestamp"`
	Payload        interface{} `json:"payload"`
}

NewEvent 会自动填 TypeSourceTimestampPublish 时再补 IDIdempotencyKey

下面这张时序图把"生产 → 消费"全链路串起来,注意每一步的参与者。

sequenceDiagram
    participant UC as 业务用例层
    participant AP as eventPublisherAdapter
    participant WP as kafkaProducer
    participant KA as Kafka Broker
    participant CS as kafkaConsumer
    participant HS as HandlerService
    participant RD as Redis
    UC->>AP: Publish(eventType, payload)
    AP->>AP: 映射 topic 加 设置 PartitionKey
    AP->>WP: Publish(topic, Event)
    WP->>KA: WriteMessages(Key 等于 用户维度)
    KA-->>WP: ack 确认
    CS->>KA: FetchMessage 拉取
    KA-->>CS: 返回消息
    CS->>RD: CheckAndSet idempotency 冒号 key
    RD-->>CS: true 首次 或 false 重复
    CS->>HS: handler 处理 event
    HS-->>CS: 成功 或 失败
    CS->>KA: CommitMessages 提交位移

四、消费者:拉消息、分发、分发到处理器

Q4. 消费者是怎么拉消息、分发到处理器的?消费者组与消费循环是怎样的?

答: 类比:消费者就像"快递柜取件员"。他不断去柜子里(Kafka 分区)取件,看包裹上写的是"文件类"还是"分享类",交给对应的处理窗口(handler)。

消费者构造在 internal/event/consumer.goNewConsumer,使用的是 kafka.Reader(kafka-go 的消费者实现,原生支持消费组):

// internal/event/consumer.go
reader := kafka.NewReader(kafka.ReaderConfig{
	Brokers:        cfg.Addrs,
	GroupID:        groupID,
	GroupTopics:    topics,
	MinBytes:       1,
	MaxBytes:       10e6,
	MaxWait:        1 * time.Second,
	CommitInterval: 1 * time.Second,
	StartOffset:    kafka.LastOffset,
})

逐条解释面试点:

  • GroupID + GroupTopics:组成消费者组。同一个 Group 内,一个分区只会被一个消费者实例消费;多个实例会自动做分区再均衡(rebalance),实现水平扩展与高可用。
  • CommitInterval: 1 * time.Second:自动提交 offset 的周期。注意这是"周期自动提交",配合代码里显式的 CommitMessages 使用。
  • StartOffset: kafka.LastOffset:消费者初次加入组时从最新位置开始,不处理历史积压。若要重放历史,需改成 FirstOffset 或指定具体 offset。
  • MaxWait: 1s:长轮询等待,没有消息时空等 1s 再返回,避免空转烧 CPU。

消费主循环 consumeLoop

// internal/event/consumer.go
func (c *kafkaConsumer) consumeLoop(ctx context.Context) {
	defer c.wg.Done()
	defer c.logger.Info("event: consumer stopped")
	for {
		select {
		case <-ctx.Done():
			return
		default:
		}
		msg, err := c.reader.FetchMessage(ctx)
		if err != nil {
			if ctx.Err() != nil {
				return
			}
			time.Sleep(1 * time.Second)
			continue
		}
		c.processMessage(ctx, msg)
	}
}

这里有个细节考FetchMessage 是阻塞拉取,但循环开头用 select 检查 ctx.Done(),保证收到取消信号能及时退出。每次拉到一条就交给 processMessage 串行处理(单分区内天然有序),处理完再拉下一条——这保证了单分区内严格按序处理

处理器注册用 RegisterHandler 把"事件类型 → 函数"存进 map:

// internal/event/consumer.go
func (c *kafkaConsumer) RegisterHandler(eventType string, handler EventHandler) {
	c.mu.Lock()
	defer c.mu.Unlock()
	c.handlers[eventType] = handler
	c.logger.Info("event: handler registered", "type", eventType)
}

RegisterAllHandlershandlers.go)把 7 种事件都注册到同一个 EventHandlerService.Handle 上,由 Handle 内部再按类型分发到 file/share/recycle/upload 四类子处理器:

// internal/event/handlers.go
func RegisterAllHandlers(consumer Consumer, handlerSvc *EventHandlerService) {
	consumer.RegisterHandler(biz.EventFileDeleted, handlerSvc.Handle)
	consumer.RegisterHandler(biz.EventFileMoved, handlerSvc.Handle)
	consumer.RegisterHandler(biz.EventFileTrashed, handlerSvc.Handle)
	consumer.RegisterHandler(biz.EventShareCreated, handlerSvc.Handle)
	consumer.RegisterHandler(biz.EventShareCancelled, handlerSvc.Handle)
	consumer.RegisterHandler(biz.EventRecycleClean, handlerSvc.Handle)
	consumer.RegisterHandler(biz.EventUploadCompleted, handlerSvc.Handle)
}

实际干活的是 NewFileEventHandler,它拿到 FileChangedPayload 后,遍历 FileIDs / FolderIDs,调用 shareRepo.DeleteByItemID 取消相关分享——这正是 Q1 里说的"删文件异步取消分享":

// internal/event/handlers.go(节选)
for _, id := range payload.FileIDs {
	if err := shareRepo.DeleteByItemID(ctx, "file", id); err != nil {
		log.Error("event: failed to cancel share for deleted file", "file_id", id, "err", err)
	}
}
for _, id := range payload.FolderIDs {
	if err := shareRepo.DeleteByItemID(ctx, "folder", id); err != nil {
		log.Error("event: failed to cancel share for deleted folder", "folder_id", id, "err", err)
	}
}

⚠️ 坑:Handle 内部对 payload 做 JSON 往返 parsePayload(先 MarshalUnmarshal 到具体结构),因为 Event.Payloadinterface{},从 Kafka 反序列化出来默认是 map[string]interface{},必须再转成 biz.FileChangedPayload。如果哪天改了 Payload 字段,记得同步改对应结构。


五、消费幂等:Kafka 会重复投递,必须幂等

Q5. Kafka 至少一次语义会重复投递,为什么必须做消费幂等?Redis SetNX 怎么实现 CheckAndSet?

答: 类比:你点了一次外卖,但平台因为网络抖动以为没收到,又发了一遍订单。商家如果没做"同一订单号只做一次"的校验,就会做两份饭。Kafka 同理——它保证"至少一次"(at-least-once),也就是说同一条消息可能投递多次

为什么会重复?看本项目的提交策略:代码里是先处理成功,再 CommitMessages 提交 offset。如果在"处理成功"和"提交 offset"之间消费者崩溃了(或 rebalance 把分区分给别人了),那么这条消息的 offset 没被提交,重启/其他实例会重新拉到这条消息再处理一次。所以"处理成功但未提交"的事件必然可能重复执行。

因此消费端必须幂等:同一个事件处理一次和处理一百次,业务结果一致。

项目用 Redis 的 SetNX(SET if Not eXists)实现幂等,键是 idempotency:<IdempotencyKey>,TTL 24 小时。看 RedisIdempotencyStore

// internal/event/handlers.go
func (s *RedisIdempotencyStore) CheckAndSet(ctx context.Context, key string, ttl time.Duration) (bool, error) {
	result, err := s.rdb.SetNX(ctx, key, "1", ttl).Result()
	if err != nil {
		return false, fmt.Errorf("event: idempotency check failed: %w", err)
	}
	return result, nil
}

SetNX 是原子的:键不存在才设置,返回 true(首次);已存在返回 false(重复)。在 processMessage 里这样用:

// internal/event/consumer.go
if c.idempotencyStore != nil {
	ok, err := c.idempotencyStore.CheckAndSet(ctx, "idempotency:"+event.IdempotencyKey, 24*time.Hour)
	if err != nil {
		c.logger.Error("event: idempotency check failed", "err", err)
	}
	if !ok {
		c.logger.Debug("event: duplicate event, skipping", "id", event.ID, "key", event.IdempotencyKey)
		c.reader.CommitMessages(ctx, msg) // 重复事件直接提交跳过
		return
	}
}

为什么 TTL 是 24 小时? 因为 Kafka 的重复投递窗口通常很短(崩溃恢复、rebalance 在秒级到分钟级),但保留 24h 能覆盖"消费者停了很久又起来重放"的边界情况,又不会让 Redis key 永久膨胀。IdempotencyKey 来自生产者的 Type:ID,保证每条消息全局唯一,重复投送时 key 相同,于是被 SetNX 拦下。

⚠️ 坑:SetNX 失败(Redis 不可用)时,代码只记日志、不拦处理。这是有意为之的"降级"——宁可可能重复处理,也不因为 Redis 挂了就让消息卡住。但你要清楚:一旦 Redis 抖动,幂等保护就失效,业务必须自身能容忍重复(比如 DeleteByItemID 本就是幂等的"按 item 删除")。这也是为什么下游取消分享用"按 ID 删除"而非"计数减一"——天然幂等。

下面用流程图展示"幂等去重"的完整判定路径。

flowchart TD
    A[收到 Kafka 消息] --> B[JSON 反序列化 Event]
    B --> C{反序列化失败 问号}
    C -->|是| Z[直接 Commit 并丢弃]
    C -->|否| D[检查幂等存储]
    D --> E{CheckAndSet
idempotency 冒号 key 24h} E -->|false 已存在| Z E -->|true 首次| F[查找事件类型处理器] F --> G{存在 handler 问号} G -->|否| Z G -->|是| H[执行业务逻辑] H --> I[成功则 CommitMessages]

六、重试与死信:处理失败怎么办

Q6. 处理失败怎么办?重试退避与死信队列 DLQ 怎么设计?

答: 类比:快递员第一次敲门没人,不会立刻放弃,而是过一会儿再来(退避重试);但如果连来三次都没人,就把包裹放到"问题件柜"(死信队列),既不丢件,也不阻塞后续派送。

processMessage 里的重试逻辑(注意与生产者侧不同,消费侧退避是 200*(i+1)ms):

// internal/event/consumer.go
var lastErr error
for i := 0; i < 3; i++ {
	if err := handler(ctx, &event); err != nil {
		lastErr = err
		c.logger.Error("event: handler failed, retrying", "type", event.Type, "id", event.ID, "attempt", i+1, "err", err)
		time.Sleep(time.Duration(200*(i+1)) * time.Millisecond)
		continue
	}
	// 成功 — 提交消息
	if err := c.reader.CommitMessages(ctx, msg); err != nil {
		c.logger.Error("event: failed to commit message", "err", err, "id", event.ID)
	}
	return
}

三次重试的等待分别是:i=0 → 200ms,i=1 → 400ms,i=2 → 600ms。这是线性退避(不是严格指数,但随次数递增,给下游恢复时间)。注意重试是在同一条消息上同步进行的——这期间该分区后续消息会被阻塞,所以重试次数不能太多、退避不能太长,否则造成积压。

三次都失败后,进入死信队列逻辑:

// internal/event/consumer.go
dlqTopic := msg.Topic + ".dlq"
if c.dlqProducer != nil {
	if derr := c.dlqProducer.Publish(ctx, dlqTopic, &event); derr != nil {
		c.logger.Error("event: failed to forward to DLQ", "dlqTopic", dlqTopic, "id", event.ID, "err", derr)
	} else {
		c.logger.Warn("event: forwarded to DLQ after exhausting retries", "dlqTopic", dlqTopic, "id", event.ID, "type", event.Type)
	}
} else {
	c.logger.Error("event: handler failed after 3 retries and no DLQ producer configured; consuming anyway", "type", event.Type, "id", event.ID, "err", lastErr)
}
c.reader.CommitMessages(ctx, msg)

几个关键点:

  • DLQ topic 命名 = 原 topic + .dlq(如 cloud_disk_file_events.dlq)。约定优于配置,监控和人工排查时一眼能对应。
  • 死信消息原样转发&event),保留完整 IdempotencyKeyPartitionKeyPayload,方便后续人工/定时任务重放。
  • 转发 DLQ 后仍然 CommitMessages:意思是"这条我处理不了,但别再重试了,标记已处理",否则它会一直卡在分区头阻塞后续。
  • 如果没有配 DLQ producer(c.dlqProducer == nil:代码只记 error 然后照常提交——这是"无 DLQ 时的妥协":消息不会丢进死信队列,但也无法再被处理,只能靠日志人工发现。这是个明显的能力短板,生产环境务必配 DLQ producer。

⚠️ 坑:死信队列只是"隔离问题消息",不等于"问题解决"。你需要额外的死信消费者 / 监控告警去消费 .dlq topic,人工或自动修复后重放。否则 DLQ 只是把坑埋起来。

下面用状态机描述"一条消息从拉取到终态"的完整生命周期。

stateDiagram-v2
    [*] --> 处理中
    处理中 --> 成功: handler 返回 nil
    成功 --> 已提交: CommitMessages
    已提交 --> [*]
    处理中 --> 重试1: 第1次失败
    重试1 --> 处理中: 等待 200ms 后退避
    处理中 --> 重试2: 第2次失败
    重试2 --> 处理中: 等待 400ms 后退避
    处理中 --> 重试3: 第3次失败
    重试3 --> 处理中: 等待 600ms 后退避
    处理中 --> 死信队列: 3次耗尽
    死信队列 --> 已提交: 转发 topic 点 dlq 后 Commit
    已提交 --> [*]

七、分区与有序:PartitionKey 为什么按用户维度

Q7. 分区与有序:PartitionKey 为什么按用户维度?怎么保证同一用户事件有序?

答: 类比:银行有很多个柜台(分区),但同一个客户的业务必须排在同一个柜台的同一个队列里,才能保证"先取钱后转账"的顺序。如果同一个客户的业务被分到不同柜台,两个柜台并行处理,顺序就乱了。

Kafka 的有序性保证是**“分区内有序,跨分区不保证”**。消息落到哪个分区由 Key 的 hash 决定:

partition = hash(Key) % 分区数

项目在 eventPublisherAdapter.Publish 里把 PartitionKey 设成 "user:%d"(用户维度):

// internal/event/biz_publisher.go
case *biz.FileChangedPayload:
	evt.PartitionKey = fmt.Sprintf("user:%d", p.UserID)
case *biz.SharePayload:
	evt.PartitionKey = fmt.Sprintf("user:%d", p.UserID)
case *biz.UploadCompletedPayload:
	evt.PartitionKey = fmt.Sprintf("user:%d", p.UserID)

于是对同一个用户hash("user:42") 永远一样,他产生的所有事件(删除、移动、上传、分享)都进同一个分区,被同一个消费者实例按序处理。而不同用户的 key 不同,散列到不同分区,可以并行处理,互不阻塞——这就是"单用户有序、多用户并行"的理想模型。

为什么选"用户维度"而不是"文件维度"?因为本项目最关心的顺序语义是"同一个人的操作序列别乱"(比如先删文件再取消分享,顺序不能反),以用户为粒度分区既能保序,又把并发度控制在"用户数"级别,足够高。

如果 PartitionKey 为空,生产者会回退到 IdempotencyKeyType:ID),此时每条消息 key 都不同,会被打散到任意分区,完全无序——所以 adapter 里主动设用户维度非常关键。

下面这张图直观展示"按用户分区,同用户有序、不同用户并行"。

graph LR
    U1[用户A 删除事件] -->|Key 等于 用户冒号1| P1[分区0]
    U2[用户A 上传事件] -->|Key 等于 用户冒号1| P1
    U3[用户B 删除事件] -->|Key 等于 用户冒号2| P2[分区1]
    U4[用户B 分享事件] -->|Key 等于 用户冒号2| P2
    P1 --> C1[消费者实例1
按序处理] P2 --> C2[消费者实例2
并行处理]

⚠️ 坑:分区数一旦定下(如 3 个),用户量增大后 hash % 分区数 会把用户重新分布,但已有数据不会迁移;如果以后要扩分区,同用户的 key 会落到新分区,可能导致短暂的顺序错乱(旧分区还有残留未消费消息)。扩容分区时要评估顺序影响,必要时按 key 做重放。


八、降级:Kafka 没配置时怎么办

Q8. Kafka 未配置时怎么办?降级 noop 怎么实现的?什么是"可选依赖"设计?

答: 类比:公司前台(消息总线)今天请假了,但你该发的文件照样发、该删的照样删——只是"通知同事"这件事暂时不做。业务主干不能因为"没有消息总线"就全体瘫痪。

项目在三个层次都做了降级,统一返回 noop(空实现):

  1. 生产者降级event/producer.go):
// internal/event/producer.go
func NewProducer(cfg *conf.Data_Kafka, writer *kafka.Writer) Producer {
	if cfg == nil || len(cfg.Addrs) == 0 {
		return noopProducer{}
	}
	return &kafkaProducer{writer: writer}
}
  1. 适配器降级event/biz_publisher.go):业务层拿到的 EventPublisher 也可能是 noop:
// internal/event/biz_publisher.go
func NewEventPublisherAdapter(producer Producer, cfg *conf.Data_Kafka) biz.EventPublisher {
	if cfg == nil || len(cfg.Addrs) == 0 {
		return noopEventPublisher{}
	}
	// ...
}
  1. 消费者降级event/consumer.go):没有 Kafka 或没有待订阅 topic,返回 noop consumer。

  2. 底层 Writer 降级data/kafka.go):即使 Addrs 为空,也返回空的 &kafka.Client{}&kafka.Writer{},不让 Wire 注入失败。

所有 noop 实现都在 internal/event/noop.go

// internal/event/noop.go
type noopProducer struct{}
func (noopProducer) Publish(ctx context.Context, topic string, event *Event) error { return nil }
func (noopProducer) Close() error                                                { return nil }

type noopConsumer struct{}
func (noopConsumer) RegisterHandler(eventType string, handler EventHandler) {}
func (noopConsumer) Start(ctx context.Context) error                        { return nil }
func (noopConsumer) Stop() error                                            { return nil }

type noopEventPublisher struct{}
func (noopEventPublisher) Publish(ctx context.Context, eventType string, payload interface{}) error {
	return nil
}

注意 noop 的 Publish 全部返回 nil——也就是说"调用方以为发出去了,其实啥也没干"。调用方不会报错,业务主流程不受影响。这就是**“可选依赖”(optional dependency)**设计:消息总线是增强能力,缺失时系统退化为"纯同步模式",核心 CRUD 照常工作。

⚠️ 坑:降级虽好,但代价是"事件驱动能力整体失效"。没有 Kafka 时,删文件不会异步取消分享,分享链接会残留——这是数据不一致。所以降级只适合"开发/测试环境"或"可接受最终不一致"的场景,生产环境必须把 Kafka 配上,并通过健康检查告警"Kafka 不可用"。

这也呼应了 Q2 里 biz 调用 Publish_ = 忽略错误:主流程从设计上就不依赖发布成功,降级与"发布失败不影响主流程"是同一套工程哲学。


九、offset 提交时机与可靠性语义

Q9. offset 提交时机怎么选?at-least-once 和 exactly-once 有什么区别?

答: 这是消息队列面试的"必考题"。先说结论:本项目是 at-least-once(至少一次)语义,靠消费端幂等兜底重复。

提交时机的两种极端:

  • 先提交 offset,再处理:如果提交后、处理前崩溃,这条消息不会再被处理 → 会丢失(at-most-once)。
  • 先处理,再提交 offset(本项目做法):如果处理后、提交前崩溃,消息会被重投 → 会重复(at-least-once)。

项目选了"先处理后提交",因为重复比丢失好处理(重复可幂等,丢失不可恢复)。看代码:handler 成功返回后才 CommitMessagesfor 循环里若重试耗尽,转发 DLQ 后也 CommitMessages。失败的路径绝不提交,从而触发重投。

// 成功路径
if err := handler(ctx, &event); err != nil { /* 重试 */ }
c.reader.CommitMessages(ctx, msg) // 仅成功才到这

但注意一个细节考:本项目 ReaderConfig 同时设了 CommitInterval: 1 * time.Second(周期自动提交)。这意味着即使你显式 CommitMessages,底层也可能在周期到点时自动提交。自动提交的好处是"崩溃后最多重复一小批",坏处是"可能在你处理失败前就把 offset 提交了"(若自动提交周期先于你的手动提交触发)。因此严格追求精确语义时,应该把 CommitInterval 设为 0(禁用自动提交),只用手动 CommitMessages。本项目因为消费端幂等做了兜底,所以自动提交带来的少量重复是可接受的。

三种语义对比(面试可直接背):

语义含义本项目是否采用实现成本
at-most-once最多一次,可能丢
at-least-once至少一次,可能重中(需消费端幂等)
exactly-once精确一次,不丢不重高(事务 / Kafka Streams / 幂等生产者 + 事务消费者)
  • exactly-once 在 Kafka 里靠"幂等生产者(enable.idempotence)+ 事务(transactional.id)+ 消费-转换-生产原子化"实现,但吞吐量下降、实现复杂。本项目没上,因为"重复可幂等"已经够用,且成本更低。

下面用图总结"提交时机 vs 语义"的取舍。

graph TD
    A[消息拉取] --> B{先提交 还是 先处理}
    B -->|先提交 offset| C[崩溃则丢失
at most once] B -->|先处理 再提交| D[崩溃则重投
at least once] D --> E[消费端幂等
Redis SetNX 兜底] E --> F[最终一致 不丢不重] F --> G[exactly once 需事务
本项目未采用]

⚠️ 坑:不要误以为"消费端幂等 = 精确一次"。幂等只解决"重复处理业务结果一致",但消息在 Kafka 里确实被处理了多次(多次网络 IO、多次 DB 操作)。真正的 exactly-once 还要在传输层用事务保证"读-处理-写"原子。


十、优雅退出:消费者生命周期

Q10. 消费者怎么优雅退出?Stop 的生命周期是怎样的?

答: 类比:商场关门(停机)时,不能把正在结账的顾客直接赶走,要等当前这笔办完、保安锁好门再走。kafkaConsumer.Stop 做的就是这件事。

Start 时创建一个可被取消的 ctx,并起一个 goroutine 跑 consumeLoop

// internal/event/consumer.go
func (c *kafkaConsumer) Start(ctx context.Context) error {
	ctx, c.cancel = context.WithCancel(ctx)
	c.wg.Add(1)
	go c.consumeLoop(ctx)
	return nil
}

consumeLoop 里每次循环开头 select 检查 ctx.Done()FetchMessage 返回 ctx.Err() != nil 时也直接 return——所以一旦 cancel,循环会尽快退出。

Stop 的三步是标准范式:

// internal/event/consumer.go
func (c *kafkaConsumer) Stop() error {
	if c.cancel != nil {
		c.cancel()        // 1. 通知 consumeLoop 退出
	}
	c.wg.Wait()           // 2. 等 goroutine 真正结束(避免 reader 被 close 时还有人在用)
	return c.reader.Close() // 3. 关闭 reader,释放连接
}

顺序很关键:先 cancel 再 wg.Wait()Close()。如果先 Close() 而循环还在 FetchMessage,会报 “connection closed” 类错误;先 cancel 让循环自然退出,再等它结束,最后才关资源,干净无竞态。

这套 ConsumerServerconsumer_server.go)还实现了 Kratos 的 transport.Server 接口,所以可以随应用一起启停——应用启动调用 Start,应用关闭(收到 SIGTERM)调用 Stop,生命周期完全托管给 Kratos:

// internal/event/consumer_server.go
func (s *ConsumerServer) Start(ctx context.Context) error {
	if s.handlerService != nil {
		RegisterAllHandlers(s.consumer, s.handlerService)
	}
	if s.idempotency != nil {
		if c, ok := s.consumer.(*kafkaConsumer); ok {
			c.SetIdempotencyStore(s.idempotency)
		}
	}
	if p, ok := s.consumer.(*kafkaConsumer); ok && s.dlqProducer != nil {
		p.SetDLQProducer(s.dlqProducer)
	}
	return s.consumer.Start(ctx)
}

注意这里的类型断言 s.consumer.(*kafkaConsumer):只有"真的 Kafka 消费者"才需要 set 幂等存储和 DLQ producer;如果是 noop 消费者,这些 set 直接跳过——因为 noop 根本不消费,设了也没用。这再次体现降级设计的自洽。

⚠️ 坑:Stopcancel() 之后,consumeLoop 可能在 FetchMessage 阻塞中,要等 MaxWait(1s)才因 ctx 取消而返回。所以优雅退出最多会多等约 1 秒,这是正常的"收尾时间",不要误判为卡死。


十一、面试延伸:不丢消息、防积压、选型

Q11. 面试延伸:如何保证不丢消息?如何避免消息积压?Kafka 与 RabbitMQ 怎么选型?

答: 这是把前面所有点串起来的"压轴题"。

1)如何保证不丢消息(生产者 + Broker + 消费者三端)

  • 生产者端:开启重试(本项目 Publish 里已有 3 次重试 + 退避),并把 RequiredAcksRequireOne 提到 RequireAll(即 acks=all),配合 min.insync.replicas >= 2,确保 leader 和足够多的 follower 都写入才返回成功。
  • Broker 端:topic 设置多副本(replication.factor >= 3),防止单点磁盘坏掉丢数据。
  • 消费者端:本项目已做"先处理后提交 offset",处理失败不提交 → 重投;再配合消费端幂等,重复也不怕。

⚠️ 一句话记忆:丢消息主要在"生产没确认 / Broker 单副本 / 消费先提交后处理"三处,本项目分别用重试+acks、多副本、先处理后提交来堵。

2)如何避免消息积压

  • 横向扩容消费者:同 Group 内增加实例,Kafka 会自动 rebalance 把分区分给新实例。本项目分区是并行单位,分区数决定了消费并发上限,所以分区数要预留(别只设 1 个分区,否则加多少实例都只有 1 个在干活)。
  • 提升单条处理速度:handler 里别做慢操作;本项目的 DeleteByItemID 是批量按 ID 删,够快。若 handler 慢,考虑批量消费、异步落库。
  • 监控 lag:监控消费组 consumer lag(最新 offset - 已提交 offset),超过阈值告警。本项目没自带监控,生产要补。
  • 死信隔离:处理不了的进 DLQ,避免单条毒消息(poison message)卡住整个分区——这正是 Q6 设计 DLQ 的意义。

3)Kafka vs RabbitMQ 选型

维度KafkaRabbitMQ
模型分区日志、拉模式队列、推模式(Exchange 路由)
顺序分区内严格有序单队列有序,多队列不保证
吞吐极高(磁盘顺序写)中(适合低延迟小消息)
消费模式消费组、可重放历史点对点 / 发布订阅,消费即删
典型场景事件流、日志、大数据管道任务队列、RPC 式请求、复杂路由
本项目契合度(事件驱动、需重放、高吞吐)

本项目选 Kafka 的理由:事件是持续流、需要按用户重放/回溯、吞吐大、且要消费组做水平扩展——这些全是 Kafka 的强项。RabbitMQ 更适合"一条任务只被一个 worker 处理一次"的传统任务队列场景。

⚠️ 坑:别用 Kafka 做"必须严格只处理一次且需要复杂路由(按 header 路由到不同队列)“的需求,那是 RabbitMQ Exchange 的强项。技术选型看场景,没有银弹。


自测题与动手练习

自测题(合上代码自问自答)

  1. 云盘项目引入事件驱动解决了哪三个问题?举一个"文件删除后取消分享"之外的事件例子(提示:回收站清理、上传完成)。
  2. biz.EventPublisher 接口只有 Publish 一个方法,为什么这样设计?eventPublisherAdapter 承担哪两项职责?
  3. 生产者 Publish 里分区键为空时回退到什么?为什么 adapter 要主动设 PartitionKey = "user:%d"
  4. 消费端幂等用 Redis 什么命令实现?键的格式和 TTL 是什么?为什么 Redis 不可用时仍然继续处理而不是阻塞?
  5. 消费失败重试几次?每次退避多久?耗尽后去哪个 topic?没有 DLQ producer 时会怎样?

动手练习

  1. 在本地用 Docker 起 Kafka,把 conf.Data_Kafka.Addrs 配上,手动构造一个 FileChangedPayload{UserID:1, FileIDs:[]uint64{100}}file.deleted 事件,验证分享被取消,并观察 Redis 里是否多了一个 idempotency:* 的 key。
  2. 故意让 shareRepo.DeleteByItemID 返回错误,连续触发 3 次,确认日志里出现 3 次 handler failed, retrying 且退避为 200/400/600ms,最后消息被转发到 *.dlq topic。
  3. conf.Data_Kafka.Addrs 清空重新启动,确认服务正常启动、删文件接口正常返回(验证 noop 降级不崩主流程),同时在日志里看到 Kafka config is empty 之类的告警。

本章小结

本文从 Kratos 云盘项目的真实代码出发,串起了"Kafka 事件驱动与可靠消费"的完整知识链:

  • 为什么用:用"删文件异步取消分享"说明了解耦、削峰、最终一致,并给出 7 种事件常量。
  • 分层解耦biz.EventPublisher 接口 + eventPublisherAdapter 做事件类型→topic 映射,依赖单向、无循环。
  • 生产segmentio/kafka-gokafka.Writer,配置 RequiredAcks/BalancerPublish 带 3 次退避重试,消息带 Key(分区键)和 Header。
  • 消费kafka.Reader 消费组 + consumeLoop 拉取 + RegisterHandler 分发,单分区内严格有序。
  • 幂等:Kafka at-least-once 会重复投递,用 Redis SetNXCheckAndSet,键 idempotency:<key>,TTL 24h)去重,重复直接 commit 跳过。
  • 重试与 DLQ:消费侧 3 次重试、退避 200*(i+1)ms,耗尽转发 topic+".dlq",无 DLQ 时仅 log 的妥协。
  • 有序PartitionKey="user:%d" 保证同用户同分区有序、不同用户并行。
  • 降级:Kafka 未配置时返回 noop 实现,业务主干不崩,体现"可选依赖”。
  • 可靠性:先处理后提交 offset(at-least-once),靠幂等兜底;exactly-once 成本高未采用。
  • 生命周期Stop = cancel + wg.Wait + Close,优雅退出;ConsumerServer 接入 Kratos 启停。
  • 面试延伸:不丢消息三端保障、防积压靠扩容分区/监控 lag/DLQ 隔离、Kafka 与 RabbitMQ 选型差异。

记住一句话核心:事件驱动把"主流程"和"副作用"拆开,可靠消费靠"幂等 + 重试 + 死信"三板斧,而分区键决定顺序边界,降级保证主干不崩。 这既是本项目的设计哲学,也是面试回答事件驱动类问题的通用框架。

复习提示:
  • 解耦的核心biz.EventPublisher 接口让业务逻辑不感知消息队列的存在——这是依赖倒置的经典实践。
  • 幂等是关键:Kafka at-least-once 语义必然重复投递,所以 CheckAndSet(Redis SETNX)是必须的补偿机制。
  • 降级的优雅:Kafka 未配置时返回 noop,业务不崩——这是"可选依赖"思想的体现。
  • 分区键 = 顺序边界user:%d 保证同一用户的消息有序,不同用户并行处理。
面试官
为什么 Kafka 消费侧要"先处理后提交 offset"而不是反过来?
候选人
因为如果先提交 offset 再处理失败,消息就永久丢失了(offset 已提交,下次不会再消费这条)。

标准流程
① Consumer 拉取消息 → 执行业务逻辑
② 业务成功后 → commitSync() / commitAsync()
③ 如果处理失败 → 触发重试逻辑

但这也带来新问题:如果处理一直失败(如 DB 挂了),会反复消费同一条消息。
所以项目设计了 重试 → DLQ 的机制:3 次重试后移入死信队列,避免阻塞消费者。

面试加分点:提到 Kafka 的 enable.auto.commit=false 必须手动提交,以及"处理中崩溃"场景下会重复消费(恰好一次需要额外事务支持)。
About Me

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

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

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

目标

学AI,加油!加油!