Ir al contenido
getnextpdf.com

Enterprise edición

Stream: procesamiento de trabajos de documento

NextPDF\Enterprise\Stream\DocumentJobStreamProcessor convierte un flujo de manifiestos de renderizado en resultados duraderos y verificables. Consume un iterable<RenderManifest> como generador, renderiza ventanas acotadas a través del motor de renderizado Pro y finaliza cada trabajo en orden de origen. Cada trabajo termina en exactamente un estado terminal: salida confirmada, reconocida como ya confirmada o enviada a dead-letter. El progreso se registra con checkpoints, de modo que una ejecución caída se reanuda sin volver a publicar nada.

La historia de Stream se divide entre dos ediciones, y la división es deliberada. Pro aporta el motor de renderizado duradero y concurrente y los almacenes locales del sistema de archivos de un solo host: la mitad en proceso. Enterprise aporta este procesador de flujos de trabajos de documento más las piezas que cruzan las fronteras de host: el confirmador de almacenamiento de objetos (ObjectStorageCommitter) y el outbox duradero de eventos terminales (FilesystemOutboxEmitter). La página de Pro Stream enuncia la misma frontera desde su lado.

Esta capacidad se distribuye en NextPDF Enterprise (nextpdf/enterprise) y se activa con un sobre de licencia de nivel Enterprise. Un despliegue sin ese derecho no carga las clases de la capacidad. Compare ediciones y obtenga una licencia.

Ventana de terminal
composer require nextpdf/enterprise

Las clases de esta página residen en NextPDF\Enterprise\Stream y NextPDF\Enterprise\Stream\Storage. Consumen los contratos Pro congelados de NextPDF\Pro\Stream: interfaces de motor, confirmador, checkpoint, idempotencia, reintento y dead-letter.

El trabajo del procesador es la semántica de entrega, no el renderizado. Agrupa el flujo de manifiestos en ventanas de desplazamiento de origen no mayores que el tamaño de lote del motor. Cada ventana se renderiza a través de RenderEngineInterface::renderBatch(), con reintento acotado y determinista de los timeouts por elemento. Luego cada elemento se finaliza en orden de desplazamiento de origen hacia un resultado terminal.

La frontera exactly-once es por elemento, y está anclada en el confirmador, no en la coordinación. El progreso duradero es una marca de agua máxima de desplazamiento basada en 1: cada desplazamiento igual o inferior al checkpoint ha alcanzado un resultado terminal. El orden de la barrera es fijo: confirmar los bytes, avanzar la marca de agua, guardar el checkpoint, vaciar las marcas de idempotencia en búfer y luego emitir los eventos terminales. Una caída entre la confirmación y el checkpoint vuelve a confirmar de forma idempotente al reanudar, porque el confirmador compara digests. Una caída después del checkpoint avanza rápido más allá del desplazamiento, de modo que nada se publica dos veces.

ObjectStorageCommitter implementa la OutputCommitterInterface de Pro contra un almacén de objetos a través de la mínima ObjectStorageClientInterface. El container del destino es el bucket y su key la clave del objeto. Volver a confirmar bytes idénticos es un no-op comparado por digest. Bytes divergentes sin overwrite generan el conflicto SPEC-COMMIT-409. Un objeto nuevo solo se crea con la escritura condicional atómica putIfAbsent(); perder esa carrera desencadena un bucle acotado de relectura y resolución. El exactly-once entre escritores se mantiene, por tanto, exactamente en la medida en que el putIfAbsent() de su adaptador sea una escritura condicional verdadera: If-None-Match: * en S3, ifGenerationMatch: 0 en GCS. Este ciclo distribuye la interfaz más el NullObjectStorageClient en memoria; el adaptador S3/GCS en vivo lo suministra el host.

