在 think-swoole 框架中融入 Custom Process(自定义进程)
在 think-swoole 框架中融入 Custom Process(自定义进程) 需要特别注意生命周期。因为 think-swoole接管了 Swoole\Server 的创建与启动,我们不能直接 $server->addProcess(),而是必须在 swoole.init 事件中,将自定义进程实例挂载进去。
同时,由于 think-swoole 使用了 Sandbox(沙箱) 隔离机制,为了确保自定义进程内能安全调用 ThinkPHP 的容器和服务,我们需要在进程初始化时对应用进行基本的引导。
以下是适配你现有代码的完整改造方案:
1. 新建自定义广播进程类
新建一个专门处理广播的进程类,便于维护和管理:app/common/swoole/process/BroadcastProcess.php
php
<?php
namespace app\common\swoole\process;
use Swoole\Process;
use Swoole\Server as SwooleServer;
use think\App;
class BroadcastProcess
{
protected $app;
protected $server;
public function __construct(App $app, SwooleServer $server)
{
$this->app = $app;
$this->server = $server;
}
public function create()
{
// 创建 Swoole 自定义进程
// 参数 2 为 false(不重定向标准输入输出),参数 3 为 1(启用支持异步非阻塞的 Unix Socket 管道)
return new Process([$this, 'run'], false, 1);
}
public function run(Process $process)
{
// 1. 设置自定义进程名称
swoole_set_process_name('swoole: broadcast process');
// 2. 异步监听管道数据
swoole_event_add($process->pipe, function ($pipe) use ($process) {
$jsonPacket = $process->read();
if (empty($jsonPacket)) {
return;
}
$packet = json_decode($jsonPacket, true);
$message = $packet['message'] ?? '';
$excludeFd = $packet['exclude_fd'] ?? null;
if (empty($message)) {
return;
}
// 3. 核心:遍历连接执行广播,彻底脱离 Worker 进程
foreach ($this->server->connections as $fd) {
// 确保是合法的 WebSocket 连接,且排除特定 FD
if ($this->server->isEstablished($fd) && $fd != $excludeFd) {
$this->server->send($fd, $message);
}
}
});
}
}
请谨慎使用此类代码。
2. 修改你的 SwooleService.php
在 swoole.init 事件中,实例化这个自定义进程,并通过 $server->addProcess() 注入到主服务。同时,我们将进程句柄保存到容器中,方便在 WebSocket 业务中直接获取并写入管道。
php
<?php
namespace app\common\swoole;
use think\Service;
use think\swoole\Manager;
use Swoole\Server as SwooleServer;
use app\common\swoole\process\BroadcastProcess;
class SwooleService extends Service
{
public function boot()
{
// 确保是在 CLI 模式下且引入了 think-swoole
if ($this->app->runningInConsole() && class_exists(Manager::class))
{
// 监听 rpc 的 receive 回调事件
$this->app->event->listen('swoole.init', function ()
{
/** @var Manager $manager */
$manager = $this->app->make(Manager::class);
$server = $manager->getServer();
// ==================== 【新增:注入自定义广播进程】 ====================
$broadcastService = new BroadcastProcess($this->app, $server);
$broadcastProcess = $broadcastService->create();
// 将进程实例单例绑定到 ThinkPHP 容器,方便在任意地方通过 app('swoole.broadcast.process') 投递广播
$this->app->bind('swoole.broadcast.process', $broadcastProcess);
// 必须在 Server start 之前加入进程
$server->addProcess($broadcastProcess);
// ===================================================================
// 获取配置
$config = $this->app->config->get('swoole.rpc_server');
// 动态添加 RPC 监听端口 (TCP 协议)
$rpcPort = $server->addListener($config['host'], $config['port'], SWOOLE_SOCK_TCP);
// 设置 RPC 端口的独立参数
$rpcPort->set([
'open_eof_check' => true,
'package_eof' => "\n",
'open_eof_split' => true,
'open_websocket_protocol' => false, // 禁用WebSocket协议
'max_wait_time' => 5, // Worker 退出最大等待时间(秒)
]);
$server->on('start', function ($server)
{
$config = $this->app->config->get('swoole.rpc_server');
swoole_set_process_name($config['process_name']);
});
$server->on('managerStart', function ($server)
{
swoole_set_process_name('swoole: manager');
});
// 获取 Sandbox 实例
$sandbox = $this->app->make(\think\swoole\Sandbox::class);
// 绑定 RPC 监听事件
$rpcPort->on('receive', function (SwooleServer $serv, int $fd, int $reactorId, string $data) use ($sandbox)
{
// 此时 RPC 运行在独立的协程中,且 Worker 进程绝不会被 WebSocket 广播阻塞
go(function () use ($serv, $fd, $reactorId, $data, $sandbox)
{
$sandbox->run(function () use ($serv, $fd, $reactorId, $data, $sandbox)
{
$rpcService = $sandbox->getApplication()->make(\app\common\swoole\rpc\RpcService::class);
$rpcService->handleRpcRequest($serv, $fd, $reactorId, $data);
});
});
});
// WebSocket 握手
$server->on('handshake', function ($request, $response) use ($sandbox)
{
go(function () use ($request, $response, $sandbox)
{
$sandbox->run(function () use ($request, $response, $sandbox)
{
$WebSocketService = $sandbox->getApplication()->make(\app\common\swoole\websocket\Service::class);
$WebSocketService->onHandshake($request, $response);
});
});
});
// WebSocket 协议
$server->on('request', function ($request, $response) use ($sandbox)
{
go(function () use ($request, $response, $sandbox)
{
$sandbox->run(function () use ($request, $response, $sandbox)
{
$WebSocketService = $sandbox->getApplication()->make(\app\common\swoole\websocket\Service::class);
$WebSocketService->handleHttpRequest($request, $response);
});
});
});
});
}
}
}
请谨慎使用此类代码。
3. 如何在 WebSocket 业务逻辑中调用?
现在,当你在 \app\common\swoole\websocket\Service 或任何地方需要触发大规模广播时,千万不要直接循环 send,而是通过容器获取进程句柄,把任务丢进管道:
php
// 在你的 WebSocket 消息处理函数中:
public function broadcastToAll($message, $excludeFd = null)
{
$payload = json_encode([
'message' => $message,
'exclude_fd' => $excludeFd
]);
/** @var \Swoole\Process $process */
// 从容器中取出我们在初始化时绑定的自定义进程
$process = app('swoole.broadcast.process');
// 写入管道,数据瞬间移交给自定义进程
// Worker 进程在此处完全不留痕迹,立刻返回并准备迎接下一个 RPC 请求
$process->write($payload);
}
请谨慎使用此类代码。
💡 架构升级后的效果
RPC 请求零延迟:无论 WebSocket 广播在自定义进程里遍历了多少万个连接、消耗了多少 CPU,主 Worker 进程的 Event Loop 始终在轻量空转,RPC 协程随调随到。
优雅契合框架:通过将其单例绑定到 TP 的
Container容器,你的其他控制器、事件监听器和沙箱内部业务,都能非常方便地通过app('swoole.broadcast.process')调度广播能力。