PHP结合WebSocket/SSE实现CRM实时消息推送方案咨询
CRM实时消息推送需求
我正在为CRM开发实时消息功能,需要实现无需频繁查询数据库(避免增加服务器与数据库负载)的前端通知推送方案。
目前已通过PHP搭建webhook,将短信消息插入MySQL数据库,代码如下:
if (isset($_POST['From']) && isset($_POST['Body'])) { $stmt = $pdo->prepare('INSERT INTO sms_messages (`number`, message) VALUES (:number, :message)'); $stmt->execute([ 'number' => $_POST['From'], 'message' => $_POST['Body'], ]); // 此处需要实现向客户端发送通知的功能 }
我希望在不使用轮询的情况下,通过JavaScript将这些通知直接推送到客户端前端。我有使用Ratchet处理PHP WebSocket的经验,但不确定如何将上述PHP事件与WebSocket集成以实现实时通知客户端。
基础WebSocket服务器代码:
require 'vendor/autoload.php'; use Ratchet\Server\IoServer; use Ratchet\Http\HttpServer; use Ratchet\WebSocket\WsServer; use Ratchet\MessageComponentInterface; use Ratchet\ConnectionInterface; class Chat implements MessageComponentInterface { public function onOpen(ConnectionInterface $conn) { echo "New connection! ({$conn->resourceId})\n"; } public function onMessage(ConnectionInterface $from, $msg) { echo sprintf('New message from %d: %s'."\n", $from->resourceId, $msg); } public function onClose(ConnectionInterface $conn) { echo "Connection {$conn->resourceId} has disconnected\n"; } public function onError(ConnectionInterface $conn, \Exception $e) { echo "An error has occurred: {$e->getMessage()}\n"; $conn->close(); } } $server = IoServer::factory(new HttpServer(new WsServer(new Chat())), 8080); $server->run();
对应的客户端JavaScript代码:
var conn = new WebSocket('ws://localhost:8080'); conn.onopen = function(e) { console.log("Connection established!"); conn.send('Hello, World!'); }; conn.onmessage = function(e) { console.log(e.data); }; conn.onerror = function(e) { console.error("Connection failed!", e); }; function sendMessage(message) { if (conn.readyState === WebSocket.OPEN) { conn.send(message); } else { console.log("WebSocket is not open. ReadyState: " + conn.readyState); } } setTimeout(function() { sendMessage('Hello again!'); }, 1000);
我也在考虑Server-Sent Events(SSE)作为替代方案,但同样不确定如何在不周期性查询数据库的情况下搭建。恳请指导如何通过WebSocket、SSE或其他合适的方案,高效实现PHP到JavaScript的实时事件推送。
解决方案
方案一:WebSocket(Ratchet 集成)
通过Redis Pub/Sub实现webhook与Ratchet服务器的通信,避免直接进程间通信的复杂性,确保消息可靠传递。
步骤1:安装依赖
composer require cboden/ratchet predis/predis
步骤2:修改Ratchet服务器,支持Redis订阅
更新WebSocket组件,让它监听Redis频道,收到消息后推送给所有在线客户端:
require 'vendor/autoload.php'; use Ratchet\Server\IoServer; use Ratchet\Http\HttpServer; use Ratchet\WebSocket\WsServer; use Ratchet\MessageComponentInterface; use Ratchet\ConnectionInterface; use Predis\Client; class Chat implements MessageComponentInterface { protected $clients; protected $redis; public function __construct() { $this->clients = new \SplObjectStorage; $this->redis = new Client(); // 异步订阅短信通知频道 $this->redis->pubSubLoop(['subscribe' => 'sms_notifications'], function($message) { if ($message->payload) { // 推送消息给所有在线客户端 foreach ($this->clients as $client) { $client->send($message->payload); } } }); } public function onOpen(ConnectionInterface $conn) { $this->clients->attach($conn); echo "New connection! ({$conn->resourceId})\n"; } public function onMessage(ConnectionInterface $from, $msg) { echo sprintf('New message from %d: %s'."\n", $from->resourceId, $msg); } public function onClose(ConnectionInterface $conn) { $this->clients->detach($conn); echo "Connection {$conn->resourceId} has disconnected\n"; } public function onError(ConnectionInterface $conn, \Exception $e) { echo "An error has occurred: {$e->getMessage()}\n"; $conn->close(); } } $server = IoServer::factory(new HttpServer(new WsServer(new Chat())), 8080); $server->run();
步骤3:更新webhook,推送消息到Redis
在插入数据库后,将消息以JSON格式发送到Redis频道:
if (isset($_POST['From']) && isset($_POST['Body'])) { $stmt = $pdo->prepare('INSERT INTO sms_messages (`number`, message) VALUES (:number, :message)'); $stmt->execute([ 'number' => $_POST['From'], 'message' => $_POST['Body'], ]); // 发送通知到Redis频道 $redis = new Predis\Client(); $notification = json_encode([ 'number' => $_POST['From'], 'message' => $_POST['Body'], 'timestamp' => date('Y-m-d H:i:s') ]); $redis->publish('sms_notifications', $notification); }
步骤4:客户端处理推送消息
修改前端代码,接收WebSocket消息并展示通知:
var conn = new WebSocket('ws://localhost:8080'); // 请求浏览器通知权限 if (Notification.permission !== 'denied') { Notification.requestPermission(); } conn.onopen = function(e) { console.log("Connection established!"); }; conn.onmessage = function(e) { const notification = JSON.parse(e.data); // 弹出系统通知 if (Notification.permission === 'granted') { new Notification('新短信消息', { body: `来自${notification.number}: ${notification.message}`, icon: '/path/to/notification-icon.png' }); } // 更新页面消息列表 const messageItem = document.createElement('div'); messageItem.textContent = `${notification.timestamp} - ${notification.number}: ${notification.message}`; document.getElementById('message-list').appendChild(messageItem); }; conn.onerror = function(e) { console.error("Connection failed!", e); };
方案二:Server-Sent Events(SSE)
如果不需要双向通信,SSE是更轻量的选择,无需单独维护WebSocket服务器,仅通过PHP长连接实现推送。
步骤1:创建SSE推送脚本(sse.php)
该脚本保持长连接,监听Redis频道并推送消息:
<?php header('Content-Type: text/event-stream'); header('Cache-Control: no-cache'); header('Connection: keep-alive'); require 'vendor/autoload.php'; use Predis\Client; $redis = new Client(); $pubSub = $redis->pubSubLoop(['subscribe' => 'sms_notifications']); foreach ($pubSub as $message) { if ($message->payload) { // 按SSE格式输出消息 echo "data: {$message->payload}\n\n"; ob_flush(); flush(); } }
步骤2:webhook代码与WebSocket方案一致
同样在插入数据库后向Redis频道发送JSON消息。
步骤3:客户端监听SSE
if (typeof EventSource !== 'undefined') { const eventSource = new EventSource('sse.php'); // 请求通知权限 if (Notification.permission !== 'denied') { Notification.requestPermission(); } eventSource.onmessage = function(e) { const notification = JSON.parse(e.data); // 展示系统通知 if (Notification.permission === 'granted') { new Notification('新短信消息', { body: `来自${notification.number}: ${notification.message}`, icon: '/path/to/notification-icon.png' }); } // 更新页面消息列表 const messageItem = document.createElement('div'); messageItem.textContent = `${notification.timestamp} - ${notification.number}: ${notification.message}`; document.getElementById('message-list').appendChild(messageItem); }; eventSource.onerror = function(e) { console.error('SSE连接出错:', e); eventSource.close(); // 5秒后尝试重连 setTimeout(() => window.location.reload(), 5000); }; } else { alert('你的浏览器不支持SSE,请切换至WebSocket方案'); }
方案对比
- WebSocket:支持双向通信,适合需要客户端主动交互的场景(如标记消息已读),但需单独维护WebSocket服务器。
- SSE:仅单向推送,实现简单,无需额外服务器,兼容性略差(现代主流浏览器均支持),仅支持文本数据。
内容的提问来源于stack exchange,提问作者tony
相关产品推荐
相关产品推荐

