Многопроцессные и многопоточные вычисления в PHP

Раздел: Веб-разработка на PHP -> Работа с данными

Основные подходы к параллельной обработке данных в PHP

Как выполнить несколько задач одновременно в PHP с максимальной производительностью?

Для параллельных вычислений внутри одного скрипта наиболее эффективным решением является расширение parallel (доступно в PHP 7.4+). Оно создаёт потоки выполнения в одном адресном пространстве, что снижает накладные расходы на межпроцессное взаимодействие. Расширение устанавливается через Composer: composer require parallel/parallel или сборкой из исходных кодов.


<?php
$runtime = new \parallel\Runtime();
$future = $runtime->run(function() {
    // Тяжелая задача
    return array_sum(range(1, 10000000));
});
$result = $future->value(); // Получение результата (блокировка)
echo "Сумма: $result";
?>

View php table (php: отображение таблиц)

Пояснение: parallel\Runtime создаёт новый поток, run() принимает замыкание, выполняемое в этом потоке. Возвращаемое значение извлекается через value(). Метод блокирует основной поток до готовности. Для параллельного запуска нескольких задач создаются несколько экземпляров Runtime.

Типичные проблемы: замыкания должны быть самодостаточными (не ссылаться на переменные извне, если они не сериализуемы). Ошибка сериализации приводит к исключению. Также потоки не могут напрямую изменять глобальные переменные. Для обмена данными используются каналы (parallel\Channel).

Решение: передавать только примитивные типы или массивы, избегать объектов с замыканиями. Для сложного обмена применять каналы.

Как организовать параллельную обработку через системные процессы с помощью pcntl_fork?

Функция pcntl_fork() создаёт дочерний процесс как точную копию родителя. Подходит для изолированных задач, требующих собственного адресного пространства. Дочерний процесс наследует все переменные и открытые файловые дескрипторы.


<?php
$pid = pcntl_fork();
if ($pid == -1) {
    die('Ошибка форка');
} elseif ($pid == 0) {
    // Дочерний процесс
    $sum = array_sum(range(1, 1000000));
    exit($sum); // Завершаемся с кодом
} else {
    // Родительский процесс
    $status = pcntl_waitpid($pid, $exitCode); // Ожидание завершения
    echo "Результат дочернего: $exitCode";
}
?>

несколько php (несколько php)

Пояснение: pcntl_fork() возвращает PID в родителе и 0 в потомке. pcntl_waitpid() ожидает завершения и получает код выхода. Для передачи большего объёма данных используются каналы или разделяемая память.

Типичные проблемы: зомби-процессы, если родитель не вызывает pcntl_wait(). Гонки данных при записи в один файл. Необходимость обработки сигналов (SIGCHLD). На платформе Windows функция недоступна.

Решение: всегда вызывать pcntl_waitpid() для каждого дочернего процесса в родителе, либо настроить обработчик сигнала. Для синхронизации использовать семафоры или каналы из библиотеки sysvmsg.

Как асинхронно выполнить множество HTTP запросов для сбора данных?

Расширение cURL предоставляет мульти-интерфейс для одновременной отправки запросов без блокировки. Это позволяет значительно ускорить сбор данных с нескольких ресурсов.


<?php
$urls = ['https://api.example.com/1', 'https://api.example.com/2'];
$mh = curl_multi_init();
$handles = [];
foreach ($urls as $i => $url) {
    $ch = curl_init($url);
    curl_setopt($ch, CURLOPT_RETURNTRANSFER, true);
    curl_multi_add_handle($mh, $ch);
    $handles[$i] = $ch;
}
$running = null;
do {
    curl_multi_exec($mh, $running);
    curl_multi_select($mh); // Блокировка до активности
} while ($running > 0);
$results = [];
foreach ($handles as $ch) {
    $results[] = curl_multi_getcontent($ch);
    curl_multi_remove_handle($mh, $ch);
}
curl_multi_close($mh);
print_r($results);
?>

Php tables (работа с таблицами в php)

