Перейти к содержимому
getnextpdf.com

Enterprise редакция

Stream: обработка документных заданий

NextPDF\Enterprise\Stream\DocumentJobStreamProcessor превращает поток манифестов рендеринга в устойчивые, учитываемые результаты. Он потребляет iterable<RenderManifest> как генератор, рендерит ограниченные окна через движок рендеринга Pro и финализирует каждое задание в порядке источника. Каждое задание завершается ровно в одном терминальном состоянии: вывод зафиксирован, распознан как уже зафиксированный или отправлен в dead-letter. Прогресс фиксируется чекпоинтами, поэтому аварийно завершившийся прогон возобновляется без повторной публикации чего-либо.

История Stream разделена между двумя редакциями, и это разделение сделано намеренно. Pro предоставляет устойчивый конкурентный движок рендеринга и локальные, одно-хостовые файловые хранилища — внутрипроцессную половину. Enterprise предоставляет этот процессор потока документных заданий плюс компоненты, пересекающие границы хостов: коммиттер объектного хранилища (ObjectStorageCommitter) и устойчивый outbox терминальных событий (FilesystemOutboxEmitter). Страница Pro Stream описывает ту же границу со своей стороны.

Эта возможность поставляется в NextPDF Enterprise (nextpdf/enterprise) и активируется лицензионным конвертом уровня Enterprise. Развёртывание без этого права не загружает классы возможности. Сравните редакции и получите лицензию.

Окно терминала
composer require nextpdf/enterprise

Классы на этой странице находятся в пространствах имён NextPDF\Enterprise\Stream и NextPDF\Enterprise\Stream\Storage. Они потребляют замороженные контракты Pro в NextPDF\Pro\Stream — интерфейсы движка, коммиттера, чекпоинта, идемпотентности, повторов и dead-letter.

Задача процессора — семантика доставки, а не рендеринг. Он группирует поток манифестов в окна по смещению источника не крупнее размера пакета движка. Каждое окно рендерится через RenderEngineInterface::renderBatch() с ограниченными, детерминированными повторами таймаутов по отдельным элементам. Затем каждый элемент финализируется в порядке смещения источника до терминального результата.

Граница exactly-once действует поэлементно и закреплена в коммиттере, а не в координации. Устойчивый прогресс — это верхняя отметка смещения, начинающегося с 1: каждое смещение на уровне чекпоинта или ниже достигло терминального результата. Порядок барьера фиксирован: зафиксировать байты, продвинуть отметку, сохранить чекпоинт, сбросить буферизованные метки идемпотентности, затем испустить терминальные события. Сбой между фиксацией и чекпоинтом при возобновлении повторно фиксирует идемпотентно, потому что коммиттер сравнивает дайджесты. Сбой после чекпоинта проматывает вперёд за смещение, поэтому ничто не публикуется дважды.

ObjectStorageCommitter реализует Pro-интерфейс OutputCommitterInterface поверх объектного хранилища через минимальный ObjectStorageClientInterface. Поле container цели — это бакет, а её key — ключ объекта. Повторная фиксация идентичных байтов — это no-op со сравнением дайджестов. Расходящиеся байты без overwrite вызывают конфликт SPEC-COMMIT-409. Новый объект создаётся только атомарной условной записью putIfAbsent(); проигрыш в этой гонке запускает ограниченный цикл повторного чтения и разрешения. Поэтому exactly-once между писателями сохраняется ровно настолько, насколько putIfAbsent() вашего адаптера является настоящей условной записью — If-None-Match: * в S3, ifGenerationMatch: 0 в GCS. В этом цикле поставляется интерфейс плюс in-memory-клиент NullObjectStorageClient; реальный адаптер S3/GCS предоставляется хостом.

