Ga naar inhoud
getnextpdf.com

Enterprise editie

Stream: verwerking van documenttaken

NextPDF\Enterprise\Stream\DocumentJobStreamProcessor zet een stream van render-manifesten om in duurzame, verantwoorde resultaten. Het consumeert een iterable<RenderManifest> als een generator, rendert begrensde vensters via de Pro render engine en finaliseert elke taak in bronvolgorde. Elke taak eindigt in precies één terminale staat: output gecommit, herkend als reeds gecommit, of dead-lettered. De voortgang wordt gecheckpoint, zodat een gecrashte run hervat zonder iets opnieuw te publiceren.

Het Stream-verhaal splitst zich over twee edities, en die splitsing is doelbewust. Pro levert de duurzame, concurrente render engine en de lokale, single-host filesystem-stores — de in-process helft. Enterprise levert deze document-job stream processor plus de onderdelen die hostgrenzen overschrijden: de object-storage committer (ObjectStorageCommitter) en de duurzame outbox van terminale events (FilesystemOutboxEmitter). De Pro Stream-pagina stelt dezelfde grens vanuit haar kant.

Deze mogelijkheid wordt geleverd in NextPDF Enterprise (nextpdf/enterprise) en activeert met een Enterprise-tier license envelope. Een deployment zonder dat entitlement laadt de klassen van de mogelijkheid niet. Vergelijk edities en verkrijg een licentie.

Terminal window
composer require nextpdf/enterprise

De klassen op deze pagina leven onder NextPDF\Enterprise\Stream en NextPDF\Enterprise\Stream\Storage. Ze consumeren de bevroren Pro-contracten in NextPDF\Pro\Stream — engine-, committer-, checkpoint-, idempotency-, retry- en dead-letter-interfaces.

De taak van de processor is leveringssemantiek, niet rendering. Het groepeert de manifeststream in bron-offset-vensters die niet groter zijn dan de batchgrootte van de engine. Elk venster rendert via RenderEngineInterface::renderBatch(), met begrensde, deterministische retry van per-item timeouts. Vervolgens wordt elk item in bron-offset-volgorde gefinaliseerd naar een terminaal resultaat.

De exactly-once-grens ligt per item, en ze is verankerd in de committer, niet in coördinatie. Duurzame voortgang is een 1-based offset high-watermark: elke offset op of onder het checkpoint heeft een terminaal resultaat bereikt. De barrièrevolgorde ligt vast: commit de bytes, verplaats de watermark, sla het checkpoint op, flush gebufferde idempotency-marks, en zend dan terminale events uit. Een crash tussen commit en checkpoint hercommit idempotent bij hervatting, omdat de committer digests vergelijkt. Een crash na het checkpoint spoelt voorbij de offset, zodat niets tweemaal publiceert.

ObjectStorageCommitter implementeert de Pro OutputCommitterInterface tegen een object store via de minimale ObjectStorageClientInterface. De container van het target is de bucket en de key de object key. Identieke bytes opnieuw committen is een digest-vergeleken no-op. Divergerende bytes zonder overwrite veroorzaken het SPEC-COMMIT-409-conflict. Een vers object wordt alleen ooit aangemaakt met de atomaire putIfAbsent() conditionele write; die race verliezen triggert een begrensde re-read-and-resolve-lus. Cross-writer exactly-once houdt daarom precies zo ver stand als de putIfAbsent() van jouw adapter een echte conditionele write is — If-None-Match: * op S3, ifGenerationMatch: 0 op GCS. Deze cyclus levert de interface plus de in-memory NullObjectStorageClient; de live S3/GCS-adapter is host-geleverd.

