Туториали

PHP Generators: Stream Large Datasets Without Eating All Your RAM

PHP генератори: Обработувајте големи збирки на податоци без да ја потрошите целата RAM меморија

Генераторот не прави голем сет на податоци мал. Тој го прави сетот на податоци инкрементален: еден запис влегува во меморијата, еден запис се валидира и еден запис излегува. Таа разлика спречува извоз од повеќе гигабајти да се претвори во итна промена на ограничувањето на меморијата.

Овој туторијал создава PHP 8.3 команда наменета за продукциска употреба, која чита JSON со разделување по нови редови, го валидира и нормализира секој запис и запишува NDJSON на стандардниот излез. Може да напојува датотека, компресор или друг процес, притоа одржувајќи ограничена меморија и природен повратен притисок.

Што градиме

Извозникот прифаќа локална NDJSON-датотека што содржи записи за клиенти. Секоја влезна линија мора да биде JSON-објект со полиња од тип низа именувани id, email и created_at. Дополнителните полиња се отфрлаат, со што се спречува проширување на горната шема да протече неочекувани податоци.

Патеката на податоците е намерно синхрона:

  1. Генераторот чита една ограничена линија.
  2. PHP го декодира, валидира и нормализира тој запис.
  3. Запишувачот целосно ја емитува резултирачката JSON-линија.
  4. Дури тогаш генераторот продолжува и чита друг запис.

Ако излезот е цевка и нејзиниот потрошувач се забави, оперативниот систем на крај го полни баферот на цевката. Следниот блокирачки fwrite() чека, што го запира генераторот да повлекува повеќе влез. Тоа е безбеден повратен притисок без редици, циклуси за анкетирање или постојано растечка низа.

Користењето меморија е константно во однос на бројот на записи, иако останува пропорционално на најголемиот дозволен запис. Затоа наметнуваме ограничување од еден мегабајт по запис.

Предуслови и структура на проектот

Потребен ви е PHP 8.3 со овозможени стандардните JSON и filter функционалности. Примерот за распоредување исто така користи Bash, gzip, systemd и Linux-домаќин на кој имате административни привилегии.

ndjson-exporter/
├── bin/
│   └── export.php
├── deploy/
│   └── customer-export.service
└── src/
    ├── NdjsonWriter.php
    └── RecordStream.php

Не е потребен менаџер за пакети ниту зависност од трета страна. Ова го одржува моделот на извршување недвосмислено синхрон и го прави однесувањето на повратниот притисок лесно за проверка.

Имплементирајте го ограничениот генератор на записи

Создадете src/RecordStream.php. Генераторот е сопственик на влезната рачка и ја затвора во блок finally, вклучително и кога валидацијата не успее или потрошувачот запре порано.

<?php

declare(strict_types=1);

final class RecordStream
{
    private const MAX_RECORD_BYTES = 1_048_576;

    public function __construct(
        private readonly string $path,
    ) {
    }

    /**
     * @return Generator<int, array{id: string, email: string, created_at: string}>
     */
    public function records(): Generator
    {
        if (!is_file($this->path) || !is_readable($this->path)) {
            throw new RuntimeException(
                sprintf('Input is not a readable local file: %s', $this->path)
            );
        }

        $handle = @fopen($this->path, 'rb');

        if ($handle === false) {
            throw new RuntimeException(
                sprintf('Could not open input: %s', $this->path)
            );
        }

        try {
            $recordNumber = 0;

            while (
                ($line = fgets($handle, self::MAX_RECORD_BYTES + 2)) !== false
            ) {
                $recordNumber++;
                $hasNewline = str_ends_with($line, "\n");
                $payload = rtrim($line, "\r\n");

                if (
                    (!$hasNewline && !feof($handle))
                    || strlen($payload) > self::MAX_RECORD_BYTES
                ) {
                    throw new RuntimeException(
                        sprintf('Record %d exceeds the byte limit', $recordNumber)
                    );
                }

                if ($payload === '') {
                    throw new RuntimeException(
                        sprintf('Record %d is empty', $recordNumber)
                    );
                }

                try {
                    $decoded = json_decode(
                        $payload,
                        false,
                        512,
                        JSON_THROW_ON_ERROR
                    );
                } catch (JsonException $exception) {
                    throw new RuntimeException(
                        sprintf(
                            'Record %d contains invalid JSON: %s',
                            $recordNumber,
                            $exception->getMessage()
                        ),
                        0,
                        $exception
                    );
                }

                if (!$decoded instanceof stdClass) {
                    throw new RuntimeException(
                        sprintf('Record %d must be a JSON object', $recordNumber)
                    );
                }

                yield $recordNumber => $this->normalize(
                    get_object_vars($decoded),
                    $recordNumber
                );
            }

            if (!feof($handle)) {
                throw new RuntimeException('The input stream failed before EOF');
            }
        } finally {
            fclose($handle);
        }
    }

