1.php 消费端

<?php
/*这是workerman+redis异步队列处理消费任务的*/
require_once __DIR__ . '/../vendor/autoload.php';
use Workerman\Worker;
use Workerman\Timer;
//创建一个进程
$worker = new Worker();
//设置进程数量
$worker->count = 4; // 4个进程并处理消费队列
//进程名称
$worker->name = 'cyg';
//进程启动时
$worker->onWorkerStart = function($worker) {
    echo "队列消费者启动"."<br />";
    //创建redis 每一个进程独立链接redis
    $redis = new Redis();
    $redis->connect('127.0.0.1', 6379);
    // 队列名称
    $queueName = 'queue';
// 定时器,每次隔 0.5 秒从队列取一个任务处理
    Timer::add(0.5, function() use ($redis, $queueName, $worker) {
        // 阻塞式:brPop 或 非阻塞式:lPop)
		//像是子弹弹夹那样从队列左边取出一个任务
        $taskJson = $redis->lPop($queueName);

        if ($taskJson) {
            $taskData = json_decode($taskJson, true);

            echo "进程 {$worker->id}开始处理任务 #{$taskData['id']}\n";

            // 模拟处理任务
            processTask($taskData, $worker->id);

            echo "进程{$worker->id}任务 #{$taskData['id']} 处理完成\n";
        } else {
           return false;
        }
    });
};

/**
 * 处理任务的业务逻辑,模拟的
 */
function processTask($taskData, $workerId) {
    $processTime = rand(1, 3); // 随机1-3秒
    sleep($processTime);

  

    echo "进程 {$workerId} 用户ID: {$taskData['user_id']}, 金额: {$taskData['amount']}, 耗时: {$processTime}秒\n";

  
    return true;
}

// 运行所有 Worker 进程
Worker::runAll();

2.php 生产端

<?php
/*像队列子弹夹推送子弹(任务)到里面*/
require_once __DIR__ . '/../vendor/autoload.php';
//连接redis
$redis = new Redis();
$redis->connect('127.0.0.1', 6379);
// 队列名称
$queueName = 'queue';
for ($i = 1; $i <= 100; $i++) {
    // 模拟任务数据
    $taskData = [
        'id' => $i,
        'user_id' => rand(1000, 9999),
        'data' => "这是第{$i}条需要处理的数据",
        'amount' => rand(100, 10000) / 100,
        'created_at' => date('Y-m-d H:i:s'),
    ];

    
	//推入数据到子弹夹中(队列)。从右边推入Redis 队列
    $redis->rPush($queueName, json_encode($taskData));
	echo "已推送任务 #{$i}\n";
 // 模拟任务产生的时间间隔(可选)
    usleep(50000); // 0.05秒
}

echo "完了 100 个任务到队列: {$queueName}\n";
echo "当前队列长度: " . $redis->lLen($queueName) . "\n";

$redis->close();

生产者快速推送任务到队列,消费者异步处理
Workerman 启动多个进程同时消费队列,提高处理效率

流程:
用户请求–生产者2.php–redis队列 --消费者1.php→ 处理完成()四个进程并发处理

更多推荐