Цикл выполняется до тех пор, пока есть активные запросы. curl_multi_select() ожидает события, уменьшая нагрузку на процессор.

Типичные проблемы: превышение лимита открытых дескрипторов, таймауты, ошибки SSL. Большое количество запросов за раз может привести к исчерпанию ресурсов.

Решение: ограничивать число одновременных соединений (например, через curl_multi_setopt($mh, CURLMOPT_MAXCONNECTS, 10)). Использовать обработку ошибок через curl_error() для каждого дескриптора.

Как организовать фоновую очередь задач для обработки данных?

Системы очередей (RabbitMQ, Redis, Beanstalkd) позволяют отделить постановку задач от их выполнения. Рабочие процессы (workers) запускаются отдельно и обрабатывают задачи из очереди. Это даёт горизонтальное масштабирование и устойчивость к сбоям.


<?php
// producer.php
$redis = new Redis();
$redis->connect('127.0.0.1', 6379);
$task = json_encode(['type' => 'resize', 'file' => 'image.jpg']);
$redis->lPush('task_queue', $task);
echo "Задача поставлена\n";
?>

поиск данных php (поиск данных в php (общие методы))


<?php
// worker.php
$redis = new Redis();
$redis->connect('127.0.0.1', 6379);
while ($task = $redis->brPop('task_queue', 5)) { // Блокирующее чтение
    $data = json_decode($task[1], true);
    // Обработка задачи
    echo "Обработка {$data['file']}\n";
}
?>

Пояснение: producer добавляет задачу в конец списка (lPush), worker блокируется на чтении (brPop) и выполняет задачу. Можно запустить несколько workers параллельно.

Типичные проблемы: потеря задач при падении worker, дублирование обработки, сложность управления состоянием. Необходимость подтверждения обработки (ack).

Решение: использовать брокеры с поддержкой подтверждений (RabbitMQ). Реализовать идемпотентность обработчиков. Для Redis применять паттерн “reliable queue” с использованием двух списков.

Расширенные примеры параллельных вычислений

Пример 1: Обработка массива изображений с помощью parallel

Задача: изменить размер 10 изображений в несколько потоков и замерить время выполнения.

Пример

<?php
require 'vendor/autoload.php';
use parallel\Runtime;
use parallel\Future;

$images = glob('images/*.jpg');
$total = count($images);
$workers = 4;
$chunkSize = ceil($total / $workers);
$futures = [];
$start = microtime(true);

for ($i = 0; $i < $workers; $i++) {
    $chunk = array_slice($images, $i * $chunkSize, $chunkSize);
    $runtime = new Runtime();
    $futures[] = $runtime->run(function($files) {
        $processed = 0;
        foreach ($files as $file) {
            $img = imagecreatefromjpeg($file);
            $resized = imagescale($img, 200, 150);
            imagejpeg($resized, 'output/' . basename($file));
            imagedestroy($img);
            imagedestroy($resized);
            $processed++;
        }
        return $processed;
    }, [$chunk]);
}

$totalProcessed = 0;
foreach ($futures as $future) {
    $totalProcessed += $future->value();
}

echo "Обработано изображений: $totalProcessed за " . (microtime(true) - $start) . " сек.\n";
?>
Обработано изображений: 10 за 2.345 сек.

Каждый поток обрабатывает свою часть списка. Время уменьшается почти пропорционально числу потоков (с учётом накладных расходов).

Пример 2: Параллельная обработка CSV файла с помощью pcntl_fork

Файл большого объёма разбивается на части, каждый дочерний процесс обрабатывает свою строку и записывает результат в отдельный временный файл.

Пример

<?php
$file = 'data.csv';
$lines = file($file, FILE_IGNORE_NEW_LINES);
$total = count($lines);
$workers = 4;
$chunkSize = ceil($total / $workers);
$childPids = [];

