事件驱动模块实现教程(Kafka)

2025-08-04T14:11:02+08:00 | 27分钟阅读 | 更新于 2025-08-04T14:11:02+08:00

@

学习目标

学完本章你应该能够:

  1. 用「事件驱动 vs 同步调用」的视角,说清为什么云盘要引入 Kafka 做异步事件(解耦、提速、最终一致),以及事件主要用在哪些场景(文件删除→取消分享、上传完成、分享创建/取消)。
  2. 讲清楚一条事件从「业务发布」到「消费端处理」的完整链路:定义 → 发布 → 订阅 → 幂等 → 重试 → 分发,并画出调用关系图。
  3. 在面试里把 Kafka 可靠性(RequiredAcks/ISR/手动提交)、Consumer Group、幂等消费、投递语义、死信队列(DLQ)、消息积压等难点讲成有逻辑的故事。
  4. 纠正一个最常见的认知错误:幂等键(去重)和分区键(保序)是两件事——本项目源码把两者混为一谈,导致"同一文件事件保序"的承诺其实不成立,能讲清正确做法。
  5. 指出本模块的 4 个商用缺口:事件静默丢失(吞错 + 无 outbox)无事务发件箱无死信队列RequireOne 丢消息风险,并给出正确实现。
  6. 指出亿级流量下分区并行、批量发送、压缩、背压、DLQ 监控等优化手段分别在解决什么问题。

前置知识

  • 本系列前 6 章(Kratos 项目骨架、biz/data 分层、Wire 依赖注入、用户/文件/分享模块)。
  • Go 基础:interface、goroutine、contextsync.WaitGroup
  • 消息队列基本概念(Producer / Consumer / Topic / Partition / Consumer Group)。

本章你会动手做的事

  1. 跟着代码走一遍 Publish 的「生成事件 ID → 序列化 → 退避重试」三步,并看清 Kafka Key 到底用的什么。
  2. 在本地用 Redis 跑一次 SetNX 幂等检查,观察重复消息被跳过;再观察"随机消息 ID 作为 Key"时,同一文件的两条事件被分到不同分区。
  3. processMessage 补一段"重试耗尽 → 发往 DLQ"的真实逻辑,体会失败处理的完整闭环。

类比:事件驱动模块像一套「公司的内部邮局」。业务部门(biz 层)只管把写好的信(事件)投进邮箱(Kafka),不用管谁去送、什么时候送到;邮差(消费者)异步取信、查重(幂等)、失败重试,最后把信送到对应部门(处理器)办后续手续。这样「写主流程」和「后续副作用」解耦,主接口能立刻返回。

flowchart LR
    B[biz 业务层] -->|Publish 事件| P[Producer 生产者]
    P -->|写入| K[(Kafka Topic)]
    K -->|拉取| C[Consumer 消费者组]
    C -->|SetNX 查重| R[(Redis 幂等存储)]
    C -->|分发| H[业务处理器
file/share/recycle/upload] H -->|副作用| DB[(数据库最终一致)]

一、技术栈与中间件

事件驱动模块围绕 Kafka 消息中间件构建,结合 Redis 实现幂等存储,并通过 Kratos 框架的生命周期钩子统一管理启停。下表列出本模块用到的所有核心技术:

技术 / 中间件用途说明
Apache Kafka分布式消息中间件,承载异步事件的存储与转发,解耦生产者与消费者
segmentio/kafka-go纯 Go 实现的 Kafka 客户端库,提供 Writer(生产者)和 Reader(消费者)抽象
Consumer Group(消费者组)Kafka 内置机制,让多个消费者实例分摊同一 Topic 的分区,实现水平扩展与负载均衡
Redis(go-redis)作为幂等性存储后端,通过 SetNX 原子命令保证「同一条 Kafka 消息」只被处理一次
重试退避(Backoff)生产者发送失败、消费者处理失败时,按退避策略重试,提升最终成功率
幂等键(IdempotencyKey)每条事件携带的唯一消息 ID,作为 Redis 幂等键,保证重复投递可去重
分区键(PartitionKey)业务实体标识(文件 ID / 分享 ID),作为 Kafka 消息 Key 控制分区路由与保序
死信队列(DLQ)重试耗尽仍失败的消息转入独立 Topic,供排查、人工干预或延迟重处理
事务发件箱(Outbox)商用补强:业务落库与事件写入同一本地事务,再由中继者投递 Kafka,杜绝静默丢失
Kratos transport.ServerKratos 框架的服务生命周期接口,将消费者包装为标准 Server,随应用一起启停
Google Wire编译时依赖注入框架,通过 ProviderSet 统一装配事件模块的所有组件
Noop 空实现模式当 Kafka 未配置时返回空实现,保证业务逻辑不因中间件缺失而崩溃

二、实现思路流程(总体)

类比:把这条链路想成「快递分拣流水线」——先按品类定义包裹(事件定义),扫码入库(生产者发布),传送带分到不同站点(消费者订阅),扫码查重避免重复派送(幂等检查),派送失败稍等再试(重试退避),还不行就进「问题件库房」(死信队列),最后送到对应柜台办手续(业务分发)。

flowchart LR
    A[事件定义
7 种事件类型] --> B[生产者发布
生成事件ID+幂等键+分区键] B -->|JSON 发送
失败退避重试| C[消费者订阅
Consumer Group 拉取] C --> D[幂等检查
Redis SetNX 按幂等键] D -->|重复 已处理| S((直接提交跳过)) D -->|首次| E[重试退避
N 次] E -->|成功| F[业务处理器分发] E -->|耗尽| G[发往 DLQ 死信队列] F -->|副作用| DB[(数据库最终一致)]

事件驱动模块的整体实现思路遵循「定义 → 发布 → 订阅 → 幂等 → 重试 → 分发」的链路,完整流程如下:

