以「用户通知中心」为例,完整拆解一个中型 Go 服务从需求分析、架构设计到代码落地的全过程。
一、案例背景:设计一个用户通知中心服务
在多数互联网系统里,通知能力并不是最显眼的业务模块,却几乎贯穿所有核心链路。用户注册成功后需要欢迎消息,订单状态变化后需要短信提醒,审批流结束后需要邮件通知,App 活动上线后需要 Push 触达。随着业务扩展,通知类型、发送渠道、发送策略和失败处理逻辑会迅速膨胀。
如果一开始只是简单写几个 SendEmail()、SendSMS() 函数,短期内可以工作,但当业务逐渐复杂后,系统通常会暴露出以下问题:
- 渠道调用逻辑散落在各业务服务中,难以统一治理
- 模板、重试、限流、状态追踪等能力重复实现
- 多渠道发送策略缺少统一编排,修改成本高
- 渠道供应商切换困难,强耦合严重
- 无法准确追踪一条通知从创建到最终送达的全过程
因此,我们设计一个中型服务规模的「用户通知中心」:
- 支持多渠道:站内信、邮件、短信、Push
- 支持同步下发与异步发送
- 支持模板渲染、路由分发、状态追踪
- 支持失败重试、死信处理、幂等控制
- 支持后续扩展更多渠道与业务场景
这个服务本质上是一个典型的中台型基础能力服务。它对外承接业务系统的通知需求,对内屏蔽不同渠道的复杂性,并通过统一模型实现扩展性和可运维性。
二、需求分析与系统边界划分
在正式设计架构之前,先把需求拆清楚。架构设计最怕的问题不是“技术选型不够高级”,而是系统边界模糊,结果导致职责混乱。
2.1 核心功能需求
用户通知中心需要提供以下核心能力:
-
统一通知接入
- 业务方通过统一 API 提交通知请求
- 支持指定单渠道或多渠道发送
- 支持立即发送、定时发送、事件触发发送
-
通知内容管理
- 支持通知模板
- 支持变量渲染
- 支持不同渠道的内容格式适配
-
消息路由能力
- 根据通知类型、用户偏好、业务优先级决定发送渠道
- 支持降级策略,例如 Push 失败后转短信
-
发送执行与结果回写
- 统一调用不同渠道服务商
- 记录发送状态:待发送、发送中、成功、失败、部分成功
- 支持供应商回执更新
-
可靠性保障
- 幂等控制
- 重试机制
- 死信队列处理
- 限流与熔断
-
运营与审计
- 通知查询
- 发送明细查看
- 异常告警
- 审计日志
2.2 非功能需求
除了功能需求,中型服务更需要明确非功能目标:
- 高可用:单节点异常不能影响整体发送
- 可扩展:新增渠道时不修改主流程核心逻辑
- 可观测:能追踪通知生命周期和失败原因
- 高吞吐:削峰填谷,支撑突发消息洪峰
- 一致性可控:发送任务创建与投递流程需要尽量可靠
- 低耦合:业务系统不应感知具体渠道供应商细节
2.3 系统边界划分
通知中心不是“包办一切”的超级系统,必须明确边界。
通知中心负责:
- 接收通知请求
- 存储通知任务与发送记录
- 路由决策与策略编排
- 异步投递与重试
- 调用渠道网关
- 追踪发送状态
- 对外提供查询接口
通知中心不负责:
- 用户主数据管理
- 邮件/短信供应商账号管理后台
- App Push SDK 客户端实现
- 复杂营销编排系统
- 上游业务事件定义本身
换句话说,通知中心是发送能力中台,而不是用户中心、营销平台或内容 CMS。
2.4 外部依赖关系
为了让边界更清晰,可以把上下游关系概括如下:
| 角色 | 示例 | 与通知中心关系 |
|---|---|---|
| 上游业务系统 | 订单系统、审批系统、营销系统 | 发起通知请求或投递业务事件 |
| 用户中心 | 用户资料、联系方式、偏好设置 | 提供接收人信息与订阅偏好 |
| 模板服务 | 模板存储与版本管理 | 提供模板内容与变量定义 |
| MQ | Kafka / RocketMQ / Pulsar | 承担异步解耦与削峰 |
| 渠道供应商 | SMTP、短信网关、Push 平台 | 执行实际发送 |
| 可观测平台 | 日志、指标、链路追踪、告警系统 | 监控通知中心运行状态 |
三、架构方案设计:分层架构 + 事件驱动
对于中型通知服务,最适合的方案通常不是“纯同步 RPC”,也不是“完全事件化的黑盒系统”,而是采用分层架构 + 事件驱动的组合设计。
- 分层架构负责保持代码组织清晰、职责明确
- 事件驱动负责解耦提交与发送,提高吞吐和可靠性
3.1 为什么采用分层架构
分层架构适合通知中心这类“业务规则明确、模块职责清楚、后续演进频繁”的系统。它能帮助我们把系统划分为:
- 接入层:面向 API/HTTP/gRPC
- 应用层:编排用例流程
- 领域层:承载通知模型和规则
- 基础设施层:数据库、缓存、MQ、外部渠道适配
这样做的好处是:
- 核心业务规则不会被外部依赖污染
- 外部系统变动时,影响被隔离在基础设施层
- 单测更容易编写
- 新增渠道、缓存或消息队列时改动范围可控
3.2 为什么采用事件驱动
通知发送天然具备异步特征:
- 业务方通常只关心“提交成功”,不愿等待所有渠道发完
- 短信、邮件、Push 的外部依赖延迟不稳定
- 高峰期通知请求量可能短时暴涨
- 部分通知允许稍后发送,但不能丢失
因此,一个合理的流程是:
- API 接收通知请求
- 落库生成通知任务
- 发布“通知已创建”事件
- 异步消费者执行路由、渲染、发送
- 状态回写并视情况继续派发重试事件
这样能将“请求接入”和“实际发送”解耦,避免上游业务被下游渠道拖慢。
3.3 整体架构图示例
下面给出一个适合中型服务的整体架构示意图:
flowchart LR
A[业务系统] --> B[通知中心 API]
B --> C[应用层 NotificationAppService]
C --> D[(MySQL 通知任务库)]
C --> E[Outbox/Event Publisher]
E --> F[(MQ)]
F --> G[发送工作节点 SendWorker]
G --> H[消息路由 Router]
H --> I[模板渲染 Renderer]
I --> J[渠道适配器 Adapters]
J --> K[站内信]
J --> L[邮件供应商]
J --> M[短信供应商]
J --> N[Push 平台]
G --> O[(Redis)]
G --> D
J --> P[回执/状态更新]
P --> D
3.4 分层职责划分
| 层级 | 核心职责 | 不应该做的事 |
|---|---|---|
| Interface 层 | 处理 HTTP/gRPC 请求、参数校验、返回响应 | 编写复杂业务规则 |
| Application 层 | 编排通知创建、发送、查询等用例 | 直接依赖具体第三方 SDK |
| Domain 层 | 定义通知实体、路由规则、状态流转、接口契约 | 关心数据库表结构细节 |
| Infrastructure 层 | 实现仓储、缓存、MQ、渠道 SDK 调用 | 修改领域规则 |
3.5 核心流程说明
请求创建流程
- 上游业务调用通知中心 API
- Interface 层完成参数校验
- Application 层创建 Notification 聚合与 Delivery 记录
- 仓储层持久化数据
- 发布通知创建事件
- 返回请求已受理结果
异步发送流程
- Worker 消费通知创建事件
- 查询通知详情与接收人信息
- 根据策略进行消息路由
- 渲染渠道内容
- 调用各渠道适配器发送
- 回写渠道发送结果
- 失败场景进入延迟重试或死信处理
状态追踪流程
- 初始状态为
PENDING - 发送中更新为
SENDING - 单渠道成功则写入子记录为
SUCCESS - 单渠道失败则写入
FAILED - 如果通知是多渠道,可根据整体规则汇总为
PARTIAL_SUCCESS或SUCCESS - 供应商异步回执到达后再更新最终投递状态
四、核心模块设计
4.1 消息路由模块
消息路由模块的目标不是“把消息发出去”,而是决定:发给谁、走什么渠道、是否需要降级、优先级如何排序。
路由输入通常包括:
- 通知类型:营销、交易、系统、安全
- 用户偏好:是否接收短信、是否允许夜间 Push
- 业务优先级:高优先级失败是否转人工或转备用渠道
- 渠道能力:某渠道是否可用、是否在限流状态
- 内容特征:超长文本是否适合短信
一个简单但实用的路由规则可以是:
| 场景 | 首选渠道 | 备用渠道 | 说明 |
|---|---|---|---|
| 安全验证码 | 短信 | Push | 强实时,高优先级 |
| 订单通知 | 站内信 + Push | 邮件 | 兼顾到达率与成本 |
| 营销活动 | Push | 站内信 | 控制短信成本 |
| 审批结果 | 站内信 + 邮件 | Push | 保证留痕 |
下面是一个简化的 Go 路由实现:
package domain
type Channel string
const (
ChannelInbox Channel = "INBOX"
ChannelEmail Channel = "EMAIL"
ChannelSMS Channel = "SMS"
ChannelPush Channel = "PUSH"
)
type NotificationType string
const (
TypeMarketing NotificationType = "MARKETING"
TypeOrder NotificationType = "ORDER"
TypeSecurity NotificationType = "SECURITY"
TypeApproval NotificationType = "APPROVAL"
)
type UserPreference struct {
AllowEmail bool
AllowSMS bool
AllowPush bool
AllowInbox bool
}
type RouteInput struct {
Type NotificationType
Preference UserPreference
IsHighRisk bool
}
type Router interface {
Resolve(input RouteInput) []Channel
}
type DefaultRouter struct{}
func (r DefaultRouter) Resolve(input RouteInput) []Channel {
switch input.Type {
case TypeSecurity:
channels := make([]Channel, 0, 2)
if input.Preference.AllowSMS {
channels = append(channels, ChannelSMS)
}
if input.Preference.AllowPush {
channels = append(channels, ChannelPush)
}
return channels
case TypeApproval:
channels := make([]Channel, 0, 3)
if input.Preference.AllowInbox {
channels = append(channels, ChannelInbox)
}
if input.Preference.AllowEmail {
channels = append(channels, ChannelEmail)
}
if input.Preference.AllowPush && input.IsHighRisk {
channels = append(channels, ChannelPush)
}
return channels
default:
channels := make([]Channel, 0, 2)
if input.Preference.AllowInbox {
channels = append(channels, ChannelInbox)
}
if input.Preference.AllowPush {
channels = append(channels, ChannelPush)
}
return channels
}
}
在工程实践中,路由模块往往不会只靠 switch-case。更常见的方式有:
- 规则配置化:把通知类型与渠道映射放入配置中心
- 策略模式:不同通知类型对应不同路由策略实现
- 规则引擎:适用于复杂条件判断较多的场景
对于一个中型服务,建议从“策略模式 + 配置化规则”起步,避免过早引入重量级规则引擎。
4.2 渠道适配器模块
渠道适配器的职责是屏蔽不同供应商差异,让应用层只面向统一接口编程。
典型差异包括:
- 请求协议不同:HTTP、SMTP、私有 SDK
- 鉴权方式不同:Token、签名、用户名密码
- 返回格式不同:同步响应、异步回执、状态轮询
- 能力差异不同:模板短信、批量发送、限速策略
因此,需要定义统一的渠道接口:
package domain
import "context"
type SendRequest struct {
NotificationID string
DeliveryID string
Target string
Title string
Content string
TemplateCode string
Metadata map[string]string
}
type SendResult struct {
ProviderMessageID string
Success bool
Code string
Message string
}
type Sender interface {
Channel() Channel
Send(ctx context.Context, req SendRequest) (SendResult, error)
}
然后每种渠道各自实现:
package adapter
import (
"context"
"fmt"
"notify/internal/domain"
)
type EmailSender struct {
endpoint string
apiKey string
}
func NewEmailSender(endpoint, apiKey string) *EmailSender {
return &EmailSender{endpoint: endpoint, apiKey: apiKey}
}
func (s *EmailSender) Channel() domain.Channel {
return domain.ChannelEmail
}
func (s *EmailSender) Send(ctx context.Context, req domain.SendRequest) (domain.SendResult, error) {
// 此处省略真实 HTTP 调用,演示统一适配模式
if req.Target == "" {
return domain.SendResult{}, fmt.Errorf("empty email target")
}
return domain.SendResult{
ProviderMessageID: "mail_123456",
Success: true,
Code: "OK",
Message: "email accepted",
}, nil
}
通过这种方式,应用层不需要关心是 SMTP 还是第三方邮件 API,只需要拿到 Sender 接口即可。
还可以继续演进:
- 一个渠道支持多个供应商,实现主备切换
- 通过工厂模式或注册表动态获取适配器
- 为每个适配器增加超时、重试、熔断包装器
例如适配器注册表:
package application
import "notify/internal/domain"
type SenderRegistry struct {
senders map[domain.Channel]domain.Sender
}
func NewSenderRegistry(items ...domain.Sender) *SenderRegistry {
m := make(map[domain.Channel]domain.Sender, len(items))
for _, item := range items {
m[item.Channel()] = item
}
return &SenderRegistry{senders: m}
}
func (r *SenderRegistry) Get(channel domain.Channel) (domain.Sender, bool) {
sender, ok := r.senders[channel]
return sender, ok
}
4.3 发送状态追踪模块
通知中心如果没有完善的状态追踪,排查问题会非常痛苦。我们不能只记录“发了没发”,而是要跟踪:
- 通知任务是否已创建
- 哪些渠道被选中
- 每个渠道当前状态如何
- 外部供应商返回了什么结果
- 是否还会继续重试
- 最终整体结果是什么
建议把状态设计成两层:
- 通知主任务状态:描述整条通知的整体进度
- 渠道投递状态:描述每个渠道实例的发送结果
推荐状态机如下:
| 维度 | 状态 | 含义 |
|---|---|---|
| Notification | PENDING | 已创建,待处理 |
| Notification | PROCESSING | 正在路由或发送 |
| Notification | SUCCESS | 全部达到成功条件 |
| Notification | PARTIAL_SUCCESS | 部分渠道成功 |
| Notification | FAILED | 最终失败 |
| Delivery | PENDING | 渠道任务待发送 |
| Delivery | SENDING | 已发起调用 |
| Delivery | SUCCESS | 渠道成功受理或送达 |
| Delivery | FAILED | 渠道发送失败 |
| Delivery | RETRYING | 等待下一次重试 |
| Delivery | DEAD | 超过重试上限,进入死信 |
下面给出一个简化的状态流转实现:
package domain
import "errors"
type Status string
const (
StatusPending Status = "PENDING"
StatusProcessing Status = "PROCESSING"
StatusSuccess Status = "SUCCESS"
StatusPartialSuccess Status = "PARTIAL_SUCCESS"
StatusFailed Status = "FAILED"
StatusRetrying Status = "RETRYING"
StatusDead Status = "DEAD"
)
type Delivery struct {
ID string
Channel Channel
Status Status
RetryCount int
MaxRetry int
FailureCode string
FailureMsg string
}
func (d *Delivery) MarkSending() error {
if d.Status != StatusPending && d.Status != StatusRetrying {
return errors.New("invalid state transition to sending")
}
d.Status = StatusProcessing
return nil
}
func (d *Delivery) MarkSuccess() {
d.Status = StatusSuccess
}
func (d *Delivery) MarkFailure(code, msg string) {
d.RetryCount++
d.FailureCode = code
d.FailureMsg = msg
if d.RetryCount > d.MaxRetry {
d.Status = StatusDead
return
}
d.Status = StatusRetrying
}
在真实项目里,状态追踪通常还会配合:
- 状态变更流水表
- 审计日志
- 回执事件表
- 幂等键与版本号控制
五、数据库与缓存设计
通知中心不是单纯的“发消息工具”,而是一个典型的数据驱动系统。通知任务、投递记录、模板快照、回执状态都需要存储。
5.1 数据库选型思路
对中型服务来说,常见做法是:
- MySQL:存储核心业务数据,保证结构化查询与事务能力
- Redis:承载缓存、幂等键、短期限流计数、热点状态
- MQ:负责异步削峰与事件投递
如果后期通知明细量级极大,也可以把历史明细归档到 ClickHouse、ES 或对象存储中,但中期阶段先用 MySQL 足够稳妥。
5.2 关键表设计
建议至少有以下几张核心表:
| 表名 | 用途 |
|---|---|
notifications |
通知主任务 |
notification_deliveries |
渠道投递记录 |
notification_events |
状态变更事件或发送流水 |
notification_templates |
模板定义 |
outbox_events |
本地消息表,保证事件可靠投递 |
notifications 表
CREATE TABLE notifications (
id BIGINT PRIMARY KEY AUTO_INCREMENT,
biz_id VARCHAR(64) NOT NULL,
notify_type VARCHAR(32) NOT NULL,
user_id BIGINT NOT NULL,
title VARCHAR(255) NOT NULL,
content TEXT NOT NULL,
status VARCHAR(32) NOT NULL,
priority INT NOT NULL DEFAULT 0,
schedule_time DATETIME NULL,
idempotent_key VARCHAR(128) NOT NULL,
created_at DATETIME NOT NULL,
updated_at DATETIME NOT NULL,
UNIQUE KEY uk_idempotent_key (idempotent_key),
KEY idx_user_id_created_at (user_id, created_at),
KEY idx_status_created_at (status, created_at)
);
notification_deliveries 表
CREATE TABLE notification_deliveries (
id BIGINT PRIMARY KEY AUTO_INCREMENT,
notification_id BIGINT NOT NULL,
channel VARCHAR(32) NOT NULL,
target VARCHAR(255) NOT NULL,
status VARCHAR(32) NOT NULL,
provider VARCHAR(64) NOT NULL,
provider_message_id VARCHAR(128) DEFAULT '',
retry_count INT NOT NULL DEFAULT 0,
max_retry INT NOT NULL DEFAULT 3,
failure_code VARCHAR(64) DEFAULT '',
failure_message VARCHAR(255) DEFAULT '',
next_retry_time DATETIME NULL,
created_at DATETIME NOT NULL,
updated_at DATETIME NOT NULL,
KEY idx_notification_id (notification_id),
KEY idx_status_next_retry_time (status, next_retry_time)
);
outbox_events 表
CREATE TABLE outbox_events (
id BIGINT PRIMARY KEY AUTO_INCREMENT,
aggregate_type VARCHAR(64) NOT NULL,
aggregate_id VARCHAR(64) NOT NULL,
event_type VARCHAR(64) NOT NULL,
payload JSON NOT NULL,
status VARCHAR(32) NOT NULL,
created_at DATETIME NOT NULL,
updated_at DATETIME NOT NULL,
KEY idx_status_created_at (status, created_at)
);
5.3 为什么需要 Outbox 模式
很多团队在实现异步架构时会遇到一个经典问题:
- 数据库写成功了
- MQ 发送失败了
- 结果通知任务存在,但发送事件没发出去
这会导致“数据存在但流程中断”。
Outbox 模式的做法是:
- 在同一个本地事务里写入业务表和
outbox_events - 后台投递器轮询未发送事件
- 成功投递到 MQ 后再更新 outbox 状态
这样可以显著降低“本地事务与消息投递不一致”的概率。
5.4 Redis 的使用场景
Redis 在通知中心里非常适合做“高频、短期、可丢失或可重建”的数据:
| 场景 | Key 示例 | 说明 |
|---|---|---|
| 幂等控制 | notify:idempotent:{key} |
防止重复创建通知 |
| 渠道限流 | notify:rate:sms:{minute} |
控制供应商调用频率 |
| 用户发送频控 | notify:user:{uid}:sms |
防止单用户被轰炸 |
| 热点模板缓存 | notify:template:{code} |
降低模板读取压力 |
| 短期状态缓存 | notify:delivery:{id} |
加速后台查询 |
但要注意:Redis 不应该替代核心状态持久化。 所有关键发送结果最终都必须落到数据库中。
六、高可用与容错设计
中型服务上线后,真正拉开质量差距的不是“能不能发出去”,而是“出问题时系统是否还能稳定工作”。
6.1 高可用设计要点
1)服务无状态化
通知 API 节点与 Worker 节点都尽量设计成无状态实例,核心状态放在 MySQL、Redis、MQ 中。这样可以实现:
- 多实例水平扩容
- 容器重建快速恢复
- 故障节点摘除简单
2)读写分离与连接池治理
数据库层面需要控制:
- 合理的连接池大小
- 慢 SQL 优化
- 按状态和时间建索引
- 热表归档策略
3)消费者集群化
发送 Worker 需要以消费者组方式部署,通过 MQ 分区或消费组实现并发处理。这样能提升吞吐,也能保证单实例宕机不导致任务停摆。
6.2 常见故障与应对策略
| 故障场景 | 表现 | 应对方案 |
|---|---|---|
| 短信供应商超时 | 请求大量积压 | 超时控制 + 熔断 + 切备用供应商 |
| Push 平台偶发失败 | 部分消息失败 | 重试队列 + 指数退避 |
| MQ 短暂不可用 | 事件积压 | Outbox 补偿投递 |
| Redis 故障 | 幂等与缓存能力下降 | 降级到 DB 校验,允许局部性能下降 |
| 数据库主库抖动 | 写入失败 | 失败快速返回 + 重试 + 告警 |
| 回执延迟 | 状态更新变慢 | 允许最终一致,前台展示“处理中” |
6.3 重试机制设计
通知系统的重试不能粗暴地“无限重发”,否则很容易造成雪崩。
建议策略:
- 区分错误类型:
- 参数错误:不重试
- 供应商限流:延迟重试
- 网络抖动:快速重试 + 指数退避
- 鉴权失败:直接告警并人工介入
- 设置最大重试次数
- 重试任务进入独立延迟队列
- 超过重试上限进入死信表或死信队列
简单的指数退避计算示例:
package retry
import "time"
func NextBackoff(retryCount int) time.Duration {
if retryCount <= 0 {
return 5 * time.Second
}
backoff := time.Duration(1<<uint(retryCount)) * 5 * time.Second
maxBackoff := 10 * time.Minute
if backoff > maxBackoff {
return maxBackoff
}
return backoff
}
6.4 幂等与去重设计
通知中心很容易因为上游重试、网络抖动、消息重复投递而出现重复发送。因此需要在多个层面做幂等控制:
- 请求幂等:基于业务幂等键避免重复创建通知
- 消费幂等:同一事件重复消费时保证不重复发出
- 渠道幂等:向供应商发送时尽量带业务请求号
一个常见做法是:
- 上游传入
biz_id + scene + receiver组合幂等键 - 创建通知时先查 Redis,再落库唯一索引校验
- Worker 处理时使用投递记录状态判断是否已发送
七、代码结构组织与关键实现片段
对于 Go 项目,中型服务的代码结构应该既能体现业务分层,又能方便多人协作。下面是一种比较常见且实用的目录组织方式:
notify-center/
├── cmd/
│ └── server/
│ └── main.go
├── internal/
│ ├── application/
│ │ ├── command/
│ │ │ └── create_notification.go
│ │ ├── query/
│ │ │ └── get_notification.go
│ │ ├── service/
│ │ │ └── notification_app_service.go
│ │ └── registry/
│ │ └── sender_registry.go
│ ├── domain/
│ │ ├── entity/
│ │ │ ├── notification.go
│ │ │ └── delivery.go
│ │ ├── repository/
│ │ │ ├── notification_repository.go
│ │ │ └── delivery_repository.go
│ │ ├── service/
│ │ │ ├── router.go
│ │ │ └── renderer.go
│ │ └── event/
│ │ └── notification_created.go
│ ├── infrastructure/
│ │ ├── persistence/
│ │ │ ├── mysql/
│ │ │ └── redis/
│ │ ├── mq/
│ │ │ ├── producer.go
│ │ │ └── consumer.go
│ │ ├── adapter/
│ │ │ ├── email_sender.go
│ │ │ ├── sms_sender.go
│ │ │ ├── push_sender.go
│ │ │ └── inbox_sender.go
│ │ └── scheduler/
│ │ └── retry_dispatcher.go
│ ├── interfaces/
│ │ ├── http/
│ │ │ ├── handler.go
│ │ │ └── dto.go
│ │ └── worker/
│ │ └── send_worker.go
│ └── pkg/
│ ├── logger/
│ ├── xerror/
│ └── trace/
└── configs/
└── config.yaml
这个结构的重点不是“目录长得多漂亮”,而是让职责边界清晰:
application负责用例编排domain负责规则与模型infrastructure负责技术实现interfaces负责入口与外部交互
7.1 创建通知用例示例
下面是应用层创建通知的简化示例:
package service
import (
"context"
"time"
"notify/internal/domain/entity"
"notify/internal/domain/event"
"notify/internal/domain/repository"
)
type EventBus interface {
Publish(ctx context.Context, evt event.NotificationCreated) error
}
type NotificationAppService struct {
notificationRepo repository.NotificationRepository
outboxRepo repository.OutboxRepository
bus EventBus
}
type CreateNotificationCmd struct {
BizID string
UserID int64
Type string
Title string
Content string
IdempotentKey string
}
func (s *NotificationAppService) CreateNotification(ctx context.Context, cmd CreateNotificationCmd) (int64, error) {
now := time.Now()
n := entity.NewNotification(
cmd.BizID,
cmd.UserID,
cmd.Type,
cmd.Title,
cmd.Content,
cmd.IdempotentKey,
now,
)
if err := s.notificationRepo.Save(ctx, n); err != nil {
return 0, err
}
evt := event.NotificationCreated{
NotificationID: n.ID,
OccurredAt: now,
}
if err := s.outboxRepo.SaveEvent(ctx, evt); err != nil {
return 0, err
}
return n.ID, nil
}
注意这里使用了“先存业务数据,再存事件”的模式。实际工程里应放在同一个数据库事务中执行。
7.2 Worker 执行发送示例
消费者收到通知创建事件后,负责执行完整发送流程:
package worker
import (
"context"
"fmt"
"notify/internal/application"
"notify/internal/domain"
)
type SendWorker struct {
NotificationRepo application.NotificationQueryRepository
DeliveryRepo application.DeliveryRepository
Router domain.Router
SenderRegistry *application.SenderRegistry
}
func (w *SendWorker) HandleCreated(ctx context.Context, notification domain.Notification) error {
channels := w.Router.Resolve(domain.RouteInput{
Type: notification.Type,
Preference: domain.UserPreference{
AllowInbox: true,
AllowEmail: true,
AllowSMS: true,
AllowPush: true,
},
})
for _, ch := range channels {
sender, ok := w.SenderRegistry.Get(ch)
if !ok {
return fmt.Errorf("sender not found, channel=%s", ch)
}
delivery := notification.NewDelivery(ch)
if err := w.DeliveryRepo.Save(ctx, delivery); err != nil {
return err
}
_ = delivery.MarkSending()
_ = w.DeliveryRepo.UpdateStatus(ctx, delivery.ID, delivery.Status)
result, err := sender.Send(ctx, domain.SendRequest{
NotificationID: notification.ID,
DeliveryID: delivery.ID,
Target: delivery.Target,
Title: notification.Title,
Content: notification.Content,
})
if err != nil {
delivery.MarkFailure("SEND_ERROR", err.Error())
if updateErr := w.DeliveryRepo.UpdateResult(ctx, delivery); updateErr != nil {
return updateErr
}
continue
}
if result.Success {
delivery.MarkSuccess()
} else {
delivery.MarkFailure(result.Code, result.Message)
}
if err := w.DeliveryRepo.UpdateResult(ctx, delivery); err != nil {
return err
}
}
return nil
}
这段代码省略了很多生产级细节,例如事务控制、日志、监控、trace、并发池、回执补偿等,但已经能说明清晰的职责分工:
- Worker 只编排流程
- Router 决定渠道
- Sender 负责实际发送
- Repository 负责持久化
7.3 HTTP 接口示例
通知中心对外通常提供一个统一创建接口:
package http
import (
"net/http"
"github.com/gin-gonic/gin"
"notify/internal/application/service"
)
type Handler struct {
app *service.NotificationAppService
}
type CreateNotificationRequest struct {
BizID string `json:"biz_id" binding:"required"`
UserID int64 `json:"user_id" binding:"required"`
Type string `json:"type" binding:"required"`
Title string `json:"title" binding:"required"`
Content string `json:"content" binding:"required"`
IdempotentKey string `json:"idempotent_key" binding:"required"`
}
func (h *Handler) CreateNotification(c *gin.Context) {
var req CreateNotificationRequest
if err := c.ShouldBindJSON(&req); err != nil {
c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()})
return
}
id, err := h.app.CreateNotification(c.Request.Context(), service.CreateNotificationCmd{
BizID: req.BizID,
UserID: req.UserID,
Type: req.Type,
Title: req.Title,
Content: req.Content,
IdempotentKey: req.IdempotentKey,
})
if err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
return
}
c.JSON(http.StatusAccepted, gin.H{
"notification_id": id,
"status": "ACCEPTED",
})
}
这里返回 202 Accepted 很符合异步通知场景,因为 API 只表示“请求已受理”,并不代表所有渠道已经发送成功。
7.4 一个更贴近生产实践的落地建议
如果这个通知中心真的要落地到中型团队环境中,我建议按以下顺序迭代:
第一阶段:先跑通核心链路
- 单一创建接口
- 支持站内信、邮件、短信、Push 四种渠道
- MySQL + Redis + MQ
- 基础状态追踪
- 简单失败重试
第二阶段:补齐治理能力
- 模板中心接入
- 用户偏好与黑名单控制
- 限流、熔断、告警
- 供应商主备切换
- 回执异步处理
第三阶段:增强运营与平台化能力
- 通知查询后台
- 发送效果统计
- 模板版本管理
- 任务重放
- 审计与灰度开关
八、总结
「用户通知中心」是一个非常典型的中型服务案例。它看似只是“发消息”,但真正做起来会同时涉及接口设计、领域建模、异步架构、状态机、数据库建模、缓存治理、容错策略和代码分层。
这类系统最重要的,不是追求“最炫的技术栈”,而是抓住几个设计核心:
- 边界清晰:通知中心只做发送能力,不侵入其他领域
- 分层明确:规则归规则,技术实现归技术实现
- 事件驱动:把接入与发送解耦,提升吞吐和稳定性
- 适配器模式:隔离外部渠道差异,降低供应商耦合
- 状态可追踪:任何通知都能查到走到了哪一步
- 容错可恢复:失败可重试、异常可降级、流程可补偿
如果把这套设计思路掌握住,后续你不仅能设计通知中心,也能举一反三地设计审批中心、任务中心、订单履约中心等中型业务服务。Go 在这类服务中的优势也非常明显:
- 并发模型适合 I/O 密集型发送场景
- 工程组织简洁,适合分层落地
- 部署轻量,适合微服务与 Worker 混合部署
对于一个中型 Go 服务来说,真正决定长期质量的,往往不是一两个“高级技巧”,而是架构层面是否提前把扩展性、可靠性和治理能力考虑进去。
📝 版权声明:本文为原创技术博客,转载请注明出处。
如文章中存在错误或不准确之处,欢迎在评论区指正,感谢您的阅读与支持!