    /**
     * @param array<string, mixed> $record
     * @return array{id: string, email: string, created_at: string}
     */
    private function normalize(array $record, int $recordNumber): array
    {
        foreach (['id', 'email', 'created_at'] as $field) {
            if (!array_key_exists($field, $record) || !is_string($record[$field])) {
                throw new RuntimeException(
                    sprintf(
                        'Record %d requires string field "%s"',
                        $recordNumber,
                        $field
                    )
                );
            }
        }

        $id = trim($record['id']);
        $email = trim($record['email']);
        $createdAt = trim($record['created_at']);

        if ($id === '' || !ctype_digit($id)) {
            throw new RuntimeException(
                sprintf('Record %d has an invalid id', $recordNumber)
            );
        }

        if (filter_var($email, FILTER_VALIDATE_EMAIL) === false) {
            throw new RuntimeException(
                sprintf('Record %d has an invalid email address', $recordNumber)
            );
        }

        $date = DateTimeImmutable::createFromFormat(
            '!Y-m-d\TH:i:sP',
            $createdAt
        );

        if (
            $date === false
            || $date->format('Y-m-d\TH:i:sP') !== $createdAt
        ) {
            throw new RuntimeException(
                sprintf(
                    'Record %d has a non-canonical created_at value',
                    $recordNumber
                )
            );
        }

        return [
            'id' => $id,
            'email' => $email,
            'created_at' => $date->format(DATE_ATOM),
        ];
    }
}

Должината проследена до fgets() остава простор за откривање запис што го надминува ограничувањето. Последната линија без нов ред останува валидна, но преголема линија се одбива пред JSON-декодирањето да може да ја зголеми нејзината цена во меморијата.

Целосно запишете го секој запис

Еден повик на fwrite() не е загарантирано дека ќе ја потроши целата низа. Запишувачот мора да напредува низ делумните запишувања и и false и запишување од нула бајти да ги третира како неуспеси.

Создадете src/NdjsonWriter.php:

<?php

declare(strict_types=1);

final class NdjsonWriter
{
    /**
     * @param iterable<int, array<string, string>> $records
     * @param resource $output
     */
    public function write(
        iterable $records,
        $output,
        int $reportEvery = 100_000
    ): int {
        $count = 0;
        $startedAt = hrtime(true);

        foreach ($records as $recordNumber => $record) {
            try {
                $json = json_encode(
                    $record,
                    JSON_THROW_ON_ERROR | JSON_UNESCAPED_SLASHES
                );
            } catch (JsonException $exception) {
                throw new RuntimeException(
                    sprintf(
                        'Could not encode record %d: %s',
                        $recordNumber,
                        $exception->getMessage()
                    ),
                    0,
                    $exception
                );
            }

            $this->writeAll($output, $json . "\n");
            $count++;

            if ($reportEvery > 0 && $count % $reportEvery === 0) {
                $this->report($count, $startedAt);
            }
        }

        if (!fflush($output)) {
            throw new RuntimeException('Could not flush the output stream');
        }

        $this->report($count, $startedAt);

        return $count;
    }

    /**
     * @param resource $output
     */
    private function writeAll($output, string $bytes): void
    {
        $offset = 0;
        $length = strlen($bytes);

        while ($offset < $length) {
            $written = @fwrite($output, substr($bytes, $offset));

            if ($written === false || $written === 0) {
                throw new RuntimeException(
                    'Output failed or the downstream consumer closed the pipe'
                );
            }

            $offset += $written;
        }
    }

    private function report(int $count, int $startedAt): void
    {
        $seconds = max((hrtime(true) - $startedAt) / 1_000_000_000, 0.001);
        $message = sprintf(
            "records=%d elapsed_seconds=%.3f records_per_second=%.1f peak_bytes=%d\n",
            $count,
            $seconds,
            $count / $seconds,
            memory_get_peak_usage(true)
        );

        @fwrite(STDERR, $message);
    }
}

Напредокот оди на стандардниот излез за грешки, никогаш на стандардниот излез, така што телеметријата не може да го оштети NDJSON-текот. Расипан компресор или затворена цевка предизвикува ненулто излегување на извозникот наместо тивко скратен успех.

Составете ја командата

Создадете bin/export.php:

<?php

declare(strict_types=1);

require dirname(__DIR__) . '/src/RecordStream.php';
require dirname(__DIR__) . '/src/NdjsonWriter.php';