┌──────────────────────────────────────────────────────────────────────┐
│                       事件驱动模块整体流程                            │
└──────────────────────────────────────────────────────────────────────┘

  ① 事件定义(7 种事件类型)
     ├─ file.deleted / file.moved / file.trashed   (文件事件)
     ├─ share.created / share.cancelled            (分享事件)
     ├─ recycle.clean                              (回收站事件)
     └─ upload.completed                           (上传事件)
                    │
                    ▼
  ② 生产者发布(biz 层 → Producer)
     ├─ eventPublisherAdapter 把事件类型映射到 Topic
     ├─ Producer.Publish 生成事件 ID + 幂等键(唯一)
     ├─ 适配器按 payload 计算分区键(业务实体,用于保序)
     ├─ JSON 序列化,Kafka Key = 分区键(保序),Value = 事件
     └─ 发送失败 → 退避重试(商用建议 RequireAll + 异步/Outbox)
                    │
                    ▼
  ③ 消费者订阅(Consumer Group)
     ├─ kafka.Reader 以 GroupID 订阅多个 Topic
     ├─ FetchMessage 拉取消息 → 反序列化为 Event
     └─ 按 event.Type 路由到对应处理器
                    │
                    ▼
  ④ 幂等检查(Redis SetNX,按"幂等键")
     ├─ key = "idempotency:" + IdempotencyKey(唯一消息ID)
     ├─ SetNX 成功(首次)→ 继续处理
     ├─ SetNX 失败(同一条消息重投)→ 直接提交,跳过
     └─ TTL = 24 小时,避免无限累积
                    │
                    ▼
  ⑤ 重试退避 + 死信队列
     ├─ 处理器执行失败 → 退避重试(200ms / 400ms / 600ms,且可被 ctx 中断)
     ├─ 成功 → CommitMessages 提交偏移量
     └─ 重试耗尽 → 发往 DLQ 死信 Topic,再提交原消息(不阻塞主队列)
                    │
                    ▼
  ⑥ 业务处理器分发
     ├─ file 事件 → 调用 shareRepo.DeleteByItemID 取消相关分享
     ├─ share 事件 → 记录日志(预留扩展:通知/统计)
     ├─ recycle 事件 → 记录日志(预留扩展)
     └─ upload 事件 → 记录日志(预留扩展:秒传索引/审核)

核心设计原则

  1. 接口隔离biz 层只依赖 EventPublisher 接口,不感知 Kafka 细节,避免循环依赖。
  2. 失败兜底:Kafka 未配置时返回 noopProducer / noopConsumer,业务逻辑正常运行。
  3. 幂等优先:消费前先做幂等检查,重复消息直接丢弃,保证同一条消息的副作用只发生一次。
  4. 优雅启停:消费者通过 ConsumerServer 包装为 transport.Server,随 Kratos 应用生命周期启停。
  5. 商用补强(修正点):发布端不吞错 + 事务发件箱保证"不丢";消费端 DLQ 保证"失败可查";RequireAll 保证"Broker 不丢"。

三、面试常问知识点与难点

1. Kafka 消息可靠性保证

Kafka 通过三层机制尽量保证消息不丢失:生产者 RequiredAcks副本机制 ISR消费者手动提交偏移量

  • 本项目 kafka.Writer 配置 RequiredAcks: kafka.RequireOne,即 Leader 副本写入即返回。风险:若 Leader 刚 ack 就宕机、且消息尚未同步到 ISR 中的 Follower,这条消息就永久丢失了。
  • 商用修正:生产环境应改为 kafka.RequireAll(WaitForAll),并要求 replication_factor >= 3min.insync.replicas >= 2,确保消息被多数副本确认后才算写入成功。
  • 消费端本项目使用 CommitMessages 在业务处理成功之后才提交偏移量,这是典型的 At-Least-Once 语义(下文第 9 节展开)。

2. Consumer Group 机制

Consumer Group 是 Kafka 实现负载均衡的核心:同一 Group 内消费者分摊 Topic 的所有分区,每个分区同一时刻只被组内一个消费者消费。本项目 GroupID = "cloud-disk-consumer",多实例部署时 Kafka 自动在实例间重分配分区,实现消费能力线性扩展。注意:实例数不应超过分区数,否则多余实例空闲。

3. 幂等消费实现方案(含认知纠偏)

本项目用 Redis SetNX 做幂等:消费前以 idempotency:{IdempotencyKey} 为 key 调用 SetNX,返回 true 表示首次、返回 false 表示重复(同一条消息重投)直接丢弃,TTL 24 小时。

⚠️ 关键认知纠偏(务必记牢):本项目源码里 IdempotencyKey = "事件类型:随机事件ID",它是逐条消息唯一的。因此这个 SetNX 只能去重"Kafka 把同一条消息重投了两次"这类情形(崩溃/Rebalance 后重投),它不是业务级幂等

  1. 它保证不了「同一文件连续两条不同事件」被去重——因为它们 ID 不同、键不同,会各处理一次(所幸 DeleteByItemID 本身是幂等操作,删两次无害)。
  2. 它保证不了「同一业务操作因 API 层重试而发了两次事件」被去重——两次发布产生两个随机 ID,键不同,会各处理一次。

真正的「业务级 Exactly-Once」要靠幂等处理器(副作用本身可重复执行,如按主键 upsert/delete)或业务键幂等(用文件 ID 当幂等键)。本项目靠"处理器副作用幂等"这条兜底,而非靠幂等键。

flowchart TD
    M[消费到一条消息] --> I[Redis SetNX
idempotency:幂等键] I -->|true 首次| P[执行业务处理器] I -->|false 同条重投| K[直接 CommitMessages 丢弃] P --> C[成功提交偏移量]

💡 幂等检查在 Redis 故障时应 fail-open 还是 fail-closed? 本项目选择 fail-open:Redis 报错时记录日志但继续处理(宁可重复也不丢)。原因与用户模块"黑名单 fail-closed"正相反——黑名单关乎安全,漏判会导致越权;而这里消费端是 At-Least-Once + 幂等处理器,重复处理最多是"多干一次本就幂等的活",不会破坏正确性。这是「重复比丢失安全」场景下的正确取舍。

4. 消息重试与死信队列(DLQ)

