本文目录导读:

在 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 管理消费者进程
- 设置合理的超时和重试机制
- 监控队列长度和处理速度
- 考虑任务失败的处理策略
选择哪种方案取决于你的具体需求:任务复杂度、实时性要求、系统规模和预算等。