Терминальные события замыкают петлю для нижестоящих систем. После барьера чекпоинта процессор пытается испустить JobTerminalEvent для каждого финализированного задания — идентификаторы, статус, квитанцию, детали ошибки, число попыток и никогда никаких байтов PDF. С обычным callback-эмиттером испускание происходит at-most-once: события после чекпоинта могут быть пропущены при возобновлении после сбоя. FilesystemOutboxEmitter делает каждое событие устойчивым после выполнения emit(): каждое событие — это один атомарный JSON-файл, названный по хешу его детерминированного eventId, поэтому повторное испускание после возобновления идемпотентно, реле доставляет at-least-once, а потребители дедуплицируют по eventId. Одна граница остаётся в любом случае: испускание происходит после барьера чекпоинта, поэтому сбой между checkpoint.save() и emit() при возобновлении пропускает терминальное событие этого элемента. Нижестоящие системы, которым требуется полный журнал событий, должны сверяться с зафиксированными объектами (хранилище — источник истины), а не только с outbox.

Лицензирование встроено в путь вывода. Фабрика withBrandingFromLicense() один раз за прогон разрешает стратегию оценочного брендинга из лицензии. Платная лицензия разрешается в тождественное преобразование. Оценочная или отсутствующая лицензия наносит водяной знак на каждый зафиксированный документ, а документ, который не удаётся брендировать, отправляется в dead-letter — процессор никогда не фиксирует небрендированные оценочные байты.

Несущее решение в том, что exactly-once опирается на объект коммиттера, создаваемый условно и сравниваемый по дайджесту, — а не на распределённые блокировки или консенсус. Условная запись объектного хранилища — единственный атомарный примитив, который требуется дизайну, и всему остальному позволено падать и восстанавливаться. Именно поэтому движок рендеринга должен оставаться без побочных эффектов, почему keyed-состояние и внутрипрогонные кеши дедупликации трактуются как пересчитываемые ускорения и почему неоднозначная фиксация прерывает прогон, а не гадает: путь возобновления сходится через то же сравнение дайджестов. Именно поэтому единственный писатель на runId — заявленное требование, а не принудительная аренда: хранилище чекпоинтов намеренно остаётся простым, а слой фиксации остаётся страховочной сеткой.

Проектный фон: Генерация документов большого объёма.

Хостам следует конструировать через фабрику, чтобы управление связкой «лицензия — брендинг» никогда не оставалось неподключённым:

public static function withBrandingFromLicense(
RenderEngineInterface $engine,
OutputCommitterInterface $committer,
IdempotencyStoreInterface $idempotency,
CheckpointStoreInterface $checkpoints,
KeyedStateStoreInterface $state,
DeadLetterStoreInterface $deadLetters,
RetryPolicy $retryPolicy,
ClockInterface $clock,
EntitlementEvaluator $entitlementEvaluator,
?LicenseKey $license,
?StreamProcessorProbe $probe = null,
?JobCompletionEmitterInterface $emitter = null,
?BrandingApplicator $brandingApplicator = null,
): self

$clock — это Symfony\Component\Clock\ClockInterface (через него засыпает backoff повторов). Лицензия null разрешается fail-closed в оценочный брендинг.

Единственная точка входа обрабатывает один прогон и возвращает его счётчики:

public function process(iterable $manifests, StreamProcessorConfig $config): ProcessingSummary

Бросает или падает с: NextPDF\Enterprise\Stream\Exception\StreamProcessorException, когда не выполнены предусловия crash-safety (прогон crashSafe с не-устойчивыми коллабораторами) или когда фиксация неоднозначна; InvalidArgumentException, когда windowSize превышает maxBatchSize() движка.

public function __construct(
public string $runId,
int $windowSize = 32,
int $checkpointIntervalJobs = 100,
public bool $crashSafe = true,
public bool $emitSkippedCompletions = false,
)

Бросает или падает с: InvalidArgumentException, когда windowSize или checkpointIntervalJobs меньше 1. $runId — стабильный идентификатор прогона с единственным писателем, по которому ключуется возобновление чекпоинта.

public function __construct(
private ObjectStorageClientInterface $client,
private string $scheme,
private ClockInterface $clock,
) {}

$scheme называет целевую схему, которую обслуживает этот коммиттер (например, s3 или gcs); $clock здесь — Psr\Clock\ClockInterface.