消费者在处理器失败时进行退避重试(本项目 200ms / 400ms / 600ms)。重试耗尽后本项目仅打日志并提交,毒丸消息直接消失——这是商用缺口。

正确做法:将最终失败的消息发送到独立 DLQ Topic(如 cloud_disk_file_events.DLQ),元数据里带上失败原因与原始消息,再由专门的「死信处理服务」消费:人工排查、延迟重投或告警。这样主队列不会被毒丸阻塞,失败也有据可查。

flowchart TD
    H[处理器失败] --> R{重试<3次?}
    R -->|是| W[退避等待 可被ctx中断]
    W --> H
    R -->|否| D[发送到 DLQ Topic]
    D --> C[提交原消息 主队列继续]

5. 事件溯源与最终一致性

事件驱动天然契合最终一致性:业务操作(如文件删除)在事务提交后立即返回,同时发布事件;消费者异步处理副作用(如取消分享)。分享状态更新有几百毫秒延迟,但最终一致。代价是要处理消息丢失、重复、乱序等边界——这正是本模块要解决的。

6. 消息顺序性保证(含事实纠错)

⚠️ 原教程此处与源码矛盾,已修正:原说"以 IdempotencyKey 作消息 Key,Kafka 按 Key 哈希将同一业务实体路由到同一分区,从而保证同一文件事件保序"。但源码的 IdempotencyKey类型:随机事件ID,随机 ID 哈希后落到哪个分区是随机的,根本无法保证同一文件保序。 这是把"幂等键"和"分区键"混为一谈导致的错误。

正确做法:分区键必须取业务实体标识(文件 ID / 分享 ID),让同一实体的事件稳定哈希到同一分区,Kafka 单分区内严格有序,从而实现"同文件事件按序消费"。本项目商用修正:Event 新增 PartitionKey 字段,生产者用它作为 Kafka 消息 Key;幂等键继续用唯一消息 ID 去重。两者职责分离:

  • Kafka Key = 分区键 → 控制分区路由与顺序;
  • Redis 幂等键 = 唯一消息 ID → 控制重复投递去重。

7. 投递语义:At-Least-Once 与"不丢"的代价

Kafka 消费者有三种语义:At Most Once(先提交后处理,可能丢)、At Least Once(先处理后提交,可能重复,本项目采用)、Exactly Once(事务性提交,最复杂)。本项目选 At-Least-Once + 幂等处理器,是可靠性与复杂度的最佳平衡。

但"不丢"不只取决于消费端,还取决于生产端。本项目生产端有两个致命缺口(下一节 Outbox 详解):① 调用点 _ = Publish(...) 把错误吞掉,Kafka 不可用时事件静默丢失;② 事件在 DB 事务提交后命令式发布,进程在"提交后、发布前"崩溃也会丢事件。

8. 消息积压处理

积压由"消费速度 < 生产速度"引起,思路:① 加消费者实例(受分区数限制);② 加分区数(提升并行度,但会破坏跨分区顺序);③ 批量消费;④ 临时救援消费者快速清积压;⑤ 降级生产端(限流或关闭非核心事件)。

9. 优雅停机与 Rebalance

消费者停机若直接退出,已拉取未处理的消息在 Rebalance 后会被其他消费者重复消费。本项目通过 context.Cancel + sync.WaitGroup 实现优雅停机:Stop() 取消 context、消费循环退出当前迭代、wg.Wait() 等正在处理的消息完成、最后关连接。配合幂等机制,Rebalance 重复消费可被安全去重。

10. 可观测性:消费堆积与延迟

生产必须监控:消费 Lag(滞后量 = 最新偏移 − 已提交偏移,反映积压)、消费速率处理耗时 P99重试次数DLQ 长度。本项目仅打日志,应接 Prometheus + Grafana,Lag 超阈值告警(见五.7)。


四、亿级流量优化思路

1. Kafka 分区并行消费

并行度上限 = 分区数。亿级流量下 Topic 分区数应设为「消费者实例数 × 单实例消费线程数」(如 50 实例 × 2 = 100 分区)。注意:分区数变更会改变 Key 哈希映射,扩容分区后同一实体的旧消息可能落到不同分区,历史顺序性无法保证,需用业务键幂等兜底。

2. 消费者水平扩展

Consumer Group 天然支持水平扩展:多实例加入同 Group,Kafka 自动分配分区,增加 Pod 副本数即可线性提升消费能力。实例数勿超分区数。

3. 消息批量发送

生产者把多条事件合并为一个 Batch,减少网络往返。本项目 kafka.Writer 配置 BatchTimeout: 10msBatchSize: 1,可优化为 BatchSize: 100BatchTimeout: 5ms。批量发送是 Kafka 高吞吐核心手段。

4. 消息压缩

Kafka 支持 gzip / snappy / lz4 / zstd,JSON 事件用 lz4 或 zstd 压缩比可达 3-5 倍且速度快,显著降低 Broker 压力与带宽成本。

5. 消费者背压(Backpressure)

下游处理跟不上时消息在 Kafka 堆积。通过限制 MaxBytes 单次拉取量、使用有界队列缓冲、基于处理耗时动态调整速率实现背压,避免雪崩。

6. 分区有序保证(修正点)

需要严格顺序的场景(同一文件状态变更)应把业务实体 ID 作为分区键。本项目修正:Event 增加 PartitionKey,生产者用其作 Kafka Key,Kafka 按 Key 哈希将同一实体事件落到同一分区从而保序(见三.6)。扩容分区需重新评估顺序影响。

7. 幂等存储优化

亿级流量下 Redis 幂等键快速膨胀。优化:① TTL 精细化(按业务时效设 1 小时而非 24 小时);② 分片存储(按业务前缀分散到不同 Redis 分片);③ 布隆过滤器前置(先快判"可能重复"再走 Redis 精确检查);④ 冷数据归档

8. 监控、告警与 DLQ 治理