Terminale events sluiten de lus voor downstream-systemen. Na de checkpoint-barrière probeert de processor een JobTerminalEvent uit te zenden voor elke gefinaliseerde taak — identifiers, status, receipt, foutdetails, aantal pogingen, en nooit enige PDF-bytes. Met een simpele callback-emitter is emissie at-most-once: events na een checkpoint kunnen bij crash-resume worden overgeslagen. FilesystemOutboxEmitter maakt elk event duurzaam zodra emit() draait: elk event is één atomair JSON-bestand dat wordt benoemd naar een hash van zijn deterministische eventId, zodat opnieuw uitzenden na een resume idempotent is, een relay at-least-once levert, en consumers dedupliceren op eventId. Eén grens blijft hoe dan ook bestaan: emissie gebeurt na de checkpoint-barrière, dus een crash tussen checkpoint.save() en emit() slaat bij resume het terminale event van dat item over. Downstream-systemen die een compleet event-grootboek vereisen, moeten reconciliëren tegen de gecommitte objecten (de store is de bron van waarheid), niet tegen de outbox alleen.

Licentiëring is in het output-pad ingebouwd. De withBrandingFromLicense()-factory herleidt eenmaal per run een evaluatie-branding-strategie uit de licentie. Een betaalde licentie herleidt naar een identiteitstransformatie. Een evaluatie- of ontbrekende licentie voorziet elk gecommit document van een watermerk, en een document dat niet gebrand kan worden, wordt dead-lettered — de processor commit nooit ongebrande evaluatie-bytes.

De dragende beslissing is dat exactly-once rust op het digest-vergeleken, conditioneel aangemaakte object van de committer — niet op distributed locks of consensus. De conditionele write van de object store is de enige atomaire primitieve die het ontwerp vereist, en al het overige mag falen en herstellen. Daarom moet de render engine side-effect-vrij blijven, daarom worden keyed state en in-run dedup-caches behandeld als herberekenbare versnellingen, en daarom breekt een dubbelzinnige commit de run af in plaats van te gokken: het resume-pad convergeert via dezelfde digest-vergelijking. Het is ook waarom single-writer per runId een uitgesproken vereiste is in plaats van een afgedwongen lease — de checkpoint store blijft doelbewust simpel, en de commit-laag blijft het vangnet.

Ontwerpachtergrond: High-volume documentgeneratie.

Hosts zouden via de factory moeten construeren, zodat de license-to-branding-control nooit ongekoppeld blijft:

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 is Symfony\Component\Clock\ClockInterface (de retry-backoff slaapt erdoorheen). Een null-licentie herleidt fail-closed naar evaluatie-branding.

Het enige toegangspunt verwerkt één run en retourneert de tellers ervan:

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

Gooit of faalt met: NextPDF\Enterprise\Stream\Exception\StreamProcessorException wanneer crash-safety-precondities falen (een crashSafe-run met niet-duurzame collaborators) of wanneer een commit dubbelzinnig is; InvalidArgumentException wanneer windowSize de maxBatchSize() van de engine overschrijdt.

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

Gooit of faalt met: InvalidArgumentException wanneer windowSize of checkpointIntervalJobs onder 1 ligt. $runId is de stabiele, single-writer run-identifier die checkpoint-resume keyt.

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

$scheme noemt het target-scheme dat deze committer bedient (bijvoorbeeld s3 of gcs); $clock is hier Psr\Clock\ClockInterface.

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

Gooit of faalt met: UnsupportedTargetException bij een scheme-mismatch; RenderManifestException wanneer de target key niet container-relatief-veilig is; CommitIntegrityException wanneer de gedeclareerde sha-256 niet overeenkomt met de bytes; OutputCommitConflictException (SPEC-COMMIT-409) bij divergerende bytes zonder overwrite; RuntimeException wanneer de create-race na 5 pogingen onder concurrente mutatie niet kan convergeren.

Het minimale adapter-oppervlak dat een live S3/GCS-integratie implementeert:

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() moet een echte atomaire conditionele create zijn (If-None-Match: * op S3, ifGenerationMatch: 0 op GCS) en retourneert alleen true wanneer deze aanroep het object schreef. put() is de onconditionele overwrite die uitsluitend gebruikt wordt wanneer het manifest overwrite aanvroeg.

JobCompletionEmitterInterface en FilesystemOutboxEmitter

Sectie met titel “JobCompletionEmitterInterface en FilesystemOutboxEmitter”
public function emit(JobTerminalEvent $event): void;

De emitter is optioneel op de processor. Events vuren alleen nadat een item duurzaam is gefinaliseerd. FilesystemOutboxEmitter is de geleverde duurzame implementatie:

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

