返回首页

Golang工程化:5.8 一个中型服务的架构设计案例拆解

以「用户通知中心」为例,完整拆解一个中型 Go 服务从需求分析、架构设计到代码落地的全过程。

一、案例背景:设计一个用户通知中心服务

在多数互联网系统里,通知能力并不是最显眼的业务模块,却几乎贯穿所有核心链路。用户注册成功后需要欢迎消息,订单状态变化后需要短信提醒,审批流结束后需要邮件通知,App 活动上线后需要 Push 触达。随着业务扩展,通知类型、发送渠道、发送策略和失败处理逻辑会迅速膨胀。

如果一开始只是简单写几个 SendEmail()SendSMS() 函数,短期内可以工作,但当业务逐渐复杂后,系统通常会暴露出以下问题:

  • 渠道调用逻辑散落在各业务服务中,难以统一治理
  • 模板、重试、限流、状态追踪等能力重复实现
  • 多渠道发送策略缺少统一编排,修改成本高
  • 渠道供应商切换困难,强耦合严重
  • 无法准确追踪一条通知从创建到最终送达的全过程

因此,我们设计一个中型服务规模的「用户通知中心」:

  • 支持多渠道:站内信、邮件、短信、Push
  • 支持同步下发与异步发送
  • 支持模板渲染、路由分发、状态追踪
  • 支持失败重试、死信处理、幂等控制
  • 支持后续扩展更多渠道与业务场景

这个服务本质上是一个典型的中台型基础能力服务。它对外承接业务系统的通知需求,对内屏蔽不同渠道的复杂性,并通过统一模型实现扩展性和可运维性。

二、需求分析与系统边界划分

在正式设计架构之前,先把需求拆清楚。架构设计最怕的问题不是“技术选型不够高级”,而是系统边界模糊,结果导致职责混乱。

2.1 核心功能需求

用户通知中心需要提供以下核心能力:

  1. 统一通知接入

    • 业务方通过统一 API 提交通知请求
    • 支持指定单渠道或多渠道发送
    • 支持立即发送、定时发送、事件触发发送
  2. 通知内容管理

    • 支持通知模板
    • 支持变量渲染
    • 支持不同渠道的内容格式适配
  3. 消息路由能力

    • 根据通知类型、用户偏好、业务优先级决定发送渠道
    • 支持降级策略,例如 Push 失败后转短信
  4. 发送执行与结果回写

    • 统一调用不同渠道服务商
    • 记录发送状态:待发送、发送中、成功、失败、部分成功
    • 支持供应商回执更新
  5. 可靠性保障

    • 幂等控制
    • 重试机制
    • 死信队列处理
    • 限流与熔断
  6. 运营与审计

    • 通知查询
    • 发送明细查看
    • 异常告警
    • 审计日志

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 的外部依赖延迟不稳定
  • 高峰期通知请求量可能短时暴涨
  • 部分通知允许稍后发送,但不能丢失

因此,一个合理的流程是:

  1. API 接收通知请求
  2. 落库生成通知任务
  3. 发布“通知已创建”事件
  4. 异步消费者执行路由、渲染、发送
  5. 状态回写并视情况继续派发重试事件

这样能将“请求接入”和“实际发送”解耦,避免上游业务被下游渠道拖慢。

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 核心流程说明

请求创建流程

  1. 上游业务调用通知中心 API
  2. Interface 层完成参数校验
  3. Application 层创建 Notification 聚合与 Delivery 记录
  4. 仓储层持久化数据
  5. 发布通知创建事件
  6. 返回请求已受理结果

异步发送流程

  1. Worker 消费通知创建事件
  2. 查询通知详情与接收人信息
  3. 根据策略进行消息路由
  4. 渲染渠道内容
  5. 调用各渠道适配器发送
  6. 回写渠道发送结果
  7. 失败场景进入延迟重试或死信处理

状态追踪流程

  1. 初始状态为 PENDING
  2. 发送中更新为 SENDING
  3. 单渠道成功则写入子记录为 SUCCESS
  4. 单渠道失败则写入 FAILED
  5. 如果通知是多渠道,可根据整体规则汇总为 PARTIAL_SUCCESSSUCCESS
  6. 供应商异步回执到达后再更新最终投递状态

四、核心模块设计

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 发送状态追踪模块

通知中心如果没有完善的状态追踪,排查问题会非常痛苦。我们不能只记录“发了没发”,而是要跟踪:

  • 通知任务是否已创建
  • 哪些渠道被选中
  • 每个渠道当前状态如何
  • 外部供应商返回了什么结果
  • 是否还会继续重试
  • 最终整体结果是什么

建议把状态设计成两层:

  1. 通知主任务状态:描述整条通知的整体进度
  2. 渠道投递状态:描述每个渠道实例的发送结果

推荐状态机如下:

维度 状态 含义
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 模式的做法是:

  1. 在同一个本地事务里写入业务表和 outbox_events
  2. 后台投递器轮询未发送事件
  3. 成功投递到 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 幂等与去重设计

通知中心很容易因为上游重试、网络抖动、消息重复投递而出现重复发送。因此需要在多个层面做幂等控制:

  • 请求幂等:基于业务幂等键避免重复创建通知
  • 消费幂等:同一事件重复消费时保证不重复发出
  • 渠道幂等:向供应商发送时尽量带业务请求号

一个常见做法是:

  1. 上游传入 biz_id + scene + receiver 组合幂等键
  2. 创建通知时先查 Redis,再落库唯一索引校验
  3. 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 服务来说,真正决定长期质量的,往往不是一两个“高级技巧”,而是架构层面是否提前把扩展性、可靠性和治理能力考虑进去。


📝 版权声明:本文为原创技术博客,转载请注明出处。

如文章中存在错误或不准确之处,欢迎在评论区指正,感谢您的阅读与支持!

上一篇

Golang工程化:5.7 稳定性设计之超时、重试、熔断、幂等

下一篇

Golang工程化:5.1 Go 项目结构与 package 设计