必须监控 Lag、消费速率、处理耗时 P99、重试次数、DLQ 长度。建议:Lag 超过阈值告警;DLQ 接入独立告警(失败消息不应长期堆积);用 Prometheus + Grafana 可视化,详见五.7。


五、详细实现流程与代码解析

以下代码片段均基于项目源码,并按商用在线服务标准修正或补全(修正处标注「修正点」)。

5.1 事件类型定义与消息格式规范(Event 结构体 + 7 种事件类型)

实现思路

事件驱动第一步是定义统一事件格式事件类型常量。本项目将业务层事件常量放 internal/biz/event.go,通用 Event 结构体放 internal/event/event.go,分层避免循环依赖。设计要点:统一 ID / IdempotencyKey / PartitionKey / Type / Source / Timestamp / Payload 字段;7 种事件类型常量集中管理;Payload 用 interface{} 由处理器按需反序列化。

修正点:原 Event 只有 IdempotencyKey 且它同时被当成 Kafka Key 用,导致"去重"与"分区保序"职责混淆(见三.6)。商用版新增 PartitionKey(业务实体键,专用于分区路由与保序),IdempotencyKey 回归"唯一消息 ID 用于去重"的本职。

关键代码

业务层事件类型常量与负载定义internal/biz/event.go):

package biz

import "context"

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

// 异步事件类型常量(共 7 种,覆盖文件 / 分享 / 回收站 / 上传四大业务域)。
const (
	EventFileDeleted    = "file.deleted"
	EventFileMoved      = "file.moved"
	EventFileTrashed    = "file.trashed"
	EventShareCreated   = "share.created"
	EventShareCancelled = "share.cancelled"
	EventRecycleClean   = "recycle.clean"
	EventUploadCompleted = "upload.completed"
)

// FileChangedPayload 文件变更事件负载
type FileChangedPayload struct {
	UserID    uint64   `json:"user_id"`
	FileIDs   []uint64 `json:"file_ids,omitempty"`
	FolderIDs []uint64 `json:"folder_ids,omitempty"`
	Action    string   `json:"action"`
	ParentID  *uint64  `json:"parent_id,omitempty"`
}

// SharePayload 分享事件负载
type SharePayload struct {
	UserID  uint64 `json:"user_id"`
	ShareID uint64 `json:"share_id"`
	UUID    string `json:"uuid"`
	Action  string `json:"action"`
}

// RecycleCleanPayload 回收站清理事件负载
type RecycleCleanPayload struct {
	Count int `json:"count"`
}

// UploadCompletedPayload 上传完成事件负载
type UploadCompletedPayload struct {
	UserID   uint64 `json:"user_id"`
	FileID   uint64 `json:"file_id"`
	FileName string `json:"file_name"`
	FileSize int64  `json:"file_size"`
	Hash     string `json:"hash,omitempty"`
}

通用 Event 结构体(修正点:新增 PartitionKey 字段)internal/event/event.go):

package event

import "time"

const (
	SourceFileUC    = "file_uc"
	SourceShareUC   = "share_uc"
	SourceRecycleUC = "recycle_uc"
	SourceUploadUC  = "upload_uc"
	SourceSystem    = "system"
)

// Event 是所有异步事件的统一消息格式。
type Event struct {
	// ID 事件唯一标识,由生产者自动生成(evt_ + 16字节随机hex)
	ID string `json:"id"`
	// IdempotencyKey 唯一消息ID(默认 "类型:事件ID"),用于消费端去重(Redis SetNX)
	IdempotencyKey string `json:"idempotency_key"`
	// PartitionKey 修正点:业务实体键(如文件ID/分享ID),专用于 Kafka 分区路由与保序。
	// 与 IdempotencyKey 解耦:前者控"顺序",后者控"去重"。
	PartitionKey string `json:"partition_key,omitempty"`
	Type         string `json:"type"`
	Source       string `json:"source"`
	Timestamp    string `json:"timestamp"`
	Payload      interface{} `json:"payload"`
}

// NewEvent 创建一个新的事件,自动填充时间戳。
func NewEvent(eventType, source string, payload interface{}) *Event {
	return &Event{
		Type:      eventType,
		Source:    source,
		Timestamp: time.Now().UTC().Format(time.RFC3339),
		Payload:   payload,
	}
}

5.2 Kafka 事件生产者(重试 + 幂等键 + 分区键 + 可靠确认)

实现思路

生产者负责把事件可靠发到 Kafka。设计要点:接口与实现分离;自动生成事件 ID 与幂等键;修正点:Kafka 消息 Key 改用 PartitionKey(保序),RequiredAcks 改为 RequireAll(不丢);Noop 兜底。

关键代码

Producer 接口与 Kafka 实现internal/event/producer.go):

package event

import (
	"cloud-disk/internal/conf"
	"context"
	"crypto/rand"
	"encoding/hex"
	"encoding/json"
	"fmt"
	"time"

	"github.com/go-kratos/kratos/v3/log"
	"github.com/segmentio/kafka-go"
)

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

type kafkaProducer struct {
	writer *kafka.Writer
}

func NewProducer(cfg *conf.Data_Kafka, writer *kafka.Writer) Producer {
	if cfg == nil || len(cfg.Addrs) == 0 {
		return noopProducer{}
	}
	return &kafkaProducer{writer: writer}
}

// Publish 序列化事件并发送到 Kafka,支持自动重试。
func (p *kafkaProducer) Publish(ctx context.Context, topic string, event *Event) error {
	// 步骤 1:唯一消息 ID + 幂等键(用于去重,与分区无关)
	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)
	if err != nil {
		return fmt.Errorf("event: failed to marshal event: %w", err)
	}

	// 修正点:Kafka Key 用 PartitionKey(保序),缺省时回退到幂等键,避免随机打散
	key := event.IdempotencyKey
	if event.PartitionKey != "" {
		key = event.PartitionKey
	}

	msg := kafka.Message{
		Topic: topic,
		Key:   []byte(key), // 同一业务实体 → 同一分区 → 单分区内有序
		Value: data,
		Headers: []kafka.Header{
			{Key: "event_type", Value: []byte(event.Type)},
			{Key: "event_source", Value: []byte(event.Source)},
		},
	}

	// 步骤 2:带退避的重试(共 3 次),且可被 ctx 中断
	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
			log.Warn("event: publish failed, retrying",
				"topic", topic, "type", event.Type, "attempt", i+1, "err", err)
			select {
			case <-ctx.Done():
				return ctx.Err()
			case <-time.After(time.Duration(100*(i+1)) * time.Millisecond):
			}
			continue
		}
		return nil
	}
	return fmt.Errorf("event: failed to publish after 3 retries: %w", lastErr)
}

