搞PHP这么多年最烦的一件事就是“用户消息到底发了没有”。尤其是业务跑起来之后订单通知、支付回调、验证码、站内信、模板消息全挤在一起最初那种“临时写个mail()函数塞进业务流程里”的做法到后面就是灾难现场。后来我干脆花了两个周末把整套消息通知系统用PHP重新捋了一遍从通道抽象、模板引擎到异步队列全部自己控制。这篇文章就把这套方案从选型到落地完整拆开讲包括建的每张表、写的每个核心类、踩过的坑给需要自己搭通知系统的朋友一个能直接参考的底子。1. 整体设计与选型思路1.1 消息通知在业务里的真实位置很多人觉得消息通知就是一个“发个邮件、发个短信”的小功能不值得单独做系统。这话在日活几百、消息量几千的时候确实是成立的。但一旦业务跑起来你会发现通知这件事横跨了所有核心链路用户注册要发验证码、下单要发确认通知、支付成功要发回执、退款要发到账通知、风控拦截要发告警、运营活动要发营销消息。如果这些逻辑全部散落在各个业务代码里每一处都直接调第三方API那你面临的就不只是重复代码问题而是四个非常现实的问题。第一是耦合问题。业务代码里到处是短信和邮件的调用点哪天要换短信服务商或者加一个站内信渠道你得把全项目的代码翻一遍。第二是可靠性问题。通知发出去是有成功失败概念的如果发失败了你重试不重试重试逻辑写在哪接口超时怎么办这些如果不统一处理用户那边就会莫名其妙地收不到消息。第三是性能问题。业务流程里同步调第三方短信通道一次请求几百毫秒甚至超时用户端的下单接口就被拖慢了。第四是审计问题。用户投诉说没收到验证码你要查到底发没发、什么时候发的、渠道返回什么错误没有统一的日志记录就得抓瞎。所以当业务量上来之后把消息通知抽成一个独立的子系统几乎是必然的。我做的这套方案核心解决的就是这四个问题通过统一消息模型解耦业务方与渠道方通过状态机和重试机制保证投递可靠性通过异步队列把第三方接口调用挪出主流程通过完整日志表保留每条消息的投递链路记录。1.2 方案选型为什么是PHP Redis WebSocket技术选型这件事没有银弹只有合不合适。我自己主力栈是PHP所以后端服务果断用PHP这不是什么情怀选择而是这套系统生在PHP项目里、长在PHP项目里保持技术栈统一能让维护成本降到最低。消息通知本身属于IO密集型场景对CPU计算要求极低PHP在CLI模式下跑长驻消费者脚本完全扛得住。这里有个关键点消息通知系统一定要脱离FPM跑不能让nginx/php-fpm来承载消费者进程因为fpm的生命周期模型不适合常驻循环跑队列会被超时回收掉。CLI模式下用pcntl或者直接while循环配合Redis阻塞读取都行我用的是CLI Redis Stream的消费者组模式稳定性实测比list类型的队列高很多。队列选Redis Stream而不是RabbitMQ原因很直接大多数中小团队已经有Redis了再为了一个通知系统引入RabbitMQ/Kafka运维成本和资源占用都不划算。Redis Stream是Redis 5.0自带的消息队列方案支持消费者组、消息确认、死信概念对于通知这个量级日峰值几十万条以内完全够用。如果你团队已经上了Kafka那当然用Kafka但从零开始搭的话Redis Stream是性价比最高的选择。实时推送这一块WebSocket是绕不开的。浏览器原生WebSocket支持已经很普及没必要用古老的轮询方案。PHP这边我用的Workerman来做WebSocket服务端因为Workerman是纯PHP实现、部署简单、文档也全。如果你不想引这个重东西还有一个轻量替代方案就是SSEServer-Sent Events基于HTTP长连接做服务端单向推送很多通知场景其实只需要服务端往客户端推SSE够用且实现更简单。我最终选了Workerman是为了后续如果要做客服系统之类的双向通信不用再换底子。1.3 消息通知系统的三种基本形态在设计架构之前我先把通知这件事拆成了三个形态分别对应不同的技术处理方式。第一种是同步通知。典型场景就是登录验证码、手机号绑定验证码。这类消息用户在线等待必须尽快返回结果等不起队列轮询再消费的过程。处理方式是直接同步调用渠道接口但要做好超时控制和降级。我的做法是给每个渠道设置超时阈值验证码通道走直连超过3秒没响应就快速失败给前端提示重新发送。第二种是异步通知。订单支付成功通知、发货通知、营销邮件这类业务方不关心发送结果用户也不是立刻就要看到完全走队列异步化。业务方只需要把消息对象丢进Redis消费者进程异步去调渠道接口发完更新状态失败进入重试队列。这样彻底不让第三方接口延迟影响主业务。第三种是实时推送。用户在网页端、管理后台、或者移动端WebView里需要页面即时弹出新消息。这类走WebSocket或者SSE服务端收到异步消息后除了走邮件/短信渠道同时向在线连接推送一份客户端收到后刷新未读数或弹提示。三个形态对应三类不同的技术组件同步走直连、异步走队列、实时走长连接。架构上它们共享的是同一套消息模型、模板体系和服务化接口。业务方的接入成本只体现在选哪个API上其余全部复用。2. 核心模块与数据结构设计2.1 消息表设计状态机驱动投递全流程消息通知系统的地基是消息表。消息一旦入表整个生命周期都要能看到。我的消息表字段设计比较全核心字段包括消息ID、业务类型、渠道、接收方、标题、内容、状态、重试次数、扩展参数、发送时间、完成时间。这里面最关键的是状态字段我用了tinyint类型状态流转是这套系统的灵魂。状态定义我设置为0待发送、1已发送、2发送失败、3发送中、4已撤回。待发送状态主要是给异步消息用的消息入队列后先标记待发送消费者取出时置为发送中成功或失败再更新终态。这里有个容易踩的坑就是不要把“发送中”状态省掉。如果没有发送中状态你很难判断一个长时间卡住的消息到底是“还没被消费”还是“已经卡死在半路”。有了发送中的时间戳你就能写一个超时巡检脚本把超过5分钟还在发送中的消息捞出来重置成待发送重新入队。接口扩展参数我用的是text字段JSON格式存储因为不同业务方可能会传入一些渠道特有的参数比如短信签名、邮件附件ID、微信模板消息的跳转路径。把这些塞进JSON里比在表里堆几十个可空字段清爽得多。另外消息表一定要按时间做分区。我遇到过的问题是消息量一上去之后全表扫描导致后台消息查询接口响应超过5秒。解决方案是按月做RANGE分区查询时带上时间范围条件就会命中分区裁剪实测查询速度提升了十倍不止。2.2 模板表设计让业务方和渠道方解耦消息模板我是单独建表的因为业务方不能直接拼内容尤其是短信有字数限制按字符数计费错了就要烧钱。模板表的思路是把内容渲染逻辑收拢到通知系统内部业务方只需要传一个模板标识符加一堆变量参数系统负责从模板表加载内容、做变量替换、生成最终消息内容。模板表核心字段包括模板ID、模板编码业务唯一、标题模板、内容模板、渠道类型、变量定义JSON格式记录这个模板要用哪些参数、创建人、状态。模板内容里的变量我用双大括号占位比如“您的验证码是{{code}}{{minutes}}分钟内有效”渲染的时候用正则扫描替换。这里要特别说一句变量解析和PHP本身的变量解析容易混淆所以我特意没有用PHP的字符串模板语法而是自己规定了一套占位符规则这样内容和代码彻底隔离业务方甚至可以自己在后台维护模板文案而不需要改代码。模板变量定义的JSON还有一个隐藏作用校验。业务方传参时系统会自动比对传进来的参数跟模板声明的变量定义是否吻合缺参数直接报错而不是渲染出来一条“您的验证码是”这种残缺消息。这个小机制帮我拦下了无数线上问题。2.3 渠道适配层面向接口编程的Channel架构渠道适配层是整个通知系统的核心抽象。我定义了一个MessageChannel接口让每条实际渠道都实现这个接口。接口里主要包含三个方法send(发送消息)、checkConfig(检查渠道配置)、getChannelName(返回渠道标识)。邮件、短信、站内信、WebSocket推送每个都是这个接口的实现类。业务方根本不需要关心消费端到底用什么渠道发出去的系统内部通过渠道映射决定。为什么这么设计举个实际场景之前项目里用的短信服务商A在某个地区网关一直抖我们需要在通知系统层面快速把短信流量切到服务商B同时对用户无感。如果业务方直接集成了服务商A的SDK这种切换就得改业务代码。有了Channel抽象层我只需要新增一个SmsProviderB类实现同一个接口然后修改配置中心的路由表把短信渠道的provider指向B流量几分钟内就切过去了。这种灵活性就是面向接口编程带来的直接收益。在Channel接口之上我还包了一层“分发器”。分发器负责读消息的channel_type字段用简单工厂模式找到对应的Channel实现类执行发送再把结果回写消息表。如果后续需要增加新的通知方式比如微信服务号模板消息只需要写一个新的Channel实现类注册到工厂里其他代码一律不用动。3. 实操从零搭建一套PHP消息通知系统3.1 环境准备与依赖清单这一节把整套系统的运行环境和使用到的库列个底方便你复现的时候有一个清晰的依赖图。我本地的开发环境是PHP 8.1 Redis 6.2 MySQL 8.0生产环境也是一致的版本组合避免开发环境和线上环境行为不一致。PHP这边需要安装这些扩展PDO_MYSQL数据库访问、Redis队列与缓存、PCNTLCLI模式下进程控制Workerman依赖、OpenSSL邮件加密传输、WebSocket的wss支持、MBString字符编码处理。Swoole扩展不做强制要求因为我们用Workerman做WebSocket服务端Workerman在纯PHP模式下也能跑只是性能不如开了event扩展的模式。如果线上对并发有要求建议还是把event扩展装上事件驱动模型能压榨出更高的连接数。用到的PHP第三方库就三个PHPMailer邮件发送的事实标准、WorkermanWebSocket服务端框架、vlucas/phpdotenv环境变量管理。PHPMailer虽然老但是稳定得吓人支持SMTP认证、SSL/TLS加密各种邮箱服务商的兼容性问题它基本都处理过了。Workerman则是PHP常驻内存应用的实用派选择生态成熟踩坑资料也多。安装方式我直接用Composer管理依赖composer.json里声明这三件套一条composer install就完了。数据库迁移建议用原生的SQL文件管理通知系统表结构变化不频繁没必要上重量级迁移工具。3.2 数据库初始化消息、模板、映射三张表的建表SQL下面是我线上正在用的三张核心表的建表SQL字段注释都写好了你直接拿去改改业务类型就能用。首先是消息表这张表承载所有消息的流转记录。CREATE TABLE message ( id bigint(20) unsigned NOT NULL AUTO_INCREMENT, msg_id varchar(64) NOT NULL COMMENT 业务消息唯一ID调用方生成用于幂等, biz_type varchar(64) NOT NULL COMMENT 业务类型如order_paid、verify_code, channel_type tinyint(4) NOT NULL COMMENT 渠道类型 1邮件 2短信 3站内信 4WebSocket推送, receiver varchar(255) NOT NULL COMMENT 接收方邮箱/手机号/用户ID, title varchar(255) NOT NULL DEFAULT COMMENT 消息标题, content text NOT NULL COMMENT 消息内容, template_code varchar(64) NOT NULL DEFAULT COMMENT 模板编码, template_params text COMMENT 模板渲染参数JSON格式, ext_params text COMMENT 扩展参数JSON格式, status tinyint(4) NOT NULL DEFAULT 0 COMMENT 0待发送 1已发送 2发送失败 3发送中 4已撤回, retry_count tinyint(4) NOT NULL DEFAULT 0 COMMENT 已重试次数, max_retry tinyint(4) NOT NULL DEFAULT 3 COMMENT 最大重试次数, channel_resp text COMMENT 渠道返回的原始响应用于排查问题, send_time datetime DEFAULT NULL COMMENT 实际发送成功时间, next_retry_time datetime DEFAULT NULL COMMENT 下次重试时间, created_at datetime NOT NULL DEFAULT CURRENT_TIMESTAMP, updated_at datetime NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP, PRIMARY KEY (id), UNIQUE KEY uk_msg_id (msg_id), KEY idx_biz_type (biz_type), KEY idx_receiver (receiver), KEY idx_status_next_retry (status,next_retry_time) ) ENGINEInnoDB DEFAULT CHARSETutf8mb4 COMMENT消息主表;msg_id是唯一的调用方可以用业务主键类型拼接生成这样同一条业务消息即便被重复提交在入库时因为唯一键约束也只会产生一条记录从源头保证了幂等。idx_status_next_retry这个联合索引是给重试扫描脚本用的能快速捞出来哪些消息已经到了重试时间。模板表的建表语句相对简洁核心是模板内容和变量定义的隔离CREATE TABLE message_template ( id int(11) unsigned NOT NULL AUTO_INCREMENT, template_code varchar(64) NOT NULL COMMENT 模板编码全局唯一, template_name varchar(128) NOT NULL COMMENT 模板名称, channel_type tinyint(4) NOT NULL COMMENT 适用渠道 1邮件 2短信 3站内信, title_tpl varchar(255) NOT NULL DEFAULT COMMENT 标题模板支持{{var}}占位符, content_tpl text NOT NULL COMMENT 内容模板支持{{var}}占位符, vars text COMMENT 变量定义JSON如[{name:code,required:true}], status tinyint(4) NOT NULL DEFAULT 1 COMMENT 1启用 0停用, created_at datetime NOT NULL DEFAULT CURRENT_TIMESTAMP, updated_at datetime NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP, PRIMARY KEY (id), UNIQUE KEY uk_template_code (template_code) ) ENGINEInnoDB DEFAULT CHARSETutf8mb4 COMMENT消息模板表;第三张表是业务方和渠道的路由映射表。它决定了哪个业务类型走哪条渠道、优先级是多少。这张表的存在让“把某类消息从一个渠道切到另一个渠道”变成了改一条DB记录的事而不是发一次版本。CREATE TABLE message_route ( id int(11) unsigned NOT NULL AUTO_INCREMENT, biz_type varchar(64) NOT NULL COMMENT 业务类型, channel_type tinyint(4) NOT NULL COMMENT 渠道类型, priority tinyint(4) NOT NULL DEFAULT 0 COMMENT 优先级数字越小越先尝试, status tinyint(4) NOT NULL DEFAULT 1 COMMENT 1启用 0停用, created_at datetime NOT NULL DEFAULT CURRENT_TIMESTAMP, PRIMARY KEY (id), UNIQUE KEY uk_biz_channel (biz_type,channel_type) ) ENGINEInnoDB DEFAULT CHARSETutf8mb4 COMMENT消息路由表;3.3 核心服务类实现通知管理器与三类渠道数据库就绪之后开始写代码。整个通知系统的核心入口是一个NotificationManager服务类业务方所有发消息的请求都走这个类的静态方法。这个类做了四件事校验参数、渲染模板、落库、决定同步直发还是投递队列。?php declare(strict_types1); namespace App\Notification; use App\Notification\Channel\ChannelFactory; use App\Notification\Template\TemplateRenderer; class NotificationManager { private TemplateRenderer $renderer; private ChannelFactory $factory; private \PDO $db; private \Redis $redis; public function __construct(\PDO $db, \Redis $redis) { $this-db $db; $this-redis $redis; $this-renderer new TemplateRenderer(); $this-factory new ChannelFactory($db); } public function send(array $params, bool $sync false): array { $msgId $params[msg_id] ?? $this-generateMsgId($params[biz_type]); // 幂等检查消息ID已存在则直接返回不重复发送 $stmt $this-db-prepare(SELECT id, status FROM message WHERE msg_id ?); $stmt-execute([$msgId]); $exists $stmt-fetch(\PDO::FETCH_ASSOC); if ($exists) { return [code 0, msg duplicated, id $exists[id]]; } // 通过路由表确定渠道 $channelType $this-resolveChannel($params[biz_type]); // 渲染模板内容 [$title, $content] $this-renderer-render( $params[template_code], $params[template_params], $channelType ); $insertSql INSERT INTO message (msg_id, biz_type, channel_type, receiver, title, content, template_code, template_params, ext_params, status, next_retry_time) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?); $stmt $this-db-prepare($insertSql); $stmt-execute([ $msgId, $params[biz_type], $channelType, $params[receiver], $title, $content, $params[template_code], json_encode($params[template_params], JSON_UNESCAPED_UNICODE), json_encode($params[ext_params] ?? [], JSON_UNESCAPED_UNICODE), $sync ? 3 : 0, $sync ? date(Y-m-d H:i:s) : date(Y-m-d H:i:s, time() 5) ]); $messageId (int)$this-db-lastInsertId(); if ($sync) { // 同步直发走渠道并更新状态 return $this-factory-syncSend($messageId); } // 异步投递到Redis Stream $this-redis-xAdd(stream:message:send, *, [ message_id $messageId, msg_id $msgId, retry 0 ]); return [code 0, msg queued, id $messageId]; } private function resolveChannel(string $bizType): int { $stmt $this-db-prepare(SELECT channel_type FROM message_route WHERE biz_type ? AND status 1 ORDER BY priority ASC LIMIT 1); $stmt-execute([$bizType]); $row $stmt-fetch(\PDO::FETCH_ASSOC); if (!$row) { throw new \RuntimeException(biz_type [{$bizType}] has no route configured); } return (int)$row[channel_type]; } private function generateMsgId(string $bizType): string { return $bizType . _ . date(YmdHis) . _ . uniqid(, true); } }这里面有两个设计点需要展开说说。第一是幂等检查前置我特意放在模板渲染之前。为什么要前置因为模板渲染是有成本的如果消息已经存在了没必要再跑一遍渲染逻辑。第二是同步和异步的状态初始化不一样同步消息直接标记为“发送中”异步消息先标记“待发送”再入队。这样做的好处是状态与处理方式严格对应排查问题时看状态就能猜到这条消息是怎么进来的。接下来是渠道抽象基类。所有渠道实现类都继承自这个基类基类负责公共逻辑比如更新消息状态、记录渠道响应子类只需要实现具体的发送方法。?php declare(strict_types1); namespace App\Notification\Channel; abstract class AbstractChannel { protected \PDO $db; public function __construct(\PDO $db) { $this-db $db; } abstract public function send(int $messageId, array $message): bool; protected function markSending(int $messageId): void { $stmt $this-db-prepare(UPDATE message SET status 3 WHERE id ?); $stmt-execute([$messageId]); } protected function markSent(int $messageId, string $resp ): void { $stmt $this-db-prepare( UPDATE message SET status 1, channel_resp ?, send_time NOW() WHERE id ? ); $stmt-execute([$resp, $messageId]); } protected function markFailed(int $messageId, string $reason): void { $stmt $this-db-prepare(UPDATE message SET status 2, channel_resp ? WHERE id ?); $stmt-execute([$reason, $messageId]); } protected function markRetry(int $messageId, int $currentRetry, int $maxRetry): bool { if ($currentRetry $maxRetry) { $this-markFailed($messageId, max retry exceeded); return false; } $nextRetry time() (int)pow(2, $currentRetry) * 60; $stmt $this-db-prepare( UPDATE message SET status 0, retry_count retry_count 1, next_retry_time ? WHERE id ? ); $stmt-execute([date(Y-m-d H:i:s, $nextRetry), $messageId]); return true; } }重试策略的指数退避是写在基类里的每次失败后的下次重试时间按 2 的当前重试次数次方乘以 60 秒递增。第一次失败1分钟后重试、第二次4分钟后、第三次9分钟后。这个策略我在生产环境跑下来比较合适不会因为重试太密集把第三方接口打死也不会让用户等太久。邮件渠道用PHPMailer实现短信渠道以阿里云为例站内信渠道直接入站内信表?php declare(strict_types1); namespace App\Notification\Channel; use PHPMailer\PHPMailer\PHPMailer; class EmailChannel extends AbstractChannel { private array $config; public function __construct(\PDO $db, array $config) { parent::__construct($db); $this-config $config; } public function send(int $messageId, array $message): bool { $this-markSending($messageId); $mail new PHPMailer(true); try { $mail-isSMTP(); $mail-Host $this-config[smtp_host]; $mail-SMTPAuth true; $mail-Username $this-config[smtp_user]; $mail-Password $this-config[smtp_pass]; $mail-SMTPSecure $this-config[smtp_secure] ?? PHPMailer::ENCRYPTION_SMTPS; $mail-Port (int)$this-config[smtp_port]; $mail-CharSet UTF-8; $mail-setFrom($this-config[from_email], $this-config[from_name]); $mail-addAddress($message[receiver]); $mail-Subject $message[title]; $mail-Body $message[content]; $mail-isHTML(true); if ($mail-send()) { $this-markSent($messageId, sent by smtp); return true; } $this-markFailed($messageId, $mail-ErrorInfo); return false; } catch (\Throwable $e) { $this-markFailed($messageId, $e-getMessage()); return false; } } }邮件这块有个经验之谈PHPMailer的isHTML(true)之后一定要把AltBody也设置上否则部分邮件客户端会显示成一片空白或者只显示HTML源码。加上$mail-AltBody strip_tags($message[content]);能有效提升兼容性。另外smtp_port和smtp_secure必须成对匹配SSL对应465TLS对应587搞混了会浪费很多时间排查超时问题。短信渠道这里要特别提醒一个坑阿里云短信接口的API版本和签名机制经常变。我遇到过最坑的是签名算法从原来的RPC风格改成了ROA风格和文档对不上。我的建议是短信类的第三方对接在Channel内部再封装一层Provider接口这样换个服务商只需要改Provider实现Channel主逻辑完全复用。站内信渠道相对简单落库到一张notice表用户在APP或网页端拉取。发送就是往里插一条记录public function send(int $messageId, array $message): bool { $this-markSending($messageId); $stmt $this-db-prepare( INSERT INTO notice (user_id, title, content, is_read, created_at) VALUES (?, ?, ?, 0, NOW()) ); $stmt-execute([$message[receiver], $message[title], $message[content]]); $this-markSent($messageId, inserted into notice table); return true; }3.4 异步队列接入Redis Stream消费者组实战异步消息是这套系统的重头戏。我用的Redis Stream消费者组模式来承载消息队列而不是简单的LPUSH/BRPOP列表。原因有两个一是消费确认消息被消费者取走后需要一个XACK确认如果没有确认服务重启后消息可以从Pending列表里捞回来重新处理不会丢二是消费者组天然支持多个消费者分摊消息每个消费者读到的消息不会重复。这些是list模式做不到的。消息入队的代码在通知管理类里已经写了就是xAdd命令。消费者进程的代码如下?php declare(strict_types1); namespace App\Notification\Consumer; use App\Notification\Channel\ChannelFactory; use App\Notification\Channel\AbstractChannel; class StreamConsumer { private \Redis $redis; private \PDO $db; private string $group notification-consumers; private string $stream stream:message:send; private string $consumerName; public function __construct(\Redis $redis, \PDO $db, string $consumerName) { $this-redis $redis; $this-db $db; $this-consumerName $consumerName; } public function run(): void { $this-createGroupIfMissing(); while (true) { try { $messages $this-redis-xReadGroup( $this-group, $this-consumerName, [$this-stream ], 10, 2000 ); if (!$messages) { continue; } foreach ($messages as $stream $items) { foreach ($items as $msgId $data) { $this-handleMessage((int)$data[message_id], $data[msg_id]); $this-redis-xAck($this-stream, $this-group, [$msgId]); } } } catch (\Throwable $e) { error_log([Consumer][ . $this-consumerName . ] . $e-getMessage() . PHP_EOL, 3, /var/log/php-notify/consumer.log); sleep(5); } } } private function handleMessage(int $messageId, string $msgId): void { $stmt $this-db-prepare(SELECT * FROM message WHERE id ?); $stmt-execute([$messageId]); $message $stmt-fetch(\PDO::FETCH_ASSOC); if (!$message) { return; } // 超出重试次数直接丢弃 if ($message[retry_count] $message[max_retry]) { $this-db-prepare(UPDATE message SET status 2, channel_resp max retry exceeded WHERE id ?) -execute([$messageId]); return; } // 到期时间未到重新放回队列延迟重试场景 if (strtotime($message[next_retry_time]) time()) { return; } $channel ChannelFactory::getInstance()-getChannel($message[channel_type]); if ($channel instanceof AbstractChannel !$channel-send($messageId, $message)) { error_log([Consumer][ . $this-consumerName . ] send failed message_id . $messageId . PHP_EOL, 3, /var/log/php-notify/consumer.log); } } private function createGroupIfMissing(): void { try { $this-redis-xGroup(CREATE, $this-stream, $this-group, 0, true); } catch (\RedisException $e) { // 组已存在忽略错误 } } }消费者进程启动之后pending状态的消息会一直被捞出来重新尝试直到达到重试上限或者发送成功。这就是XACK确认机制的价值所在消息只要没被确认服务挂了重启之后重新消费不会出现“消息从Redis里消失但实际没发出去”的情况。部署消费者进程我用的是Supervisor配置好process数量后自动拉起。几个consumer实例并联消费同一组StreamRedis会做好消息分配不会重复处理。我跑下来发现开启多个消费者实例之后如果某个消费者处理某条消息特别慢它会一直占着那条消息去重试其他消费者会领别的消息。这样能避免单个慢消息阻塞整个队列吞吐。3.5 模板渲染与实时推送实现细节模板渲染这块我实现了一个TemplateRenderer类核心逻辑很简单但有一些细节值得说一说。渲染分两步先从消息模板表加载模板然后做变量替换。变量替换我用的正则表达式模式是/{{\s*(\w)\s*}}/把模板里的占位符替换成对应的参数值。因为要防止用户传过来带HTML标签的内容直接注入到邮件或页面里替换出来的值默认htmlspecialchars处理一遍如果模板里需要输出富文本用特定标记区分。?php declare(strict_types1); namespace App\Notification\Template; class TemplateRenderer { private \PDO $db; public function __construct(\PDO $db) { $this-db $db; } public function render(string $templateCode, array $params, int $channelType): array { $stmt $this-db-prepare( SELECT title_tpl, content_tpl, vars FROM message_template WHERE template_code ? AND channel_type ? AND status 1 ); $stmt-execute([$templateCode, $channelType]); $template $stmt-fetch(\PDO::FETCH_ASSOC); if (!$template) { throw new \RuntimeException(template [{$templateCode}] not found); } $vars json_decode($template[vars], true) ?? []; foreach ($vars as $var) { $name $var[name]; $required $var[required] ?? false; if ($required !isset($params[$name])) { throw new \RuntimeException(template [{$templateCode}] missing required param [{$name}]); } } $title $this-replace($template[title_tpl], $params); $content $this-replace($template[content_tpl], $params); return [$title, $content]; } private function replace(string $template, array $params): string { $result preg_replace_callback( /\{\{\s*(\w)\s*\}\}/, function ($matches) use ($params) { $key $matches[1]; $value $params[$key] ?? ; return htmlspecialchars((string)$value, ENT_QUOTES, UTF-8); }, $template ); return $result ?? $template; } }实时推送部分Workerman做WebSocket服务端客户端连接上来后维护一个在线连接池。当异步消息消费成功且渠道类型是站内信或需要实时推送时投递到一个WebSocket广播topic服务端把消息推给指定用户的在线连接。?php declare(strict_types1); namespace App\Notification\Push; use Workerman\Connection\TcpConnection; class PushHub { private static array $connections []; public static function addConnection(int $userId, TcpConnection $connection): void { $connection-userId $userId; self::$connections[$userId][] $connection; } public static function removeConnection(TcpConnection $connection): void { $userId $connection-userId ?? null; if ($userId null) { return; } if (!isset(self::$connections[$userId])) { return; } self::$connections[$userId] array_filter( self::$connections[$userId], fn(TcpConnection $conn) $conn ! $connection ); } public static function pushToUser(int $userId, array $data): void { $payload json_encode($data, JSON_UNESCAPED_UNICODE); if (!isset(self::$connections[$userId])) { return; } foreach (self::$connections[$userId] as $conn) { if ($conn-getStatus() TcpConnection::STATUS_ESTABLISHED) { $conn-send($payload); } } } }Workerman的events.php里注册onConnect、onClose、onMessage三个回调分别做连接登记、连接移除、心跳检测。心跳检测很重要我每30秒从客户端发一次ping服务端超过90秒没收到心跳就主动断开连接防止死连接占用内存。连接池是静态数组挂在内存里的单机部署没问题。如果要做多机负载均衡就得把这个连接池状态挪到Redis里每台机器启动一个本地消费者订阅Redis的channel把推送消息转发给本地连接。这块如果展开就是另一个话题了先给个方向。4. 常见问题与排查技巧实录4.1 SMTP发信失败超时与认证问题邮件渠道最容易遇到的问题就是SMTP超时和认证失败。我在生产环境遇到过一个诡异现象本地开发环境邮件秒发到了线上服务器就卡住最后超时。排查了半天发现是线上服务器的25端口被运营商封了而SMTP配置里没有指定SSL端口PHPMailer默认走了25。解决方案是把端口显式改为465并开启SMTPSecure问题立刻消失。这类问题的排查思路是先看error log里PHPMailer的ErrorInfo区分是连接超时还是认证失败。连接超时大概率是网络或端口问题可以用telnet smtp.xxx.com 465测一下通不通认证失败要检查用户名密码是否带上了域名后缀很多邮箱服务商要求用户名写完整邮箱地址而不是单纯的账号前缀。4.2 队列消息没有被消费Stream消费者组的神奇Bug我在上线异步队列之后遇到过一整个队列的消息堆积Redis里的消息明明有几千条消费者进程却一条也不处理。排查下来发现是创建消费者组时用了xGroup(CREATE, stream, group, 0, true)的最后一个参数true这是让组从流的开头开始消费。但后续xReadGroup的时候我读的是最新的消息中间一旦出现消费者崩溃Pending列表里会有大量未确认的旧消息而新的消费者实例又只读两边就对不上了。这个问题的解决办法是把xReadGroup的读取起点设为0而不是先把Pending历史消息处理完再继续读新消息。我在消费者代码里加了Pending检查$pending $this-redis-xPending($this-stream, $this-group); if (isset($pending[pending]) $pending[pending] 0) { $messages $this-redis-xReadGroup($this-group, $this-consumerName, [$this-stream 0], 10); // 处理Pending消息 }这里要注意不要死循环因为处理完的Pending是被XACK掉的总数会慢慢降下来直到归零再切回模式。4.3 模板变量转义与XSS风险模板渲染加了htmlspecialchars之后有一个副作用如果业务方确实需要传入一段富文本HTML作为消息内容比如运营人员编辑的活动详情就会被转义成纯文本显示不出来效果。我一开始一刀切全部转义结果运营反馈活动邮件里的图片链接全变成了纯文本。后来我增加了白名单机制模板定义里声明某个变量是safe_html类型渲染的时候不对它做转义但是后台会经过一份HTML过滤器只允许p、strong、a、img等少量标签其他的通通剥掉。这个过滤器虽然简单但能挡住绝大多数XSS攻击。如果是让业务方直接在后台填HTML那风险太大了不建议放开。4.4 时区问题导致定时消息错乱通知系统里有一个定时消息场景比如订单发货后24小时自动发送评价邀请。一开始定时任务直接用PHP的date(Y-m-d H:i:s)生成时间戳用的PHP时区是UTC数据库和前端页面又是北京时间结果定时消息全部提前了8小时发出用户凌晨收到消息体验非常差。这个问题就是时区没有统一。我的解决办法是数据库连接串里设置time_zone08:00PHP侧date_default_timezone_set(Asia/Shanghai)Redis侧的过期时间用时间戳而不是时间字符串。三处对齐之后再没出现过时区错乱。4.5 连接超时与重试风暴通知系统最怕的重试风暴是某个渠道彻底挂了比如短信服务商宕机此时所有消费线程都在疯狂重试失败的短信第三方接口被打得更惨同时Redis队列不断堆积CPU和内存持续飙升。我后来加了一个熔断机制。在ChannelFactory里加了一个简单的计数器如果某个渠道在1分钟内失败超过50次直接熔断该渠道10分钟期间新消息直接标记为发送失败并记录reason不执行实际发送。这样第三方服务恢复之前系统不会再去打它避免了雪上加霜。熔断状态存在Redis里多个消费者实例共享不会出现每个消费者各自为政的情况。我在熔断逻辑里还做了一个优雅降级对于有多个渠道的业务类型优先尝试降级渠道。比如原本走短信的验证码短信通道熔断之后如果该业务允许就降级到邮件渠道发送保证用户消息不丢失。5. 扩展方向从能用走向好用5.1 更灵活的多租户与业务线隔离如果你们公司有多个业务线共用一个用户体系消息通知系统也需要做多租户隔离。最简单的方式是给消息表加一个tenant_id字段配合路由表做租户维度隔离。模板表同样要加租户维度不同租户即便是同一个模板编码文案和渠道配置也可以不同。我实际认证过这种方案的可行性只需要在Query层统一加上tenant条件其他逻辑保持不变。如果你需要用同一套底层同时服务内部运营系统、C端用户系统和外部客户系统这个改造非常值得做。5.2 频率限制与用户偏好设置消息发太多会被用户投诉甚至导致短信通道被运营商限制。所以频率限制是这个系统的必备能力。我在发送入口加了一个RateLimiter基于Redis的INCR和EXPIRE实现支持三个维度同一用户1分钟内最多发N条、同一手机号1小时内最多发M条、同一业务类型每天最多发K条。超过限制直接拒绝并记录拦截日志。用户偏好设置也很重要。很多用户不想要营销邮件只接收订单通知。我在用户偏好表里记录每个用户每个业务类型的接收开关发送前首先检查用户是否订阅了该类型的消息。这个东西在冷启动阶段可以先不做全局的偏好界面但至少要在后端留一张表配合业务方在注册流程或设置页埋入口收集偏好。5.3 消息撤回、重推与全链路追踪状态机里的“已撤回”状态不是摆设。某些场景比如用户取消订单后之前发的订单确认通知就没意义了需要能撤回。短信和邮件一旦发出去了技术上很难真正撤回只能在文案上做补偿。但站内信和WebSocket推送是能撤回的。我给通知系统加了一个撤回接口调用方传msg_id系统会把站内信状态改成已撤回并且通过WebSocket推送一条撤回指令客户端收到后自动把消息移除。全链路追踪我是这样做的在消息表里不只有渠道响应还有一个trace_id从业务方调用入口生成贯穿消息入队、消费、渠道发送、回执整个链路。在日志系统里搜索trace_id就能看到这条消息完整的时间线和每一步的状态变更。这个对排查用户投诉“我没收到消息”非常有帮助每条消息走到哪一步一目了然不用靠猜。5.4 个人使用体会整套系统从设计到上线我最大的收获不是代码本身而是对“异步”这两个字的理解变深了。消息通知表面上是发一条消息背后牵扯的是你如何理解业务边界、如何应对失败、如何让系统在异常情况下保持体面。我做这套系统之前觉得消息发不出去大不了用户没收到做完之后才意识到一个通知系统的可靠性直接影响用户对产品信任度的判断——验证码丢了可以重发但支付回执丢了用户会怀疑自己钱付没付成功。如果你也在计划自建消息通知系统我给的建议是别一上来就追求大而全先把邮件、短信、站内信三个渠道跑通用队列把异步链路搭起来把重试、幂等、日志这三件套做扎实就已经能解决你80%的问题了。实时推送、多租户、频率控制这些等业务真需要了再加也不迟。关键是底子要打对状态机、渠道抽象、模板引擎这三个核心一旦设计合理后续扩展就是不断加实现类的事而不是推倒重来。