Los eventos terminales cierran el bucle para los sistemas downstream. Tras la barrera del checkpoint, el procesador intenta emitir un JobTerminalEvent por cada trabajo finalizado: identificadores, estado, recibo, detalles de error, recuento de intentos, y nunca ningún byte de PDF. Con un emisor de callback simple, la emisión es at-most-once: los eventos posteriores a un checkpoint pueden omitirse al reanudar tras una caída. FilesystemOutboxEmitter hace que cada evento sea duradero una vez que se ejecuta emit(): cada evento es un archivo JSON atómico nombrado por un hash de su eventId determinista, de modo que reemitir tras una reanudación es idempotente, un relay entrega at-least-once y los consumidores deduplican por eventId. Una frontera permanece de cualquier modo: la emisión ocurre después de la barrera del checkpoint, así que una caída entre checkpoint.save() y emit() omite el evento terminal de ese elemento al reanudar. Los sistemas downstream que requieran un libro mayor de eventos completo deben reconciliar contra los objetos confirmados (el almacén es la fuente de verdad), no contra el outbox por sí solo.

Las licencias están integradas en la ruta de salida. La fábrica withBrandingFromLicense() resuelve una estrategia de branding de evaluación una vez por ejecución a partir de la licencia. Una licencia de pago se resuelve en una transformación de identidad. Una licencia de evaluación o ausente pone una marca de agua en cada documento confirmado, y un documento que no puede recibir branding se envía a dead-letter: el procesador nunca confirma bytes de evaluación sin marca.

La decisión estructural es que exactly-once descansa en el objeto creado de forma condicional y comparado por digest del confirmador, no en bloqueos distribuidos ni consenso. La escritura condicional del almacén de objetos es la única primitiva atómica que el diseño requiere, y todo lo demás puede fallar y recuperarse. Por eso el motor de renderizado debe permanecer libre de efectos secundarios, por eso el estado con clave y las cachés de deduplicación en ejecución se tratan como aceleraciones recalculables, y por eso una confirmación ambigua aborta la ejecución en lugar de adivinar: la ruta de reanudación converge a través de la misma comparación de digest. También es por eso que el único escritor por runId es un requisito enunciado y no una concesión (lease) impuesta: el almacén de checkpoints se mantiene deliberadamente simple, y la capa de confirmación queda como la red de seguridad.

Antecedentes de diseño: Generación de documentos de alto volumen.

Los hosts deben construir a través de la fábrica, de modo que el control de licencia a branding nunca quede sin cablear:

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 es Symfony\Component\Clock\ClockInterface (el backoff de reintento duerme a través de él). Una licencia null se resuelve fail-closed a branding de evaluación.

El único punto de entrada procesa una ejecución y devuelve sus contadores:

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

Lanza o falla con: NextPDF\Enterprise\Stream\Exception\StreamProcessorException cuando fallan las precondiciones de seguridad ante caídas (una ejecución crashSafe con colaboradores no duraderos) o cuando una confirmación es ambigua; InvalidArgumentException cuando windowSize excede el maxBatchSize() del motor.

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

Lanza o falla con: InvalidArgumentException cuando windowSize o checkpointIntervalJobs es inferior a 1. $runId es el identificador de ejecución estable y de un solo escritor que sirve de clave para la reanudación por checkpoint.

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

$scheme nombra el esquema de destino que atiende este confirmador (por ejemplo s3 o gcs); $clock aquí es Psr\Clock\ClockInterface.

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

Lanza o falla con: UnsupportedTargetException ante un desajuste de esquema; RenderManifestException cuando la clave de destino no es segura relativa al contenedor; CommitIntegrityException cuando el sha-256 declarado no coincide con los bytes; OutputCommitConflictException (SPEC-COMMIT-409) ante bytes divergentes sin overwrite; RuntimeException cuando la carrera de creación no puede converger tras 5 intentos bajo mutación concurrente.

La superficie mínima de adaptador que implementa una integración S3/GCS en vivo:

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() debe ser una creación condicional atómica real (If-None-Match: * en S3, ifGenerationMatch: 0 en GCS) y devuelve true solo cuando esta llamada escribió el objeto. put() es la sobrescritura incondicional que se usa únicamente cuando el manifiesto solicitó overwrite.

JobCompletionEmitterInterface y FilesystemOutboxEmitter

Sección titulada «JobCompletionEmitterInterface y FilesystemOutboxEmitter»
public function emit(JobTerminalEvent $event): void;

El emisor es opcional en el procesador. Los eventos se disparan solo después de que un elemento se finaliza de forma duradera. FilesystemOutboxEmitter es la implementación duradera distribuida:

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