for ($i = 0; $i < $workers; $i++) {
    $pid = pcntl_fork();
    if ($pid == -1) { die('Fork failed'); }
    if ($pid == 0) {
        // Дочерний процесс
        $startLine = $i * $chunkSize;
        $chunk = array_slice($lines, $startLine, $chunkSize);
        $outFile = tempnam(sys_get_temp_dir(), "chunk_$i");
        $fh = fopen($outFile, 'w');
        foreach ($chunk as $line) {
            $processed = strtoupper($line); // Пример обработки
            fwrite($fh, $processed . PHP_EOL);
        }
        fclose($fh);
        // Передаём имя файла обратно через exit-код (ограничение) или через канал
        exit();
    } else {
        $childPids[] = $pid;
    }
}

// Родитель ожидает всех
foreach ($childPids as $pid) {
    pcntl_waitpid($pid, $status);
}
echo "Все дочерние процессы завершены.\n";
?>

Здесь не показан объединение результатов. Для передачи данных между процессами лучше использовать именованные каналы или разделяемую память.

Все дочерние процессы завершены.

Пример 3: Асинхронная загрузка нескольких URL с помощью curl_multi и прогресс бар

Пример

<?php
$urls = [];
for ($i = 0; $i < 20; $i++) {
    $urls[] = "https://httpbin.org/delay/1?num=$i";
}
$mh = curl_multi_init();
$handles = [];
foreach ($urls as $i => $url) {
    $ch = curl_init($url);
    curl_setopt_array($ch, [
        CURLOPT_RETURNTRANSFER => true,
        CURLOPT_TIMEOUT => 10,
        CURLOPT_CONNECTTIMEOUT => 5
    ]);
    curl_multi_add_handle($mh, $ch);
    $handles[$i] = $ch;
}
$active = null;
do {
    $status = curl_multi_exec($mh, $active);
    $info = curl_multi_info_read($mh);
    if ($info) {
        $ch = $info['handle'];
        $index = array_search($ch, $handles, true);
        echo "Завершён запрос #$index\n";
    }
    if ($active) {
        curl_multi_select($mh);
    }
} while ($active);

$contents = [];
foreach ($handles as $ch) {
    $contents[] = curl_multi_getcontent($ch);
    curl_multi_remove_handle($mh, $ch);
    curl_close($ch);
}
curl_multi_close($mh);
echo "Загружено " . count($contents) . " страниц.\n";
?>
Завершён запрос #0
Завершён запрос #3
...
Загружено 20 страниц.

Пример 4: Очередь задач с RabbitMQ (amqp расширение)

Пример

<?php
// producer.php
$connection = new AMQPConnection([
    'host' => 'localhost',
    'port' => 5672,
    'login' => 'guest',
    'password' => 'guest'
]);
$connection->connect();
$channel = new AMQPChannel($connection);
$exchange = new AMQPExchange($channel);
$exchange->setName('tasks');
$exchange->setType(AMQP_EX_TYPE_DIRECT);
$exchange->declare();

for ($i = 0; $i < 100; $i++) {
    $msg = json_encode(['id' => $i, 'data' => "task_$i"]);
    $exchange->publish($msg, 'processing');
}
echo "Опубликовано 100 задач\n";
$connection->disconnect();
?>
Пример

<?php
// worker.php
$connection = new AMQPConnection([
    'host' => 'localhost',
    'port' => 5672,
    'login' => 'guest',
    'password' => 'guest'
]);
$connection->connect();
$channel = new AMQPChannel($connection);
$queue = new AMQPQueue($channel);
$queue->setName('processing');
$queue->setFlags(AMQP_DURABLE);
$queue->declare();
$queue->bind('tasks', 'processing');

while (true) {
    $envelope = $queue->get();
    if ($envelope) {
        $msg = json_decode($envelope->getBody(), true);
        echo "Обрабатывается задача #{$msg['id']}\n";
        // Имитация работы
        usleep(100000);
        $queue->ack($envelope->getDeliveryTag());
    } else {
        sleep(1);
    }
}
$connection->disconnect();
?>
Обрабатывается задача #0
Обрабатывается задача #1
...

Запуск нескольких worker'ов позволит обрабатывать задачи параллельно. Каждая задача гаратированно будет обработана (благодаря подтверждению ack).

Несколько PHP - comments

En
несколько php (php)