public function commit(
string $jobId,
OutputObjectKey $target,
string $bytes,
string $sha256,
bool $overwrite = false,
): CommitReceipt

Бросает или падает с: UnsupportedTargetException при несовпадении схемы; RenderManifestException, когда целевой ключ не является безопасным относительно контейнера; CommitIntegrityException, когда объявленный sha-256 не соответствует байтам; OutputCommitConflictException (SPEC-COMMIT-409) при расходящихся байтах без overwrite; RuntimeException, когда гонка создания не может сойтись после 5 попыток при конкурентной мутации.

Минимальная поверхность адаптера, которую реализует реальная интеграция с S3/GCS:

public function shaOf(string $bucket, string $key): ?string;
public function put(string $bucket, string $key, string $bytes, string $sha256): void;
public function putIfAbsent(string $bucket, string $key, string $bytes, string $sha256): bool;

putIfAbsent() должен быть настоящим атомарным условным созданием (If-None-Match: * в S3, ifGenerationMatch: 0 в GCS) и возвращать true только тогда, когда именно этот вызов записал объект. put() — безусловная перезапись, используемая исключительно когда манифест запросил overwrite.

public function emit(JobTerminalEvent $event): void;

Эмиттер для процессора опционален. События срабатывают только после того, как элемент устойчиво финализирован. FilesystemOutboxEmitter — поставляемая устойчивая реализация:

public function __construct(string $directory, ?AtomicFileWriter $writer = null)

Бросает или падает с: InvalidArgumentException, когда каталог не существует; emit() бросает RuntimeException, если событие не удаётся закодировать в JSON. hasEvent(string $eventId): bool проверяет outbox; count(): int сообщает число недоставленных событий.

public function __construct(
public string $eventId,
public string $runId,
public int $sourceOffset,
public string $jobId,
public string $idempotencyKeyValue,
public JobTerminalStatus $status,
public ?CommitReceipt $receipt,
public ?string $errorCode,
public ?string $errorMessage,
public int $attempts,
public DateTimeImmutable $occurredAt,
) {}

eventId детерминирован — runId:sourceOffset:idempotencyKey:status — и именно это делает возможной дедупликацию в outbox. toArray() сериализует событие для транспорта; оно не несёт байтов PDF. JobTerminalStatus — строковый enum: Committed (committed), DeadLettered (dead_lettered), Skipped (skipped).

Неизменяемые счётчики, возвращаемые process(): runId, sourceRead, fastForwardedByCheckpoint, skippedByIdempotency, windows, renderBatchCalls, renderRetries, commitReceipts, deadLettered, checkpointSaves и finalCommittedOffset (финальная терминальная верхняя отметка).

Exactly-once-фиксация в объектное хранилище в изоляции. In-memory-клиент NullObjectStorageClient замещает ваш адаптер S3/GCS; наблюдаемая вами семантика — та, которую реальный адаптер обязан сохранить.

stream-object-commit-quickstart.php
<?php
declare(strict_types=1);
require __DIR__ . '/vendor/autoload.php';
use NextPDF\Enterprise\Stream\Storage\NullObjectStorageClient;
use NextPDF\Enterprise\Stream\Storage\ObjectStorageCommitter;
use NextPDF\Manifest\OutputObjectKey;
use NextPDF\Pro\Stream\Exception\OutputCommitConflictException;
use Symfony\Component\Clock\NativeClock;
$committer = new ObjectStorageCommitter(
client: new NullObjectStorageClient(), // swap in your S3/GCS adapter
scheme: 's3',
clock: new NativeClock(),
);
$target = new OutputObjectKey(scheme: 's3', container: 'invoices', key: '2026/07/inv-1001.pdf');
$bytes = '%PDF-1.7 example-rendered-bytes';
$sha = hash('sha256', $bytes);
$first = $committer->commit('inv-1001', $target, $bytes, $sha);
$replay = $committer->commit('inv-1001', $target, $bytes, $sha); // crash-resume replay
printf("first : reuse=%s, %d bytes\n", var_export($first->idempotentReuse, true), $first->bytesWritten);
printf("replay: reuse=%s\n", var_export($replay->idempotentReuse, true));
try {
$divergent = '%PDF-1.7 different-bytes';
$committer->commit('inv-1001', $target, $divergent, hash('sha256', $divergent));
} catch (OutputCommitConflictException $conflict) {
echo 'conflict: ' . $conflict->specCode() . "\n"; // no silent clobber
}