Lanza o falla con: InvalidArgumentException cuando el directorio no existe; emit() lanza RuntimeException si un evento no puede codificarse en JSON. hasEvent(string $eventId): bool comprueba el outbox; count(): int informa de los eventos no entregados.

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 es determinista — runId:sourceOffset:idempotencyKey:status — que es lo que hace posible la deduplicación del outbox. toArray() serializa el evento para el transporte; no lleva bytes de PDF. JobTerminalStatus es un enum de tipo string: Committed (committed), DeadLettered (dead_lettered), Skipped (skipped).

Contadores inmutables devueltos por process(): runId, sourceRead, fastForwardedByCheckpoint, skippedByIdempotency, windows, renderBatchCalls, renderRetries, commitReceipts, deadLettered, checkpointSaves y finalCommittedOffset (la marca de agua terminal máxima final).

Confirmación de almacenamiento de objetos exactly-once de forma aislada. El NullObjectStorageClient en memoria hace las veces de su adaptador S3/GCS; la semántica que observa es la que un adaptador en vivo debe preservar.

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
}

Salida esperada:

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

Una ejecución completa segura ante caídas: almacenes Pro duraderos, el confirmador de almacenamiento de objetos, un outbox duradero y branding resuelto por licencia. Volver a ejecutar el mismo runId tras una caída avanza rápido y converge.

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

Salida de ejemplo (los contadores dependen de su flujo de trabajos):

run nightly-invoices-2026-07-03: read=1200 committed=1187 dedup-skipped=13 dead-lettered=0 checkpoints=12 final-offset=1200
  • El único escritor por runId es su responsabilidad. El almacén de checkpoints no tiene lease ni compare-and-swap. Dos escritores concurrentes sobre un mismo runId están fuera del contrato; imponga la exclusividad en su planificador.
  • crashSafe: true falla rápido ante colaboradores no duraderos. El confirmador y los almacenes de checkpoint, idempotencia y dead-letter deben implementar todos el marcador DurableCapability, o process() lanza StreamProcessorException nombrando a los infractores. El almacén de estado con clave está deliberadamente exento: el estado con clave perdido se recalcula desde el checkpoint hacia adelante.
  • windowSize debe caber en el motor. Una ventana mayor que maxBatchSize() lanza InvalidArgumentException antes de que comience cualquier trabajo.
  • Una confirmación ambigua aborta; un conflicto no. SPEC-COMMIT-409 es un conflicto terminal determinista: el elemento se envía a dead-letter y la ejecución continúa. Cualquier otro fallo de confirmación es ambiguo: el prefijo finalizado se registra con checkpoint y la ejecución lanza para una reanudación segura.
  • Los fallos de renderizado nunca abortan la ejecución. Un resultado Failed por elemento, un presupuesto de reintentos agotado o bytes de evaluación no marcables envían todos ese elemento a dead-letter y continúan.
  • Los valores de jobId duplicados son seguros; el trabajo duplicado se distingue por idempotencyKey. Los resultados se correlacionan con los elementos por desplazamiento de origen único, nunca por jobId. Una clave de idempotencia duplicada se reconoce incluso dentro del mismo intervalo de barrera, antes de volver a renderizar.
  • La durabilidad del emisor decide la semántica de eventos. Un emisor de callback simple es solo observador y at-most-once ante una caída. FilesystemOutboxEmitter hace que el outbox sea duradero y con clave de deduplicación; la entrega por relay es entonces at-least-once, y el exactly-once downstream requiere deduplicación del consumidor por eventId. Su directorio (como el de cada almacén del sistema de archivos) debe preexistir, o el constructor lanza InvalidArgumentException.
  • Los eventos omitidos están desactivados por defecto. Establezca emitSkippedCompletions: true para emitir también un evento terminal Skipped para los elementos cortocircuitados por deduplicación.
  • Las claves de salida fallan en cerrado. commit() reafirma que la clave de destino es segura relativa al contenedor: sin traversal .., sin escape absoluto, sin byte nulo, sin esquema de stream-wrapper embebido y sin dos puntos (que cierra el vector de flujo de datos alternativo NTFS). Las claves inseguras lanzan antes de cualquier llamada al almacenamiento.
  • La integridad se reverifica en la frontera. El confirmador recalcula el sha-256 sobre los bytes reales y rechaza un desajuste con CommitIntegrityException, de modo que un traspaso corrupto no puede asentarse silenciosamente.
  • Los eventos no llevan contenido de documento. JobTerminalEvent y las filas del outbox contienen solo identificadores, digests, marcas de tiempo y cadenas de error. Los mensajes de error pueden hacer eco de diagnósticos del motor; depúrelos, junto con cualquier esquema de jobId que identifique al inquilino, antes de enviar los archivos del outbox a destinos de terceros.
  • La salida de evaluación nunca se publica sin marca. Cuando el branding es obligatorio y no puede aplicarse, el elemento se envía a dead-letter en lugar de confirmarse.
  • El exactly-once entre escritores es solo tan fuerte como su adaptador. Si putIfAbsent() no es una escritura condicional atómica verdadera, la garantía degrada a semántica de un solo escritor. Las credenciales del almacén de objetos y la política del bucket son asuntos del host; el módulo nunca los gestiona.

