PHP项目多进程并发如何代码实现

wen PHP项目 25

本文目录导读:

PHP项目多进程并发如何代码实现

  1. PCNTL扩展(最常用)
  2. 使用Process类(Symfony Process组件)
  3. 并发执行命令(更高级的用法)
  4. 使用共享内存通信
  5. 使用队列系统(高级方案)
  6. 注意事项

在PHP中实现多进程并发主要有以下几种方式:

PCNTL扩展(最常用)

基础示例

<?php
// 确保使用CLI模式
if (php_sapi_name() !== 'cli') {
    die("This script must be run from CLI");
}
$tasks = [
    ['name' => 'Task 1', 'time' => 3],
    ['name' => 'Task 2', 'time' => 5],
    ['name' => 'Task 3', 'time' => 2],
];
$processes = [];
foreach ($tasks as $task) {
    $pid = pcntl_fork();
    if ($pid == -1) {
        // 创建进程失败
        die("Could not fork");
    } elseif ($pid) {
        // 父进程
        $processes[] = $pid;
        echo "Started process: $pid for {$task['name']}\n";
    } else {
        // 子进程
        executeTask($task);
        exit(0); // 子进程结束
    }
}
// 父进程等待所有子进程完成
foreach ($processes as $pid) {
    pcntl_waitpid($pid, $status);
    echo "Process $pid completed with status: $status\n";
}
function executeTask($task) {
    echo "[" . getmypid() . "] Starting {$task['name']}\n";
    sleep($task['time']);
    echo "[" . getmypid() . "] Finished {$task['name']}\n";
}

带信号处理的完整示例

<?php
declare(ticks = 1);
// 信号处理
function signalHandler($signo) {
    switch ($signo) {
        case SIGTERM:
            echo "Received SIGTERM\n";
            exit;
        case SIGCHLD:
            echo "Child process terminated\n";
            // 回收子进程
            while (($pid = pcntl_waitpid(-1, $status, WNOHANG)) > 0) {
                echo "Child $pid exited\n";
            }
            break;
    }
}
// 注册信号处理器
pcntl_signal(SIGTERM, 'signalHandler');
pcntl_signal(SIGCHLD, 'signalHandler');
// 创建子进程
$maxProcesses = 3;
$currentProcesses = 0;
$tasks = range(1, 10);
foreach ($tasks as $taskId) {
    while ($currentProcesses >= $maxProcesses) {
        // 等待一个子进程结束
        $pid = pcntl_wait($status, WUNTRACED);
        if ($pid > 0) {
            $currentProcesses--;
        }
    }
    $pid = pcntl_fork();
    if ($pid == -1) {
        die("Could not fork");
    } elseif ($pid) {
        $currentProcesses++;
        echo "Started process $pid for task $taskId\n";
    } else {
        // 子进程执行任务
        echo "[" . getmypid() . "] Processing task $taskId\n";
        sleep(rand(1, 3));
        echo "[" . getmypid() . "] Completed task $taskId\n";
        exit(0);
    }
}
// 等待所有子进程完成
while ($currentProcesses > 0) {
    $pid = pcntl_wait($status);
    if ($pid > 0) {
        $currentProcesses--;
    }
}
echo "All tasks completed\n";

使用Process类(Symfony Process组件)

composer require symfony/process
<?php
require_once 'vendor/autoload.php';
use Symfony\Component\Process\Process;
use Symfony\Component\Process\Exception\ProcessFailedException;
// 并发执行多个进程
$processes = [];
// 创建多个进程
for ($i = 1; $i <= 3; $i++) {
    $process = new Process(['php', '-r', "echo 'Process $i running\n'; sleep(2); echo 'Process $i done\n';"]);
    $process->start();
    $processes[] = $process;
}
// 等待所有进程完成
foreach ($processes as $process) {
    $process->wait();
    if (!$process->isSuccessful()) {
        throw new ProcessFailedException($process);
    }
    echo $process->getOutput();
}

并发执行命令(更高级的用法)

