【问题标题】:Swoole with RabbitMQSwoole 与 RabbitMQ
【发布时间】:2018-08-19 23:12:49
【问题描述】:

我正在尝试使用 websockets 将一些数据从 php 应用程序发送到用户的浏览器。因此,我决定将Swoole 与 RabbitMQ 结合使用。

这是我第一次使用 websockets,在阅读了一些关于 Socket.IO、Ratchet 等的帖子之后,我决定停止使用 Swoole,因为它是用 C 语言编写的,并且可以方便地与 php 一起使用。

这就是我理解使用 websocket 启用数据传输的想法的方式: 1) 在 CLI 中启动 RabbitMQ worker 和 Swoole 服务器 2)php应用程序向RabbitMQ发送数据 3) RabbitMQ 向 worker 发送带有数据的消息 4) Worker 接收到带有数据的消息 + 与 Swoole 套接字服务器建立套接字连接。 5) Swoole 服务器向所有连接广播数据

问题是如何绑定Swoole socket server和RabbitMQ?或者如何让 RabbitMQ 与 Swoole 建立连接并向其发送数据?

代码如下:

Swoole 服务器 (swoole_sever.php)

$server = new \swoole_websocket_server("0.0.0.0", 2345, SWOOLE_BASE);

$server->on('open', function(\Swoole\Websocket\Server $server, $req)
{
    echo "connection open: {$req->fd}\n";
});

$server->on('message', function($server, \Swoole\Websocket\Frame $frame)
{
    echo "received message: {$frame->data}\n";
    $server->push($frame->fd, json_encode(["hello", "world"]));
});

$server->on('close', function($server, $fd)
{
    echo "connection close: {$fd}\n";
});

$server->start();

Worker 从 RabbitMQ 接收消息,然后连接到 Swoole 并通过套接字连接(worker.php)广播消息

$connection = new AMQPStreamConnection('0.0.0.0', 5672, 'guest', 'guest');
$channel = $connection->channel();

$channel->queue_declare('task_queue', false, true, false, false);

echo ' [*] Waiting for messages. To exit press CTRL+C', "\n";

$callback = function($msg){
    echo " [x] Received ", $msg->body, "\n";
    sleep(substr_count($msg->body, '.'));
    echo " [x] Done", "\n";
    $msg->delivery_info['channel']->basic_ack($msg->delivery_info['delivery_tag']);


    // Here I'm trying to make connection to Swoole server and sernd data
    $cli = new \swoole_http_client('0.0.0.0', 2345);

    $cli->on('message', function ($_cli, $frame) {
        var_dump($frame);
    });

    $cli->upgrade('/', function($cli)
    {
        $cli->push('This is the message to send to Swoole server');
        $cli->close();
    });
};

$channel->basic_qos(null, 1, null);
$channel->basic_consume('task_queue', '', false, false, false, false, $callback);

while(count($channel->callbacks)) {
    $channel->wait();
}

$channel->close();
$connection->close();

消息将发送到 RabbitMQ (new_task.php) 的新任务:

$connection = new AMQPStreamConnection('0.0.0.0', 5672, 'guest', 'guest');
$channel = $connection->channel();

$channel->queue_declare('task_queue', false, true, false, false);

$data = implode(' ', array_slice($argv, 1));
if(empty($data)) $data = "Hello World!";
$msg = new AMQPMessage($data,
    array('delivery_mode' => AMQPMessage::DELIVERY_MODE_PERSISTENT)
);

$channel->basic_publish($msg, '', 'task_queue');

echo " [x] Sent ", $data, "\n";

$channel->close();
$connection->close();

在启动 swoole 服务器和 worker 后,我正在从命令行触发 new_task.php:

php new_task.php

在运行 RabbitMQ Worker (worker.php) 的命令行提示符中,我可以看到一条消息已传递给该 worker(“[x] Received Hello World!”消息正在出现)。

但是在运行 Swoole 服务器的命令行提示符中什么也没有发生。

所以问题是: 1)这种方法的想法正确吗? 2) 我做错了什么?

【问题讨论】:

  • 我从来没有用过这个swoole,但我用过rabbit,只是粗略地看了一下他们的文档,看起来你错过了github.com/swoole/swoole-src/blob/master/examples/…的一些东西,比如$cli->setData(...)$cli->execute(...)
  • 我建议把事情分解一下,例如编写一个命令行 php 脚本来测试向这个 swoole 发送一个简单的消息,然后一旦你开始工作,集成 RabbitMq 的东西,这样你就可以隔离问题。
  • 我已经这样做了。我可以使用 Swoole 进行套接字连接并以两种方式发送数据。我还可以使用 RabbitMQ 向 Worker 发送消息。所以,我唯一的问题是将 RabbitMQ 与 Swoole 一起使用。
  • 确保你的工人有适当的权限,这是我唯一能想到的。也许您可以以 root 身份连接到swoole,但工作人员不能。如果代码分开工作应该没有区别,所以它必须是环境。
  • 我注意到奇怪的事情。我正在做的步骤: 1)启动 Swoole 服务器 2)启动 Worker 3)执行 new_task.php 4)在 worker.php 中得到消息 5)在 swoole_server.php 中没有看到 6)手动停止 worker.php 在终端 7 中运行) 在 swoole_server.php 的终端窗口中,我收到消息:“connection close: 1”。

标签: php sockets websocket rabbitmq swoole


【解决方案1】:

在收到消息时触发的回调(在 worker.php 中)中,您使用的是仅异步的 swoole_http_client。这似乎导致代码永远不会完全执行,因为回调函数在触发异步代码之前返回。

做同样事情的同步方法将解决问题。这是一个简单的例子:

$client = new WebSocketClient('0.0.0.0', 2345);
$client->connect();
$client->send('This is the message to send to Swoole server');
$recv = $client->recv();
print_r($recv);
$client->close();

查看WebSocketClient 类和github 的示例用法。

您也可以将其包装在 coroutine 中,如下所示:

go(function () {
    $client = new WebSocketClient('0.0.0.0', 2345);
    $client->connect();
    $client->send('This is the message to send to Swoole server');
    $recv = $client->recv();
    print_r($recv);
    $client->close();
});

【讨论】:

  • 非常感谢!刚刚测试了第一种方法,它按预期工作!
  • 很高兴听到。祝你的项目好运。 :)
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2022-10-17
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2017-12-05
相关资源
最近更新 更多