Ningún estándar publicado define el comportamiento de este módulo. Las garantías de exactly-once, checkpoint y outbox de esta página son contratos de ingeniería de la API de NextPDF Enterprise, enunciados aquí como comportamiento observable externamente; no son conformidad con, ni certificación frente a, ningún estándar. El uso interno de SHA-256 como digest de integridad es igualmente plomería, no una afirmación de cumplimiento. Como en todo NextPDF: el soporte no es conformidad, y la conformidad no es certificación. NextPDF no posee certificación alguna ni la otorga; si un despliegue construido sobre este módulo cumple sus obligaciones regulatorias o contractuales es una determinación que corresponde a sus evaluadores.

Los eventos terminales no forman parte de la transacción del checkpoint: la emisión se ejecuta después de checkpoint.save(), de modo que el outbox conserva de forma duradera cada evento emitido, pero no es un libro mayor completo a través de las caídas. Los objetos confirmados siguen siendo la fuente de verdad.

  • Cada desplazamiento de origen igual o inferior a finalCommittedOffset ha alcanzado exactamente un resultado terminal: Committed, Skipped o DeadLettered.
  • Los elementos se finalizan en orden de desplazamiento de origen; el orden de la barrera es confirmación, guardado de checkpoint, vaciado de marca de idempotencia y luego emisión de evento.
  • Volver a ejecutar una ejecución con el mismo runId nunca publica dos veces: los desplazamientos con checkpoint avanzan rápido, y las reconfirmaciones byte a byte idénticas son no-ops comparados por digest con idempotentReuse: true.
  • Un objeto nuevo solo se crea a través de la creación condicional atómica; los bytes divergentes en una clave ocupada sin overwrite son un dead-letter SPEC-COMMIT-409 determinista, nunca una sobrescritura.
  • Una confirmación ambigua registra con checkpoint el prefijo finalizado y aborta con StreamProcessorException; el desplazamiento fallido no avanza.
  • Una ejecución crashSafe rechaza colaboradores de confirmación, checkpoint, idempotencia o dead-letter no duraderos antes de leer cualquier entrada.
  • Los identificadores de evento son una función pura de ejecución, desplazamiento, clave de idempotencia y estado, de modo que un outbox duradero conserva como máximo una fila por evento.

NextPDF Core renderiza un documento a la vez a través del writer y el contrato de manifiesto de renderizado — véase Writer. Core por sí solo no tiene flujos de trabajos duraderos, ni reanudación por checkpoint, ni deduplicación por idempotencia, ni confirmación de almacenamiento de objetos, ni outbox de eventos terminales. NextPDF Pro añade el motor de renderizado duradero y concurrente y los almacenes del sistema de archivos de un solo host (Stream en Pro). La mitad entre hosts —este procesador, el confirmador de almacenamiento de objetos y el outbox duradero— requiere NextPDF Enterprise.

Esta página documenta únicamente el comportamiento observable externamente y la superficie de API pública soportada. Las rutas de espacio de nombres internas, las clases auxiliares, las tablas de mecanismos, los nombres de archivo de runbooks y los prefijos de tickets quedan fuera de alcance.