function main(array $arguments): int
{
    if (count($arguments) !== 2) {
        @fwrite(
            STDERR,
            sprintf("Usage: php %s /path/to/customers.ndjson\n", $arguments[0])
        );

        return 64;
    }

    try {
        $stream = new RecordStream($arguments[1]);
        $writer = new NdjsonWriter();
        $writer->write($stream->records(), STDOUT);

        return 0;
    } catch (Throwable $exception) {
        @fwrite(STDERR, 'export_error=' . $exception->getMessage() . "\n");

        return 1;
    }
}

exit(main($argv));

Командата прифаќа само обична, читлива локална датотека. Тоа ги исклучува PHP URL-обвивките и случајните мрежни читања. Исто така, ги зачувува идентификаторите како низи, така што водечките нули опстануваат при извозот.

Тестирајте ги успешните и неуспешните патеки

Извршете го овој интеграциски тест од коренот на проектот под Bash. Тој користи уникатен привремен директориум, проверува точен излез, го проверува бројот на линии и потврдува дека неправилниот влез не успева.

set -euo pipefail

test_dir="$(mktemp -d)"

cat > "$test_dir/input.ndjson" <<'EOF'
{"id":"001","email":"[email protected]","created_at":"2026-01-01T00:00:00+00:00","internal_note":"discard me"}
{"id":"002","email":"[email protected]","created_at":"2026-01-02T12:30:00+01:00"}
EOF

php -d memory_limit=64M bin/export.php \
  "$test_dir/input.ndjson" > "$test_dir/actual.ndjson"

diff -u \
  <(printf '%s\n' \
    '{"id":"001","email":"[email protected]","created_at":"2026-01-01T00:00:00+00:00"}' \
    '{"id":"002","email":"[email protected]","created_at":"2026-01-02T12:30:00+01:00"}') \
  "$test_dir/actual.ndjson"

test "$(wc -l < "$test_dir/actual.ndjson")" -eq 2

printf '%s\n' '{"id":"003","email":"not-an-email","created_at":"2026-01-03T00:00:00+00:00"}' \
  > "$test_dir/invalid.ndjson"

if php bin/export.php "$test_dir/invalid.ndjson" \
  > "$test_dir/rejected.ndjson"; then
    printf '%s\n' 'Expected invalid input to fail' >&2
    exit 1
fi

printf 'Tests passed; temporary files remain at %s\n' "$test_dir"

За голема пробна датотека, проверете ја врвната резидентна меморија и пропусноста без задржување на излезот:

/usr/bin/time -v \
  php -d memory_limit=64M bin/export.php \
  /srv/import/customers.ndjson > /dev/null

Бројот на записи треба да има мало влијание врз врвната меморија. Големината на записот, JSON-декодирањето, трошокот на PHP-извршувањето и конфигурираната максимална големина на линија сè уште се важни.

Распоредете како зацврстена сериска услуга

Услугата подолу стримува во gzip. Бидејќи Bash работи со pipefail, неуспех на PHP или gzip го спречува конечното преместување. Привремените и конечните датотеки се наоѓаат на истиот датотечен систем, што го прави успешното преместување атомско за читателите.

Создадете deploy/customer-export.service:

[Unit]
Description=Stream and compress the normalized customer export
After=local-fs.target

[Service]
Type=oneshot
User=customer-export
Group=customer-export
UMask=0077
ExecStart=/bin/bash -o pipefail -c '/usr/bin/php -d memory_limit=64M /opt/customer-export/bin/export.php /srv/import/customers.ndjson | /usr/bin/gzip -c > /srv/export/customers.ndjson.gz.part && /usr/bin/mv -f /srv/export/customers.ndjson.gz.part /srv/export/customers.ndjson.gz'
TimeoutStopSec=30s
KillMode=control-group
NoNewPrivileges=yes
PrivateDevices=yes
PrivateTmp=yes
ProtectHome=yes
ProtectSystem=strict
ReadOnlyPaths=/opt/customer-export /srv/import
ReadWritePaths=/srv/export
RestrictAddressFamilies=AF_UNIX
RestrictSUIDSGID=yes
LockPersonality=yes
MemoryDenyWriteExecute=yes
MemoryMax=128M

[Install]
WantedBy=multi-user.target

Инсталирајте ја со изречна сопственост и дозволи. Извршете useradd само кога за првпат ја создавате услужната сметка.

sudo useradd --system \
  --home-dir /nonexistent \
  --shell /usr/sbin/nologin \
  customer-export

sudo install -d -o root -g root -m 0755 /opt/customer-export
sudo install -d -o root -g root -m 0755 /opt/customer-export/bin
sudo install -d -o root -g root -m 0755 /opt/customer-export/src
sudo install -d -o root -g customer-export -m 0750 /srv/import
sudo install -d -o customer-export -g customer-export -m 0750 /srv/export

sudo install -o root -g root -m 0644 \
  src/RecordStream.php /opt/customer-export/src/RecordStream.php
