You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.29 10:53:13