<?php
class ProcessManager {
    private $maxProcesses;
    private $running = [];
    public function __construct($maxProcesses = 5) {
        $this->maxProcesses = $maxProcesses;
    }
    public function executeCommands(array $commands) {
        $results = [];
        $commandQueue = $commands;
        $completed = 0;
        $total = count($commands);
        while ($completed < $total) {
            // 启动新的进程
            while (count($this->running) < $this->maxProcesses && !empty($commandQueue)) {
                $command = array_shift($commandQueue);
                $process = $this->startProcess($command);
                $this->running[] = [
                    'process' => $process,
                    'command' => $command
                ];
            }
            // 检查运行的进程
            foreach ($this->running as $key => $running) {
                $process = $running['process'];
                if (!$process->isRunning()) {
                    $results[] = [
                        'command' => $running['command'],
                        'output' => $process->getOutput(),
                        'exitCode' => $process->getExitCode()
                    ];
                    unset($this->running[$key]);
                    $completed++;
                }
            }
            // 避免CPU过载
            usleep(10000); // 10ms
        }
        return $results;
    }
    private function startProcess($command) {
        $process = new Process($command);
        $process->setTimeout(3600); // 1小时超时
        $process->start();
        return $process;
    }
}
// 使用示例
$manager = new ProcessManager(3);
$commands = [
    ['php', 'script1.php'],
    ['php', 'script2.php'],
    ['php', 'script3.php'],
    ['php', 'script4.php'],
    ['php', 'script5.php'],
];
$results = $manager->executeCommands($commands);
print_r($results);

使用共享内存通信

<?php
$shmKey = ftok(__FILE__, 't');
$shmId = shm_attach($shmKey, 1024, 0666);
$processes = [];
for ($i = 1; $i <= 5; $i++) {
    $pid = pcntl_fork();
    if ($pid == -1) {
        die("Could not fork");
    } elseif ($pid) {
        $processes[] = $pid;
    } else {
        // 子进程写入共享内存
        $result = processData($i);
        shm_has_var($shmId, $i) ? shm_remove_var($shmId, $i) : null;
        shm_put_var($shmId, $i, $result);
        exit(0);
    }
}
// 等待所有子进程完成
foreach ($processes as $pid) {
    pcntl_waitpid($pid, $status);
}
// 读取共享内存结果
for ($i = 1; $i <= 5; $i++) {
    if (shm_has_var($shmId, $i)) {
        $result = shm_get_var($shmId, $i);
        echo "Process $i result: $result\n";
    }
}
shm_remove($shmId);

使用队列系统(高级方案)

对于复杂的并发需求,推荐使用消息队列系统:

// 使用Redis作为队列
$redis = new Redis();
$redis->connect('127.0.0.1', 6379);
// 生产者
function addToQueue($job) {
    global $redis;
    $redis->rPush('job_queue', serialize($job));
}
// 消费者(可以启动多个实例)
function processQueue() {
    global $redis;
    while (true) {
        $job = $redis->blPop('job_queue', 10);
        if ($job) {
            $jobData = unserialize($job[1]);
            echo "[" . getmypid() . "] Processing job: " . $jobData['name'] . "\n";
            sleep($jobData['time']);
            echo "[" . getmypid() . "] Job completed\n";
        }
    }
}

注意事项

  1. 环境要求:PCNTL扩展只能在CLI模式下使用,不能在Web服务器环境(如Apache/FPM)中使用

  2. 资源管理:注意文件描述符限制、内存限制等资源管理

  3. 信号处理:正确处理SIGCHLD信号,避免僵尸进程

  4. 调试困难:多进程程序的调试比较困难,建议做好日志记录

  5. 替代方案:对于Web应用,考虑使用消息队列(RabbitMQ、Redis Queue)、任务调度器(如Gearman)等更成熟的方案

选择哪种方案取决于你的具体需求,简单的脚本可以使用PCNTL,复杂的生产环境建议使用消息队列系统。

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