func generateEventID() string {
	b := make([]byte, 16)
	if _, err := rand.Read(b); err != nil {
		return fmt.Sprintf("evt_%d", time.Now().UnixNano())
	}
	return "evt_" + hex.EncodeToString(b)
}

业务层适配器(修正点:按 payload 计算 PartitionKey 业务实体键)internal/event/biz_publisher.go):

// businessKeyOf 根据事件类型与负载推导分区键(业务实体),用于保序。
// 修正点:原实现没有这一层,导致 Kafka Key 是随机消息ID,无法按实体保序。
func businessKeyOf(eventType string, payload interface{}) string {
	switch eventType {
	case biz.EventFileDeleted, biz.EventFileMoved, biz.EventFileTrashed:
		if p, ok := payload.(*biz.FileChangedPayload); ok && len(p.FileIDs) > 0 {
			return fmt.Sprintf("file:%d", p.FileIDs[0]) // 同一文件的多事件稳定同分区
		}
	case biz.EventShareCreated, biz.EventShareCancelled:
		if p, ok := payload.(*biz.SharePayload); ok {
			return fmt.Sprintf("share:%d", p.ShareID)
		}
	case biz.EventUploadCompleted:
		if p, ok := payload.(*biz.UploadCompletedPayload); ok {
			return fmt.Sprintf("file:%d", p.FileID)
		}
	}
	return "" // 空 → 生产者回退到幂等键(顺序性不保证,但可用)
}

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)
	evt.PartitionKey = businessKeyOf(eventType, payload) // 修正点:注入分区键
	// 修正点:不再吞错,错误向上返回(或交给 Outbox 兜底,见 5.6)
	return a.producer.Publish(ctx, topic, evt)
}

Kafka Writer 初始化(修正点:RequireAll + 关闭自动建 Topic)internal/data/kafka.go):

writer := &kafka.Writer{
	Addr:                   kafka.TCP(addrs...),
	Balancer:               &kafka.Hash{},          // 按 Key 哈希分区(保序需要)
	RequiredAcks:           kafka.RequireAll,       // 修正点:多数副本确认,防丢
	AllowAutoTopicCreation: false,                  // 修正点:生产禁用自动建Topic,由运维预建
	BatchTimeout:           5 * time.Millisecond,   // 优化:批量刷新
	BatchSize:              100,                     // 优化:批量大小
	Compression:            kafka.Snappy,            // 优化:压缩降带宽
}

5.3 Kafka 事件消费者(幂等检查 + 重试退避 + 死信队列)

实现思路

类比consumeLoop 像「永不下班的值班前台」——不断取信,先查重,再办手续,办完归档;收到「下班通知」(context 取消)后,会把手头这封信处理完才走。办不成三次,就转交「问题件库房」(DLQ)。

flowchart TD
    S[启动 goroutine
consumeLoop] --> F[FetchMessage 拉取] F -->|context 取消| X[退出循环] F -->|拿到消息| D[反序列化] D -->|失败 毒丸| CK[提交并跳过] D -->|成功| I[幂等检查 SetNX 按幂等键] I -->|重复| CK I -->|首次| H[查找处理器] H -->|无处理器| CK H -->|有| R[执行+重试 可ctx中断] R -->|成功| CM[提交偏移量] R -->|耗尽| DL[发送到 DLQ] DL --> CM

消费者承担「订阅 → 拉取 → 幂等 → 重试 → 分发」全链路。设计要点:Consumer Group 订阅;阻塞式消费循环;幂等检查按幂等键去重;重试退避且可被 ctx 中断;成功后手动提交;修正点:重试耗尽发往 DLQ;Redis 故障时 fail-open。

关键代码

Consumer 接口与幂等存储接口internal/event/consumer.go):

type EventHandler func(ctx context.Context, event *Event) error

type Consumer interface {
	RegisterHandler(eventType string, handler EventHandler)
	Start(ctx context.Context) error
	Stop() error
}

// IdempotencyStore 幂等检查接口(抽象便于替换 Redis / 内存 / DB)
type IdempotencyStore interface {
	CheckAndSet(ctx context.Context, key string, ttl time.Duration) (bool, error)
}

// DLQSender 死信队列发送接口(修正点:新增)
type DLQSender interface {
	Send(ctx context.Context, event *Event, reason string) error
}

消费循环与消息处理(核心逻辑,含 DLQ 与 ctx 感知退避)

// processMessage 处理单条消息:反序列化 → 幂等 → 分发 → 重试 → (成功提交 | 发DLQ)
func (c *kafkaConsumer) processMessage(ctx context.Context, msg kafka.Message) {
	var event Event
	if err := json.Unmarshal(msg.Value, &event); err != nil {
		c.logger.Error("event: failed to unmarshal, skipping", "err", err, "key", string(msg.Key))
		c.commit(ctx, msg) // 毒丸消息:提交跳过,避免无限重试
		return
	}

	// 幂等检查(按幂等键)。Redis 故障 fail-open:记录日志但继续处理(宁可重复)
	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, continue(fail-open)", "err", err)
		}
		if !ok {
			c.logger.Debug("event: duplicate, skip", "id", event.ID)
			c.commit(ctx, msg)
			return
		}
	}

	c.mu.RLock()
	handler, ok := c.handlers[event.Type]
	c.mu.RUnlock()
	if !ok {
		c.logger.Warn("event: no handler for type", "type", event.Type)
		c.commit(ctx, msg)
		return
	}

	// 重试退避(修正点:sleep 可被 ctx 中断,停机时不空等)
	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)
			select {
			case <-ctx.Done():
				return // 停机:交给优雅退出逻辑,不强行提交
			case <-time.After(time.Duration(200*(i+1)) * time.Millisecond):
			}
			continue
		}
		c.commit(ctx, msg)
		return
	}

	// 修正点:重试耗尽 → 发往 DLQ(不再静默丢弃),再提交原消息让主队列前进
	c.logger.Error("event: handler failed after retries, to DLQ",
		"type", event.Type, "id", event.ID, "err", lastErr)
	if c.dlq != nil {
		if derr := c.dlq.Send(ctx, &event, lastErr.Error()); derr != nil {
			c.logger.Error("event: send to DLQ failed", "err", derr)
		}
	}
	c.commit(ctx, msg)
}