Ожидаемый вывод:

first : reuse=false, 31 bytes
replay: reuse=true
conflict: SPEC-COMMIT-409

Полный crash-safe-прогон: устойчивые хранилища Pro, коммиттер объектного хранилища, устойчивый outbox и брендинг, разрешённый из лицензии. Повторный запуск того же runId после сбоя проматывает вперёд и сходится.

stream-run-production.php
<?php
declare(strict_types=1);
require __DIR__ . '/vendor/autoload.php';
use NextPDF\Enterprise\Licensing\EntitlementEvaluator;
use NextPDF\Enterprise\Stream\DocumentJobStreamProcessor;
use NextPDF\Enterprise\Stream\Exception\StreamProcessorException;
use NextPDF\Enterprise\Stream\FilesystemOutboxEmitter;
use NextPDF\Enterprise\Stream\Storage\ObjectStorageCommitter;
use NextPDF\Enterprise\Stream\StreamProcessorConfig;
use NextPDF\Manifest\Render\SingleDocumentRenderer;
use NextPDF\Manifest\RenderManifest;
use NextPDF\Pro\Stream\Checkpoint\FilesystemCheckpointStore;
use NextPDF\Pro\Stream\Dedup\FilesystemIdempotencyStore;
use NextPDF\Pro\Stream\Engine\InProcessRenderEngine;
use NextPDF\Pro\Stream\Retry\FilesystemDeadLetterStore;
use NextPDF\Pro\Stream\Retry\RetryPolicy;
use NextPDF\Pro\Stream\State\InMemoryKeyedStateStore;
use Symfony\Component\Clock\NativeClock;
// Production requires a host-supplied adapter whose putIfAbsent() is a TRUE
// atomic conditional create (S3 If-None-Match: *, GCS ifGenerationMatch: 0)
// and whose shaOf() reads durable object state. NullObjectStorageClient is
// for the quick start only - it keeps nothing across processes.
$s3Client = new \Aws\S3\S3Client(['region' => 'eu-central-1', 'version' => 'latest']);
$objectClient = new \Acme\Storage\S3ObjectStorageClient($s3Client); // implements ObjectStorageClientInterface
$stateDir = '/var/lib/nextpdf/stream';
foreach (['checkpoints', 'idempotency', 'dead-letters', 'outbox'] as $sub) {
if (!is_dir($stateDir . '/' . $sub)) {
mkdir($stateDir . '/' . $sub, 0770, true);
}
}
// One manifest per JSONL line; the generator never materialises the batch.
$manifests = (static function (string $path): Generator {
$handle = fopen($path, 'rb');
if ($handle === false) {
throw new RuntimeException('Cannot open job stream: ' . $path);
}
try {
while (($line = fgets($handle)) !== false) {
if (trim($line) !== '') {
yield RenderManifest::fromJson(trim($line));
}
}
} finally {
fclose($handle);
}
})('/var/spool/nextpdf/jobs.jsonl');
$license = null; // your licensing bootstrap yields a LicenseKey; null = evaluation branding
$processor = DocumentJobStreamProcessor::withBrandingFromLicense(
engine: new InProcessRenderEngine(SingleDocumentRenderer::standalone()),
// For a live bucket, implement ObjectStorageClientInterface over your S3/GCS SDK.
committer: new ObjectStorageCommitter($objectClient, 's3', new NativeClock()),
idempotency: new FilesystemIdempotencyStore($stateDir . '/idempotency'),
checkpoints: new FilesystemCheckpointStore($stateDir . '/checkpoints'),
state: new InMemoryKeyedStateStore(), // recomputable; durability not required here
deadLetters: new FilesystemDeadLetterStore($stateDir . '/dead-letters'),
retryPolicy: new RetryPolicy(maxAttempts: 3, baseDelayMs: 200, maxDelayMs: 5_000),
clock: new NativeClock(),
entitlementEvaluator: new EntitlementEvaluator(),
license: $license,
emitter: new FilesystemOutboxEmitter($stateDir . '/outbox'),
);
$config = new StreamProcessorConfig(
runId: 'nightly-invoices-2026-07-03',
windowSize: 32,
checkpointIntervalJobs: 100,
crashSafe: true,
);
try {
$summary = $processor->process($manifests, $config);
} catch (StreamProcessorException $e) {
// Ambiguous commit or a non-durable collaborator: the finalized prefix is
// checkpointed. Re-run the SAME runId; the committer converges by digest.
fwrite(STDERR, 'Run aborted for safe resume: ' . $e->getMessage() . PHP_EOL);
exit(1);
}
printf(
"run %s: read=%d committed=%d dedup-skipped=%d dead-lettered=%d checkpoints=%d final-offset=%d\n",
$summary->runId,
$summary->sourceRead,
$summary->commitReceipts,
$summary->skippedByIdempotency,
$summary->deadLettered,
$summary->checkpointSaves,
$summary->finalCommittedOffset,
);

