订单状态变化、审批待办、设备告警和账户安全提醒都需要向用户发送消息。系统规模较小时,业务代码直接调用邮件、短信或站内信接口似乎最省事;随着渠道和业务增加,这种做法会迅速产生三个问题:
- 每个业务模块重复实现模板、收件人解析和异常重试。
- 渠道故障会拖慢业务事务,甚至导致业务成功但接口返回失败。
- 同一事件可能因消息队列重复消费、网络超时或人工补偿而被多次发送。
统一消息中心的目标不是简单封装几个发送接口,而是建立一条可追踪的投递流水线:业务系统只表达“发生了什么、通知谁”,消息中心负责决定“用什么模板、经过哪些渠道、何时重试以及如何查询结果”。
本文采用 Spring Boot 描述实现方式,数据库和消息队列产品可以按现有基础设施替换。示例强调接口与事务设计,不绑定某个云厂商;短信、邮件等外部渠道仍需依据供应商文档实现适配器。
先划清三个数据层次
统一消息中心最容易出现的建模错误,是用一条记录同时表示业务事件和渠道发送结果。一个审批事件可能同时生成站内信、邮件和 WebSocket 推送,因此至少需要区分三个概念:
- 消息请求:业务侧提交的一次通知意图,包含业务幂等键、模板和接收人。
- 用户消息:用户在站内看到的逻辑消息,可包含已读、撤回等状态。
- 投递任务:某条消息通过某个渠道的一次发送任务,记录重试次数和渠道回执。
可使用以下核心表结构:
CREATE TABLE message_request (
id BIGINT PRIMARY KEY,
idempotency_key VARCHAR(128) NOT NULL,
biz_type VARCHAR(64) NOT NULL,
template_code VARCHAR(64) NOT NULL,
recipient_id VARCHAR(64) NOT NULL,
variables_json TEXT NOT NULL,
created_at TIMESTAMP NOT NULL,
UNIQUE (idempotency_key)
);
CREATE TABLE delivery_task (
id BIGINT PRIMARY KEY,
request_id BIGINT NOT NULL,
channel VARCHAR(32) NOT NULL,
destination VARCHAR(256),
status VARCHAR(24) NOT NULL,
attempt_count INT NOT NULL DEFAULT 0,
next_attempt_at TIMESTAMP NULL,
provider_msg_id VARCHAR(128),
last_error_code VARCHAR(64),
last_error_text VARCHAR(512),
created_at TIMESTAMP NOT NULL,
updated_at TIMESTAMP NOT NULL,
UNIQUE (request_id, channel)
);
CREATE INDEX idx_delivery_schedule
ON delivery_task (status, next_attempt_at);
idempotency_key应由业务身份构成,例如approval:9421:assignee:u1008:v3,而不是每次请求随机生成 UUID。随机值只能区分请求,无法识别同一业务动作的重复提交。若同一事件允许多次提醒,应把提醒轮次或业务版本纳入键中。
投递链路与事务边界
推荐的处理顺序如下:
- 业务服务完成自身事务,并写入本地事件表(Outbox)。
2.后台发布器把事件投递到消息队列。 - 消息中心消费事件,在同一数据库事务中创建消息请求和各渠道任务。
- 独立工作线程领取投递任务,调用渠道适配器。
- 根据返回结果将任务标记为成功、待重试或永久失败。
Outbox 的价值在于避免“业务数据已提交,但消息队列发送失败”的双写不一致。业务记录与事件记录处于同一事务,发布器即使暂时停机,也可以稍后继续扫描未发布事件。不要在业务数据库事务中直接调用外部邮件或短信接口,因为不可控的网络等待会延长锁持有时间,而且数据库回滚无法撤销已经发出的通知。
消息中心的消费端仍要执行幂等插入。消息队列通常只能降低重复概率,不能替代业务幂等。数据库唯一约束是最终防线,应用层的“先查询再插入”在并发下并不充分。
@Transactional
public long accept(NotificationCommand command) {
try {
MessageRequest request = requestRepository.insert(command);
List<Channel> channels = policyService.resolve(
command.bizType(), command.recipientId());
for (Channel channel : channels) {
deliveryTaskRepository.insertPending(
request.id(), channel, resolveDestination(command, channel));
}
return request.id();
} catch (DuplicateKeyException duplicate) {
return requestRepository
.findIdByIdempotencyKey(command.idempotencyKey())
.orElseThrow(() -> duplicate);
}
}
这里捕获唯一键冲突后返回原请求编号,使上游重试得到稳定结果。实际项目还应校验重复请求的模板、接收人等关键字段是否与原请求一致,防止错误复用幂等键掩盖数据冲突。
用策略和适配器隔离渠道差异
渠道编排负责选择渠道,渠道适配器负责执行发送,两者不应混在一起。例如高优先级设备告警可以选择“站内信加短信”,普通状态更新只写站内信;用户离线时是否补发邮件,也属于策略而非邮件适配器的职责。
public interface ChannelSender {
Channel channel();
SendResult send(DeliveryContext context);
}
public record SendResult(
boolean accepted,
boolean retryable,
String providerMessageId,
String errorCode,
String errorMessage
) {
}
@Component
public final class WebSocketSender implements ChannelSender {
private final SimpMessagingTemplate messagingTemplate;
private final OnlineUserRegistry onlineUsers;
public WebSocketSender(SimpMessagingTemplate messagingTemplate,
OnlineUserRegistry onlineUsers) {
this.messagingTemplate = messagingTemplate;
this.onlineUsers = onlineUsers;
}
@Override
public Channel channel() {
return Channel.WEBSOCKET;
}
@Override
public SendResult send(DeliveryContext context) {
if (!onlineUsers.isOnline(context.recipientId())) {
return new SendResult(false, false, null,
"USER_OFFLINE", "recipient is offline");
}
messagingTemplate.convertAndSendToUser(
context.recipientId(), "/queue/notifications", context.payload());
return new SendResult(true, false, null, null, null);
}
}
需要注意,convertAndSendToUser成功通常只说明消息已交给当前应用内的消息通道,并不天然等价于客户端已经展示。若业务要求“送达确认”,客户端必须回传包含messageId的 ACK,服务端再记录确认时间。高可靠场景还要为 ACK 设置超时,并明确超时后是转离线信箱、切换其他渠道,还是仅记录未确认。
邮件或短信适配器的凭据应来自环境变量,配置中只引用变量名:
notification:
mail:
host: ${
MAIL_HOST}
username: ${
MAIL_USERNAME}
password: ${
MAIL_PASSWORD}
connect-timeout: 3s
read-timeout: 5s
worker:
batch-size: 50
lease-duration: 30s
max-attempts: 6
生产环境应由部署平台注入变量或挂载密钥文件,并限制日志输出,避免把凭据、完整手机号、邮箱地址和模板变量写入普通应用日志。
安全领取任务,避免并发重复发送
多实例工作线程不能同时处理同一任务。一种通用做法是给任务增加lease_owner和lease_until字段,领取时使用数据库行锁或原子条件更新,把短时间处理权租给某个实例。支持SKIP LOCKED的数据库可以在短事务中批量领取:
BEGIN;
SELECT id
FROM delivery_task
WHERE status IN ('PENDING', 'RETRY_WAIT')
AND (next_attempt_at IS NULL OR next_attempt_at <= CURRENT_TIMESTAMP)
ORDER BY id
FOR UPDATE SKIP LOCKED
LIMIT 50;
-- 将查询到的任务更新为 PROCESSING,并写入租约信息。
COMMIT;
外部发送必须在领取事务提交后进行,不能持有行锁等待网络响应。工作进程崩溃后,扫描器可以回收租约已过期的任务。数据库是否支持SKIP LOCKED以及对应语法取决于具体产品和版本;不支持时,可以使用带状态条件的原子UPDATE逐条竞争任务。
即使任务领取做到了互斥,外部调用仍可能出现“渠道已经接受请求,但本地在记录成功前断网”的不确定状态。若渠道支持幂等请求号,应把稳定的任务编号作为渠道幂等键;若不支持,就不能承诺严格的恰好一次发送,只能在“可能重复”和“可能漏发”之间按业务风险取舍。
重试不是简单循环
应先把错误分类,再决定是否重试:
- 网络超时、连接失败、渠道限流通常可以延迟重试。
- 地址格式错误、模板不存在、收件人退订通常应永久失败。
- 鉴权失败可能来自错误配置,快速循环重试只会扩大故障,应暂停渠道并触发告警。
退避时间可以采用指数增长并加入随机扰动:
Duration nextDelay(int attempt) {
long baseSeconds = Math.min(300, 1L << Math.min(attempt, 8));
long jitter = ThreadLocalRandom.current().nextLong(0, 6);
return Duration.ofSeconds(baseSeconds + jitter);
}
最大重试次数、最大间隔和可重试错误集合都应配置化。达到上限后进入DEAD状态,由人工或补偿任务处理,不能无限重试。人工重放应生成审计记录,并继续沿用原任务身份或显式增加重放轮次,避免绕过幂等规则。
模板治理与数据最小化
模板建议保存主题、正文、渠道类型、语言和版本。投递任务创建时应冻结模板版本;否则模板发布后,历史任务重试可能发送与首次尝试不同的内容。变量渲染应采用模板引擎的结构化上下文,不要用连续的字符串替换实现。
进入模板的变量必须有白名单。模板渲染失败应在调用渠道前暴露,并记录缺失变量名称,但日志中不要输出密码、验证码、身份证号等值。对于 WebSocket 下发的富文本,前端仍需进行安全渲染;消息来自内部服务并不意味着内容天然可信。
可观测性与验收步骤
上线前至少建立以下指标:待处理任务数、最老任务等待时间、各渠道成功率、重试率、永久失败数和发送耗时分布。成功率需要按渠道和错误码拆分,否则一个渠道故障会被总体数据稀释。
可以按下面的步骤进行验收,而不依赖虚构的吞吐结论:
- 使用相同幂等键并发提交,确认只产生一条请求和每渠道一条任务。
- 让模拟渠道先返回超时再成功,确认任务经过
RETRY_WAIT并最终完成。 - 在渠道返回成功后、数据库更新前终止进程,检查恢复策略及重复风险。
- 启动两个工作实例,确认同一任务没有被同时领取。
- 提交缺少模板变量的请求,确认不会调用外部渠道。
- 让 WebSocket 客户端断线重连,验证未读消息从持久化存储补齐,而不是依赖内存队列。
- 检查日志、追踪和告警内容,确认敏感变量已脱敏。
压测结果与线程数、连接池、渠道配额、数据库延迟和消息大小直接相关,应在目标环境中测量。不能仅凭单机测试推导生产容量。
常见问题
WebSocket 能否替代站内信存储?
不能。WebSocket 是在线传输通道,连接可能随网络切换、页面关闭或实例重启而中断。需要离线可见和已读状态时,消息必须持久化,WebSocket 只负责通知客户端有新内容。
是否需要保证所有渠道同时成功?
通常不应把多个外部渠道放进一个原子事务。更实际的语义是每个投递任务独立成功,并由业务策略定义“任一渠道成功即可”还是“指定渠道必须成功”。这种规则需要体现在聚合状态中,而不是依赖数据库事务强行绑定。
为什么不能只依靠消息队列重试?
队列重试通常面向消费失败,难以表达每个渠道不同的退避、暂停、人工重放和回执状态。把投递状态持久化后,队列负责唤醒处理,数据库负责保存事实,两者职责更清楚。
用户修改邮箱后,历史任务发到哪里?
这取决于业务语义。若要求按事件发生时的地址发送,应在创建任务时冻结地址;若要求始终使用最新联系方式,则发送前解析。安全通知等场景还需防止账户联系方式刚被恶意修改后立即接收敏感消息,因此必须由产品和安全规则共同决定。
总结
统一消息中心的核心是把通知意图与渠道执行解耦,并承认外部调用无法与本地数据库形成天然原子事务。可落地的方案需要业务幂等键、请求与投递任务分层、短事务领取、错误分类重试、模板版本冻结以及可查询的审计状态。WebSocket、邮件和短信只是适配器;真正决定系统能否稳定运行的,是明确的状态机、故障边界和恢复流程。