Gooit of faalt met: InvalidArgumentException wanneer de directory niet bestaat; emit() gooit RuntimeException als een event niet als JSON kan worden gecodeerd. hasEvent(string $eventId): bool controleert de outbox; count(): int rapporteert niet-geleverde events.

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 is deterministisch — runId:sourceOffset:idempotencyKey:status — wat outbox-dedup mogelijk maakt. toArray() serialiseert het event voor transport; het draagt geen PDF-bytes. JobTerminalStatus is een string-enum: Committed (committed), DeadLettered (dead_lettered), Skipped (skipped).

Onveranderlijke tellers geretourneerd door process(): runId, sourceRead, fastForwardedByCheckpoint, skippedByIdempotency, windows, renderBatchCalls, renderRetries, commitReceipts, deadLettered, checkpointSaves, en finalCommittedOffset (de finale terminale high-watermark).

Exactly-once object-storage commit in isolatie. De in-memory NullObjectStorageClient staat in voor jouw S3/GCS-adapter; de semantiek die je waarneemt, is die welke een live adapter moet behouden.

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
}

Verwachte output:

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

Een volledige crash-safe run: duurzame Pro-stores, de object-storage committer, een duurzame outbox, en license-herleide branding. Dezelfde runId opnieuw draaien na een crash spoelt vooruit en convergeert.

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,
);

Voorbeeldoutput (tellers hangen af van jouw job stream):

run nightly-invoices-2026-07-03: read=1200 committed=1187 dedup-skipped=13 dead-lettered=0 checkpoints=12 final-offset=1200
  • Single-writer per runId is jouw verantwoordelijkheid. De checkpoint store heeft geen lease of compare-and-swap. Twee concurrente writers op één runId vallen buiten het contract; dwing exclusiviteit af in je scheduler.
  • crashSafe: true faalt snel op niet-duurzame collaborators. De committer-, checkpoint-, idempotency- en dead-letter-stores moeten allemaal de DurableCapability-marker implementeren, anders gooit process() een StreamProcessorException die de overtreders benoemt. De keyed state store is doelbewust uitgezonderd: verloren keyed state wordt vanaf het checkpoint voorwaarts herberekend.
  • windowSize moet passen bij de engine. Een venster groter dan maxBatchSize() gooit een InvalidArgumentException voordat enig werk begint.
  • Een dubbelzinnige commit breekt af; een conflict niet. SPEC-COMMIT-409 is een deterministisch terminaal conflict: het item wordt dead-lettered en de run gaat door. Elke andere commit-fout is dubbelzinnig: het gefinaliseerde prefix wordt gecheckpoint en de run gooit voor veilige hervatting.
  • Render-fouten breken de run nooit af. Een per-item Failed-resultaat, een uitgeput retry-budget, of onbrandbare evaluatie-bytes dead-letteren dat item allemaal en gaan door.
  • Dubbele jobId-waarden zijn veilig; dubbel werk wordt gekeyed op idempotencyKey. Resultaten correleren met items via unieke bron-offset, nooit via jobId. Een dubbele idempotency-key wordt herkend zelfs binnen hetzelfde barrière-interval, vóór het opnieuw renderen.
  • Emitter-duurzaamheid bepaalt de event-semantiek. Een simpele callback-emitter is observer-only en at-most-once over een crash. FilesystemOutboxEmitter maakt de outbox duurzaam en dedup-keyed; relay-levering is dan at-least-once, en downstream exactly-once vereist consumer-dedup op eventId. De directory ervan (zoals die van elke filesystem-store) moet vooraf bestaan, anders gooit de constructor een InvalidArgumentException.
  • Skipped events staan standaard uit. Zet emitSkippedCompletions: true om ook een Skipped-terminaal event uit te zenden voor dedup-kortgesloten items.
  • Output keys falen gesloten. commit() her-verifieert dat de target key container-relatief-veilig is: geen ..-traversal, geen absolute escape, geen null-byte, geen ingebedde stream-wrapper-scheme, en geen dubbele punt (die de NTFS alternate-data-stream-vector sluit). Onveilige keys gooien vóór enige storage-aanroep.
  • Integriteit wordt aan de grens opnieuw geverifieerd. De committer herberekent sha-256 over de daadwerkelijke bytes en verwerpt een mismatch met een CommitIntegrityException, zodat een gecorrumpeerde overdracht niet stilzwijgend kan landen.
  • Events dragen geen documentinhoud. JobTerminalEvent en de outbox-rijen bevatten alleen identifiers, digests, timestamps en foutstrings. Foutmeldingen kunnen engine-diagnostiek echoën; schrob ze, en elk tenant-identificerend jobId-schema, voordat je outbox-bestanden naar externe sinks verstuurt.
  • Evaluatie-output wordt nooit ongebrand gepubliceerd. Wanneer branding vereist is en niet toegepast kan worden, wordt het item dead-lettered in plaats van gecommit.
  • Cross-writer exactly-once is slechts zo sterk als jouw adapter. Als putIfAbsent() geen echte atomaire conditionele write is, degradeert de garantie tot single-writer-semantiek. Object-store-credentials en bucket policy zijn host-zaken; de module beheert ze nooit.