func (c *kafkaConsumer) commit(ctx context.Context, msg kafka.Message) {
	if err := c.reader.CommitMessages(ctx, msg); err != nil {
		c.logger.Error("event: commit failed", "err", err, "id", string(msg.Key))
	}
}

Redis 幂等存储实现internal/event/handlers.go):

// RedisIdempotencyStore 基于 Redis SetNX 的原子幂等检查。
// 返回 true = 首次(可处理);false = 重复(应跳过)。TTL 防无限膨胀。
type RedisIdempotencyStore struct{ rdb *redis.Client }

func NewRedisIdempotencyStore(rdb *redis.Client) *RedisIdempotencyStore {
	return &RedisIdempotencyStore{rdb: rdb}
}

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
}

var _ IdempotencyStore = (*RedisIdempotencyStore)(nil)

5.4 业务事件处理器(文件删除→取消分享、上传完成→异步处理)

实现思路

处理器是「事件 → 业务副作用」的桥梁:集中式 Handle 按类型 switch 分发;工厂函数便于注入依赖;文件处理器调用 shareRepo.DeleteByItemID 实现最终一致;分享/回收站/上传处理器预留扩展。修正点:文件处理器循环里单文件失败已做到"不中断其他文件",但整个处理器返回 nil 意味着"部分成功也被当成功"——若需精确可观测,应统计失败数并上报,这里保持幂等删除即可。

关键代码

事件处理器服务与分发internal/event/handlers.go):

type EventHandlerService struct {
	fileHandler    func(context.Context, *Event) error
	shareHandler   func(context.Context, *Event) error
	recycleHandler func(context.Context, *Event) error
	uploadHandler  func(context.Context, *Event) error
	mu             sync.RWMutex
}

func NewEventHandlerService(shareRepo biz.ShareRepo) *EventHandlerService {
	svc := &EventHandlerService{}
	svc.SetFileHandler(NewFileEventHandler(shareRepo))
	svc.SetShareHandler(NewShareEventHandler())
	svc.SetRecycleHandler(NewRecycleEventHandler())
	svc.SetUploadHandler(NewUploadEventHandler())
	return svc
}

// Handle 统一分发入口
func (s *EventHandlerService) Handle(ctx context.Context, event *Event) error {
	s.mu.RLock()
	defer s.mu.RUnlock()
	switch event.Type {
	case biz.EventFileDeleted, biz.EventFileMoved, biz.EventFileTrashed:
		if s.fileHandler != nil {
			return s.fileHandler(ctx, event)
		}
	case biz.EventShareCreated, biz.EventShareCancelled:
		if s.shareHandler != nil {
			return s.shareHandler(ctx, event)
		}
	case biz.EventRecycleClean:
		if s.recycleHandler != nil {
			return s.recycleHandler(ctx, event)
		}
	case biz.EventUploadCompleted:
		if s.uploadHandler != nil {
			return s.uploadHandler(ctx, event)
		}
	default:
		return fmt.Errorf("event: unknown event type: %s", event.Type)
	}
	return nil
}

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)
}

// parsePayload 通过 JSON 往返将 interface{} 负载转为强类型
func parsePayload(payload interface{}, out interface{}) error {
	data, err := json.Marshal(payload)
	if err != nil {
		return err
	}
	return json.Unmarshal(data, out)
}

文件事件处理器(核心:删除/移回收站 → 取消相关分享)

// NewFileEventHandler 文件删除/移回收站时异步取消相关分享(最终一致)。
func NewFileEventHandler(shareRepo biz.ShareRepo) func(context.Context, *Event) error {
	return func(ctx context.Context, event *Event) error {
		var payload biz.FileChangedPayload
		if err := parsePayload(event.Payload, &payload); err != nil {
			log.Warn("event: parse file payload failed", "err", err)
			return nil // 负载解析失败不重试(非瞬态)
		}
		// DeleteByItemID 本身幂等:重复取消无害,配合幂等键保证副作用只生效一次
		for _, id := range payload.FileIDs {
			if err := shareRepo.DeleteByItemID(ctx, "file", id); err != nil {
				log.Error("event: cancel share for file failed", "file_id", id, "err", err)
			}
		}
		for _, id := range payload.FolderIDs {
			if err := shareRepo.DeleteByItemID(ctx, "folder", id); err != nil {
				log.Error("event: cancel share for folder failed", "folder_id", id, "err", err)
			}
		}
		log.Info("event: file event handled", "type", event.Type, "id", event.ID,
			"files", len(payload.FileIDs), "folders", len(payload.FolderIDs))
		return nil
	}
}

func NewShareEventHandler() func(context.Context, *Event) error {
	return func(ctx context.Context, event *Event) error {
		log.Info("event: share event handled", "type", event.Type, "id", event.ID)
		return nil
	}
}
func NewRecycleEventHandler() func(context.Context, *Event) error {
	return func(ctx context.Context, event *Event) error {
		log.Info("event: recycle event handled", "type", event.Type, "id", event.ID)
		return nil
	}
}
func NewUploadEventHandler() func(context.Context, *Event) error {
	return func(ctx context.Context, event *Event) error {
		log.Info("event: upload completed handled", "type", event.Type, "id", event.ID)
		return nil
	}
}