sudo install -o root -g root -m 0644 \
  src/NdjsonWriter.php /opt/customer-export/src/NdjsonWriter.php
sudo install -o root -g root -m 0644 \
  bin/export.php /opt/customer-export/bin/export.php
sudo install -o root -g root -m 0644 \
  deploy/customer-export.service \
  /etc/systemd/system/customer-export.service

sudo install -o root -g customer-export -m 0640 \
  customers.ndjson /srv/import/customers.ndjson.next
sudo mv -f \
  /srv/import/customers.ndjson.next \
  /srv/import/customers.ndjson

sudo systemctl daemon-reload
sudo systemctl start customer-export.service
sudo systemctl status customer-export.service
sudo journalctl -u customer-export.service --no-pager

Услугата не отвора мрежен слушач, па не ѝ е потребно правило за заштитен ѕид. Нејзиното ограничување на адресните семејства исто така спречува обични интернет-врски. Влезот е само за читање, излезот е изолиран, а рестриктивниот umask ги штити извезените податоци. Ако подоцна се додаде испорака надолу по текот, повторно разгледајте ги и мрежниот песочник и ракувањето со ингеренциите, наместо неформално да ги ослабувате.

Перформанси, набљудливост и компромиси

Извозникот ги евидентира бројот на записи, поминатото време, пропусноста и врвната алоцирана меморија. Неговиот излезен код ги разликува грешките при користење од грешките при обработка, додека systemd и shell-цевководот го зачувуваат тој статус.

Компресијата може да стане тесното грло. Тоа е прифатливо: повратниот притисок го принудува PHP да чека наместо да акумулира записи. Ако времето на CPU е поважно од големината на излезот, изберете соодветно ниво на gzip-компресија по мерење со репрезентативни податоци.

Овој дизајн дава предност на точноста и ограничената меморија пред паралелната пропусност. Повеќе работници би барале поделени влезни податоци и детерминистичко составување на излезот. Додавањето асинхрона редица без строго ограничување би ја уништило централната гаранција за меморијата.

Повторните обиди го започнуваат извозот од почеток. Тие не додаваат кон објавената датотека, а читателите го гледаат или претходниот завршен извоз или новиот завршен извоз. Неуспешно извршување може да остави .part-датотека, која следното извршување безбедно ја скратува пред запишувањето.

Вообичаени продукциски неуспеси

  • Меморијата сè уште расте: побарајте iterator_to_array(), задржани записи, неограничени серии или логирање што го баферира излезот.
  • Командата изгледа заглавена: проверете го процесот надолу по текот и датотечниот систем. Блокирањето на полна цевка е очекуван повратен притисок, не нужно мртва блокада.
  • Излезот е скратен: барајте pipefail и објавувајте само откако целосниот цевковод успешно ќе заврши.
  • Еден запис ја исцрпува меморијата: задржете го ограничувањето на бајти и ограничувањето на PHP-меморијата. Константната меморија не значи дека неограничена поединечна вредност е безопасна.
  • Датумите се менуваат неочекувано: прифатете една канонска претстава на временската ознака и одбивајте нормализирани, но невалидни календарски вредности.
  • Исклучувањето остава делумни податоци: прекинете ја целата процесна група и никогаш не го преименувајте делумниот артефакт при неуспех.

Конечна контролна листа за проверка

  • Влезот е читлива локална NDJSON-датотека со еден објект по линија.
  • Секој запис е ограничен пред JSON-декодирањето.
  • Генераторот дава еден нормализиран запис во даден момент.
  • Запишувачот ракува со делумни запишувања и затворање од надолниот потрошувач.
  • Дијагностиката користи стандарден излез за грешки, а податоците користат стандарден излез.
  • Цевководот за компресија работи со pipefail.
  • Само завршен артефакт се објавува атомски.
  • Дозволите на услугата ги штитат и изворните и извезените податоци.
  • Врвната меморија останува стабилна додека бројот на записи се зголемува.
  • Неправилните, преголемите, прекинатите извршувања и извршувањата со полн диск враќаат неуспех.

Генераторите се највредни кога целиот цевковод ја почитува нивната мрзеливост. Штом секоја фаза обработува една ограничена ставка и чека следната фаза да заврши, големите датотеки престануваат да бидат проблем на управување со меморијата. Резултатот не е само умна итерација; тоа е предвидлива продукциска патека за податоци што безбедно се забавува, видливо не успева и објавува само завршена работа.

Портрет на автор на блогот

Mihajlo

Јас сум Михајло - развивач поттикнат од љубопитност, дисциплина и постојаната желба да создадам нешто значајно. Споделувам увиди, упатства и бесплатни услуги за да им помогнам на другите да ја поедностават својата работа и да растат во постојано развивачкиот свет на софтверот и вештачката интелигенција.