Geen gepubliceerde standaard definieert het gedrag van deze module. De exactly-once-, checkpoint- en outbox-garanties op deze pagina zijn engineering-contracten van de NextPDF Enterprise API, hier vermeld als extern waarneembaar gedrag — ze zijn geen conformiteit met, of certificering tegen, enige standaard. Het interne gebruik van SHA-256 als integriteitsdigest is eveneens leidingwerk, geen compliance-claim. Zoals overal in NextPDF: support is geen conformiteit, en conformiteit is geen certificering. NextPDF houdt geen certificering en verleent er geen; of een deployment gebouwd op deze module voldoet aan jouw regulatoire of contractuele verplichtingen is een bepaling voor jouw assessors.

Terminale events maken geen deel uit van de checkpoint-transactie: emissie draait na checkpoint.save(), dus de outbox houdt elk uitgezonden event duurzaam vast maar is geen compleet grootboek over crashes heen. Gecommitte objecten blijven de bron van waarheid.

  • Elke bron-offset op of onder finalCommittedOffset heeft precies één terminaal resultaat bereikt: Committed, Skipped, of DeadLettered.
  • Items worden in bron-offset-volgorde gefinaliseerd; de barrièrevolgorde is commit, checkpoint save, idempotency-mark flush, en dan event-emissie.
  • Een run opnieuw draaien met dezelfde runId publiceert nooit dubbel: gecheckpointe offsets spoelen vooruit, en byte-identieke re-commits zijn digest-vergeleken no-ops met idempotentReuse: true.
  • Een vers object wordt alleen ooit aangemaakt via de atomaire conditionele create; divergerende bytes op een bezette key zonder overwrite zijn een deterministische SPEC-COMMIT-409-dead-letter, nooit een clobber.
  • Een dubbelzinnige commit checkpoint het gefinaliseerde prefix en breekt af met een StreamProcessorException; de gefaalde offset wordt niet verplaatst.
  • Een crashSafe-run weigert niet-duurzame committer-, checkpoint-, idempotency- of dead-letter-collaborators voordat enige input wordt gelezen.
  • Event-ids zijn een pure functie van run, offset, idempotency-key en status, zodat een duurzame outbox ten hoogste één rij per event bevat.

NextPDF Core rendert één document tegelijk via de writer en het render-manifest-contract — zie Writer. Core alleen heeft geen duurzame job streams, geen checkpoint-resume, geen idempotency-dedup, geen object-storage commit, en geen terminal-event-outbox. NextPDF Pro voegt de duurzame, concurrente render engine en de single-host filesystem-stores toe (Stream in Pro). De cross-host-helft — deze processor, de object-storage committer, en de duurzame outbox — vereist NextPDF Enterprise.

Deze pagina documenteert alleen extern waarneembaar gedrag en het ondersteunde publieke API-oppervlak. Interne namespace-paden, helper-klassen, mechanisme-tabellen, runbook-bestandsnamen en ticket-prefixen vallen buiten de scope.