PHP 怎么异步沟通

wen PHP项目 3

本文目录导读:

PHP 怎么异步沟通

  1. 消息队列(最常用)
  2. 进程异步(pcntl_fork)
  3. 使用 Swoole(高性能方案)
  4. 使用 HTTP 异步请求
  5. 使用 Guzzle 异步客户端
  6. 使用 ReactPHP(事件驱动)
  7. 最佳实践建议

在 PHP 中实现异步通信主要有以下几种方式,从简单到复杂:

消息队列(最常用)

RabbitMQ

// 生产者(发送消息)
require_once __DIR__ . '/vendor/autoload.php';
use PhpAmqpLib\Connection\AMQPStreamConnection;
use PhpAmqpLib\Message\AMQPMessage;
$connection = new AMQPStreamConnection('localhost', 5672, 'guest', 'guest');
$channel = $connection->channel();
$channel->queue_declare('task_queue', false, true, false, false);
$data = json_encode(['task' => 'send_email', 'to' => 'user@example.com']);
$msg = new AMQPMessage($data, ['delivery_mode' => 2]);
$channel->basic_publish($msg, '', 'task_queue');
echo " [x] Sent message\n";
$channel->close();
$connection->close();
// 消费者(处理消息)
require_once __DIR__ . '/vendor/autoload.php';
use PhpAmqpLib\Connection\AMQPStreamConnection;
$connection = new AMQPStreamConnection('localhost', 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(1);
    echo " [x] Done\n";
    $msg->ack();
};
$channel->basic_qos(null, 1, null);
$channel->basic_consume('task_queue', '', false, false, false, false, $callback);
while ($channel->is_consuming()) {
    $channel->wait();
}

Redis 队列

// 生产者
$redis = new Redis();
$redis->connect('127.0.0.1', 6379);
$task = ['type' => 'email', 'data' => ['to' => 'user@example.com']];
$redis->lpush('task_queue', json_encode($task));
// 消费者(配合 Supervisor 运行)
$redis = new Redis();
$redis->connect('127.0.0.1', 6379);
while (true) {
    $task = $redis->brpop('task_queue', 0);
    if ($task) {
        $data = json_decode($task[1], true);
        // 处理任务
        processTask($data);
    }
}

进程异步(pcntl_fork)

// 异步执行不需要等待结果的任务
$pid = pcntl_fork();
if ($pid == -1) {
    // 创建进程失败
    die('Could not fork');
} elseif ($pid) {
    // 父进程 - 继续执行,不等待子进程
    echo "Parent process continues...\n";
    // 返回响应给用户
} else {
    // 子进程 - 处理耗时任务
    sleep(5);
    file_put_contents('/tmp/task_result.txt', 'Task completed');
    exit(0);
}

使用 Swoole(高性能方案)

// 异步任务示例
$server = new Swoole\Server('127.0.0.1', 9501);
$server->set([
    'task_worker_num' => 4, // 任务进程数
]);
$server->on('Receive', function ($server, $fd, $reactor_id, $data) {
    // 投递异步任务
    $task_id = $server->task($data);
    echo "Async task created: $task_id\n";
});
$server->on('Task', function ($server, $task_id, $reactor_id, $data) {
    // 处理任务
    echo "Processing task: $data\n";
    $result = processTask($data);
    $server->finish($result);
});
$server->on('Finish', function ($server, $task_id, $data) {
    echo "Task $task_id completed: $data\n";
});
$server->start();

使用 HTTP 异步请求

// 使用 cURL 进行异步 HTTP 请求
function sendAsyncRequest($url, $data) {
    $ch = curl_init();
    curl_setopt($ch, CURLOPT_URL, $url);
    curl_setopt($ch, CURLOPT_POST, 1);
    curl_setopt($ch, CURLOPT_POSTFIELDS, http_build_query($data));
    curl_setopt($ch, CURLOPT_RETURNTRANSFER, true);
    curl_setopt($ch, CURLOPT_TIMEOUT, 1); // 超时时间设短
    curl_setopt($ch, CURLOPT_NOSIGNAL, 1);
    // 不等待响应
    curl_exec($ch);
    curl_close($ch);
}
// 使用
sendAsyncRequest('http://api.example.com/task', ['data' => 'test']);

使用 Guzzle 异步客户端

require_once __DIR__ . '/vendor/autoload.php';
use GuzzleHttp\Client;
use GuzzleHttp\Promise;
$client = new Client();
// 创建多个异步请求
$promises = [
    'api1' => $client->getAsync('http://api1.example.com/data'),
    'api2' => $client->getAsync('http://api2.example.com/data'),
    'api3' => $client->getAsync('http://api3.example.com/data'),
];
// 等待所有请求完成
$results = Promise\Utils::settle($promises)->wait();
foreach ($results as $key => $result) {
    if ($result['state'] === 'fulfilled') {
        echo "$key: " . $result['value']->getBody() . "\n";
    }
}

使用 ReactPHP(事件驱动)

require_once __DIR__ . '/vendor/autoload.php';
use React\EventLoop\Factory;
use React\ChildProcess\Process;
$loop = Factory::create();
// 执行异步命令
$process = new Process('php long_task.php');
$process->start($loop);
$process->on('exit', function ($exitCode) {
    echo "Process exited with code: $exitCode\n";
});
$loop->run();

最佳实践建议

根据场景选择合适方案:

场景 推荐方案
简单任务(发邮件、通知) Redis/MQ 队列
大量消息、可靠性要求高 RabbitMQ/Kafka
需要高性能、低延迟 Swoole
微服务间通信 HTTP+队列
简单演示/快速实现 pcntl_fork

生产环境示例架构:

// 1. API 接口层 - 快速返回响应
public function createOrder()
{
    $orderId = $this->createOrderInDB();
    // 投递异步任务
    $this->queueService->push('send_notification', [
        'order_id' => $orderId,
        'action' => 'order_created'
    ]);
    // 立即返回
    return response()->json(['order_id' => $orderId]);
}
// 2. 消费者进程(独立部署)
public function consumeNotification()
{
    while ($task = $this->queueService->pop('send_notification')) {
        $this->sendEmail($task);
        $this->sendSMS($task);
    }
}

注意事项:

  • 确保队列服务配置了持久化,防止任务丢失
  • 使用 Supervisor 或 systemd 管理消费者进程
  • 设置合理的超时和重试机制
  • 监控队列长度和处理速度
  • 考虑任务失败的处理策略

选择哪种方案取决于你的具体需求:任务复杂度、实时性要求、系统规模和预算等。

抱歉,评论功能暂时关闭!