5.5 ConsumerServer 生命周期管理(实现 transport.Server 接口)

实现思路

Kratos 通过 transport.Server 接口统一管理 HTTP/gRPC/自定义服务生命周期。把消费者包装为 ConsumerServer 后,应用启停自动启停消费。要点:实现 Start/StopStartRegisterAllHandlers 并注入幂等存储(修正点:同时注入 DLQ 发送器);Wire 装配。

关键代码

ConsumerServerinternal/event/consumer_server.go):

var _ transport.Server = (*ConsumerServer)(nil)

type ConsumerServer struct {
	consumer       Consumer
	handlerService *EventHandlerService
	idempotency    IdempotencyStore
	dlq            DLQSender // 修正点:注入死信发送器
	topics         []string
}

func NewConsumerServer(consumer Consumer, handlerService *EventHandlerService,
	idempotency IdempotencyStore, dlq DLQSender, topics []string) *ConsumerServer {
	return &ConsumerServer{consumer: consumer, handlerService: handlerService,
		idempotency: idempotency, dlq: dlq, topics: topics}
}

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)
			c.dlq = s.dlq // 修正点:注入 DLQ
		}
	}
	if err := s.consumer.Start(ctx); err != nil {
		log.Error("event: failed to start consumer", "err", err)
		return err
	}
	log.Info("event: consumer server started", "topics", s.topics)
	return nil
}

func (s *ConsumerServer) Stop(ctx context.Context) error {
	if err := s.consumer.Stop(); err != nil {
		log.Error("event: failed to stop consumer", "err", err)
		return err
	}
	log.Info("event: consumer server stopped")
	return nil
}

Topic 配置与消费者组internal/event/kafka_config.go):

const defaultConsumerGroup = "cloud-disk-consumer"

func KafkaConsumerGroup() string { return defaultConsumerGroup }

func KafkaTopics(cfg *conf.Data_Kafka) []string {
	if cfg == nil || cfg.Topics == nil {
		return []string{"file-events", "share-events", "recycle-events", "upload-events"}
	}
	prefix := cfg.TopicPrefix
	if prefix == "" {
		prefix = "cloud_disk_"
	}
	return []string{
		prefix + cfg.Topics.FileEvents,
		prefix + cfg.Topics.ShareEvents,
		prefix + cfg.Topics.RecycleEvents,
		prefix + cfg.Topics.UploadEvents,
	}
}

Wire 装配internal/event/biz_publisher.go):

var ProviderSet = wire.NewSet(
	NewProducer,
	NewConsumer,
	NewEventHandlerService,
	NewRedisIdempotencyStore,
	NewEventPublisherAdapter,
	NewConsumerServer,
	KafkaConsumerGroup,
	KafkaTopics,
	wire.Bind(new(IdempotencyStore), new(*RedisIdempotencyStore)),
	// 修正点:绑定 DLQ 发送器
	wire.Bind(new(DLQSender), new(*KafkaDLQSender)),
)

Noop 空实现兜底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
}

5.6 事务发件箱(Outbox):根治"事件静默丢失"(商用补强,重点)

为什么需要

本模块原有最大商用缺口:业务在 DB 事务提交之后命令式 Publish(且调用点 _ = Publish(...) 吞错)。这有两个致命后果:

  1. 进程在「DB 提交成功 → Publish 执行前」崩溃,事件永远丢失,导致文件已删但分享没取消——业务不一致且无任何报错
  2. Kafka 临时不可用 / 3 次重试仍失败,错误被 _ = 吞掉,事件静默消失。

类比:发件箱像「先写纸条塞进自家信箱(和业务一起落库),再由邮差定时来取寄出」。纸条和业务单据写在同一本账里(同一事务),要么都成、要么都不成,绝不会出现"单据已记、信却没寄"的半吊子状态。

flowchart TD
    A[业务事务开始] --> B[写业务表]
    B --> C[同一事务写 outbox 表 事件未发送]
    C --> D[提交事务]
    D --> E[事务成功 事件已落库]
    E --> F[中继 goroutine 读 outbox]
    F --> G[发往 Kafka]
    G -->|成功| H[标记 outbox 已发送]
    G -->|失败| F2[退避后重试]
    F2 --> G

关键代码

Outbox 表模型与仓储接口internal/data/outbox.go):

// Outbox 事件发件箱表:与业务在同一 DB 事务写入,保证原子性。
type Outbox struct {
	ID        uint64 `gorm:"primaryKey"`
	EventType string `gorm:"index"`
	Payload   string // JSON 序列化后的事件(含 PartitionKey)
	Status    int    // 0=待发送 1=已发送
	CreatedAt time.Time
}

// OutboxRepo 发件箱仓储
type OutboxRepo interface {
	// Save 在业务事务内调用,把事件作为"待发送"写入
	Save(ctx context.Context, eventType string, payload []byte) error
	// Pending 取出待发送事件(中继者轮询)
	Pending(ctx context.Context, limit int) ([]*Outbox, error)
	// MarkSent 标记已发送
	MarkSent(ctx context.Context, id uint64) error
}

业务层:事务内写 outbox(以分享创建为例)

// 原错误写法(已修正):_ = uc.eventPublisher.Publish(...)  // 吞错 + 非原子
//
// 修正点:改为在 DB 事务内写入 outbox,由中继者异步投递,保证"业务落库=事件落库"
func (uc *ShareUsecase) CreateShare(ctx context.Context, userID uint64, itemType string, itemID uint64, ...) (...) {
	// ... 构造 share 并 shareRepo.Create(事务内)...
	// 修正点:同一事务内落 outbox
	payload, _ := json.Marshal(&SharePayload{UserID: userID, ShareID: created.ID, UUID: created.UUID, Action: "created"})
	if err := uc.outbox.Save(ctx, EventShareCreated, payload); err != nil {
		return nil, "", err // 落 outbox 失败 = 整笔事务回滚,绝不会"分享了却没发事件"
	}
	return created, url, nil
}