Пример вывода (счётчики зависят от вашего потока заданий):

run nightly-invoices-2026-07-03: read=1200 committed=1187 dedup-skipped=13 dead-lettered=0 checkpoints=12 final-offset=1200
  • Единственный писатель на runId — ваша ответственность. У хранилища чекпоинтов нет ни аренды, ни compare-and-swap. Два конкурентных писателя на одном runId находятся вне контракта; обеспечивайте эксклюзивность в своём планировщике.
  • crashSafe: true быстро падает на не-устойчивых коллабораторах. Хранилища коммиттера, чекпоинта, идемпотентности и dead-letter должны все реализовывать маркер DurableCapability, иначе process() бросает StreamProcessorException с именами нарушителей. Хранилище keyed-состояния намеренно исключено: потерянное keyed-состояние пересчитывается от чекпоинта вперёд.
  • windowSize должен вписываться в движок. Окно крупнее maxBatchSize() бросает InvalidArgumentException до начала любой работы.
  • Неоднозначная фиксация прерывает прогон; конфликт — нет. SPEC-COMMIT-409 — детерминированный терминальный конфликт: элемент отправляется в dead-letter, а прогон продолжается. Любой другой сбой фиксации неоднозначен: финализированный префикс фиксируется чекпоинтом, а прогон бросает исключение ради безопасного возобновления.
  • Сбои рендеринга никогда не прерывают прогон. Поэлементный результат Failed, исчерпанный бюджет повторов или небрендируемые оценочные байты — всё это отправляет элемент в dead-letter и продолжает работу.
  • Дубликаты значений jobId безопасны; дублирование работы ключуется по idempotencyKey. Результаты соотносятся с элементами по уникальному смещению источника, а не по jobId. Дубликат ключа идемпотентности распознаётся даже в пределах одного барьерного интервала, до повторного рендеринга.
  • Устойчивость эмиттера определяет семантику событий. Обычный callback-эмиттер только наблюдает и работает at-most-once при сбое. FilesystemOutboxEmitter делает outbox устойчивым и с ключом дедупликации; доставка через реле тогда at-least-once, а exactly-once ниже по потоку требует дедупликации потребителем по eventId. Его каталог (как у любого файлового хранилища) должен существовать заранее, иначе конструктор бросает InvalidArgumentException.
  • Пропущенные события по умолчанию выключены. Установите emitSkippedCompletions: true, чтобы также испускать терминальное событие Skipped для элементов, закороченных дедупликацией.
  • Ключи вывода отказывают безопасно. commit() повторно утверждает, что целевой ключ безопасен относительно контейнера: никакого обхода .., никакого абсолютного выхода, никакого нулевого байта, никакой встроенной схемы stream-wrapper и никакого двоеточия (которое закрывает вектор альтернативных потоков данных NTFS). Небезопасные ключи бросают исключение до любого обращения к хранилищу.
  • Целостность повторно проверяется на границе. Коммиттер пересчитывает sha-256 по фактическим байтам и отклоняет несовпадение с CommitIntegrityException, поэтому повреждённая передача не может пройти незаметно.
  • События не несут содержимого документа. JobTerminalEvent и строки outbox содержат только идентификаторы, дайджесты, метки времени и строки ошибок. Сообщения об ошибках могут отражать диагностику движка; вычищайте их, а также любую идентифицирующую арендатора схему jobId, прежде чем отправлять файлы outbox в сторонние приёмники.
  • Оценочный вывод никогда не публикуется небрендированным. Когда брендинг требуется, но не может быть применён, элемент отправляется в dead-letter, а не фиксируется.
  • Exactly-once между писателями настолько силён, насколько силён ваш адаптер. Если putIfAbsent() не является настоящей атомарной условной записью, гарантия деградирует до семантики единственного писателя. Учётные данные объектного хранилища и политика бакета — забота хоста; модуль ими никогда не управляет.