中继者:读 outbox → 发 Kafka → 标记已发

// OutboxRelay 后台 goroutine:轮询 outbox,投递到 Kafka,成功则标记已发。
func (r *outboxRelay) Run(ctx context.Context) {
	ticker := time.NewTicker(500 * time.Millisecond)
	defer ticker.Stop()
	for {
		select {
		case <-ctx.Done():
			return
		case <-ticker.C:
			batch, err := r.outbox.Pending(ctx, 100)
			if err != nil || len(batch) == 0 {
				continue
			}
			for _, row := range batch {
				// 构造事件并发布(带分区键)
				evt, _ := rebuildEvent(row) // 从 payload 还原 Event + PartitionKey
				if err := r.producer.Publish(ctx, r.topicOf(row.EventType), evt); err != nil {
					log.Error("outbox: publish failed", "id", row.ID, "err", err)
					continue // 不标记,下次轮询重试
				}
				_ = r.outbox.MarkSent(ctx, row.ID)
			}
		}
	}
}

采用 Outbox 后,原调用点的 _ = Publish(...) 改为「写 outbox」,Kafka 短暂不可用也不会丢事件——事件在 outbox 表里等待中继者重试,直到成功投递。这是商用系统"业务 + 事件"原子性的标准答案。


5.7 可观测性:Lag / 重试 / DLQ 监控(商用补强)

实现思路

本模块原仅打日志,生产必须量化。建议在消费循环与 DLQ 发送处埋点,暴露给 Prometheus:

  • event_consume_lag{topic}:消费者滞后量(可由 kafka-goReadLag/Stat 获取,或外部 exporter)。
  • event_handler_retries_total{type}:重试次数累计。
  • event_dlq_total{type}:进入死信队列的消息数(核心告警指标)。
  • event_process_duration_seconds{type}:处理耗时直方图(P99)。
flowchart LR
    C[消费循环] -->|lag/耗时| P[Prometheus 指标]
    R[重试] -->|计数| P
    D[DLQ 发送] -->|计数 告警| P
    P --> G[Grafana 看板 + 告警]
// 示例:在重试耗尽发 DLQ 处计数(伪代码)
eventDLQTotal.WithLabelValues(event.Type).Inc()
// 在 processMessage 入口记录耗时
defer func(start time.Time) {
	eventProcessDuration.WithLabelValues(event.Type).Observe(time.Since(start).Seconds())
}(time.Now())

商用告警建议:DLQ 长度 > 0 即告警(失败消息不应长期堆积);Lag 持续上涨触发扩容复核;P99 处理耗时突增排查下游(如 DeleteByItemID 慢查询)。


自测题与动手练习

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

  1. 本项目的幂等键 IdempotencyKey = "类型:随机事件ID",它到底在防什么?它防不住什么?真正的"业务级 Exactly-Once"靠什么兜底?
  2. 原教程说"以幂等键作 Kafka Key 可保证同一文件事件保序",这句话错在哪?正确做法应该用什么当 Kafka Key?为什么?
  3. 调用点写 _ = uc.eventPublisher.Publish(...) 有什么商用风险?进程在"DB 提交后、Publish 前"崩溃会发生什么?事务发件箱(Outbox)如何解决?
  4. 消费者重试 3 次仍失败,本项目原实现会怎样?正确的商用做法是什么(DLQ)?DLQ 里的消息后续怎么处理?
  5. RequiredAcks = RequireOneRequireAll 在可靠性上差在哪?什么场景会丢消息?
  6. 消费端幂等检查在 Redis 故障时应该 fail-open 还是 fail-closed?为什么和用户模块"黑名单 fail-closed"相反?
  7. 监控看到消费 Lag 持续上涨,你作为值班 SRE 第一步排查什么?扩容为什么不能超过分区数?
  8. 幂等键 TTL 设 24 小时,同一消息在 25 小时后被重投会怎样?对幂等/非幂等处理器分别有什么影响?

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

  1. Event 加上 PartitionKey,在适配器里为 file/share/upload 三类事件计算业务键;起一个本地 Kafka,用两个分区观察"同文件两条事件"是否落到同一分区(看 msg.Key 哈希)。
  2. processMessage 补上"重试耗尽 → 发往 DLQ Topic"的逻辑,起一个消费者故意让 handler 返回错误,观察 DLQ 里出现对应消息。
  3. 实现最小版 Outbox:建一张 outbox 表,业务事务内写事件,另起一个 goroutine 每 500ms 轮询并发 Kafka,验证"Kafka 宕机期间事件不丢、恢复后自动补发"。
  4. 用 Redis SetNX 跑一次幂等检查,观察重复键返回 false;再把 Redis 停掉,确认消费端 fail-open 仍继续处理(并理解为何这样是安全的)。

本章小结

  • 事件驱动模块用 Kafka 承载异步事件,实现"主流程"与"副作用"(如删文件→取消分享)解耦,靠最终一致性提升响应速度。
  • 认知纠偏(最重要):幂等键(去重)与分区键(保序)是两件事。本项目原实现把随机消息 ID 同时当两者用,导致"同文件事件保序"的承诺不成立;商用修正为 PartitionKey(业务实体)作 Kafka Key 控顺序、IdempotencyKey(唯一消息 ID)作 Redis 键控去重。
  • 可靠性三重补强:① 生产端 RequiredAll + 多副本,防 Broker 丢消息;② 消费端 At-Least-Once + 幂等处理器,重复投递安全;③ 事务发件箱(Outbox) 根治"吞错 + 非原子发布"导致的事件静默丢失。
  • 失败闭环:消费端退避重试(可被 ctx 中断),重试耗尽转入 DLQ 死信队列,主队列不阻塞、失败可查可重投。
  • 可观测:Lag / 重试次数 / DLQ 长度 / P99 耗时须接 Prometheus + Grafana 并告警,这是生产上线的必要能力。
  • 下一篇可进入「可观测性模块」,看如何把 Lag、重试、DLQ 指标统一接入监控与告警体系。
About Me

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

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

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

目标

学AI,加油!加油!