Никакой опубликованный стандарт не определяет поведение этого модуля. Гарантии exactly-once, чекпоинта и outbox на этой странице — инженерные контракты API NextPDF Enterprise, изложенные здесь как внешне наблюдаемое поведение — они не являются соответствием какому-либо стандарту или сертификацией по нему. Внутреннее использование SHA-256 в качестве дайджеста целостности — тоже техническая деталь, а не заявление о соответствии. Как и везде в NextPDF: поддержка — не соответствие, а соответствие — не сертификация. NextPDF не имеет сертификации и не выдаёт таковой; отвечает ли развёртывание, построенное на этом модуле, вашим нормативным или договорным обязательствам, определяют ваши оценщики.

Терминальные события не являются частью транзакции чекпоинта: испускание выполняется после checkpoint.save(), поэтому outbox устойчиво хранит каждое испущенное событие, но не является полным журналом при сбоях. Зафиксированные объекты остаются источником истины.

  • Каждое смещение источника на уровне finalCommittedOffset или ниже достигло ровно одного терминального результата: Committed, Skipped или DeadLettered.
  • Элементы финализируются в порядке смещения источника; порядок барьера — фиксация, сохранение чекпоинта, сброс метки идемпотентности, затем испускание события.
  • Повторный запуск прогона с тем же runId никогда не публикует дважды: смещения из чекпоинта проматываются вперёд, а побайтно идентичные повторные фиксации — это no-op со сравнением дайджестов и idempotentReuse: true.
  • Новый объект создаётся только через атомарное условное создание; расходящиеся байты по занятому ключу без overwrite — детерминированный dead-letter SPEC-COMMIT-409, а не перезапись.
  • Неоднозначная фиксация фиксирует чекпоинтом финализированный префикс и прерывается с StreamProcessorException; провалившееся смещение не продвигается.
  • Прогон crashSafe отвергает не-устойчивые коллабораторы коммиттера, чекпоинта, идемпотентности или dead-letter до чтения любого ввода.
  • Идентификаторы событий — чистая функция от прогона, смещения, ключа идемпотентности и статуса, поэтому устойчивый outbox хранит не более одной строки на событие.

NextPDF Core рендерит по одному документу за раз через writer и контракт манифеста рендеринга — см. Writer. Один только Core не имеет ни устойчивых потоков заданий, ни возобновления по чекпоинту, ни дедупликации по идемпотентности, ни фиксации в объектное хранилище, ни outbox терминальных событий. NextPDF Pro добавляет устойчивый конкурентный движок рендеринга и одно-хостовые файловые хранилища (Stream в Pro). Меж-хостовая половина — этот процессор, коммиттер объектного хранилища и устойчивый outbox — требует NextPDF Enterprise.

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