Bỏ qua để đến nội dung
getnextpdf.com

Enterprise phiên bản

Stream: xử lý document-job

NextPDF\Enterprise\Stream\DocumentJobStreamProcessor biến một luồng render manifest thành các kết quả bền vững, có thể quy trách nhiệm. Nó tiêu thụ một iterable<RenderManifest> như một generator, render các cửa sổ có giới hạn qua render engine Pro, và hoàn tất mọi job theo thứ tự nguồn. Mỗi job kết thúc ở đúng một terminal state: output đã commit, được nhận biết là đã commit từ trước, hoặc bị dead-lettered. Tiến độ được checkpoint, nên một lần chạy bị sập sẽ resume mà không re-publish bất cứ thứ gì.

Câu chuyện Stream được chia đôi giữa hai edition, và sự phân chia này là có chủ đích. Pro cung cấp render engine bền vững, chạy đồng thời và các store filesystem cục bộ, single-host — nửa in-process. Enterprise cung cấp bộ xử lý document-job stream này cộng với các phần vượt qua ranh giới máy chủ: object-storage committer (ObjectStorageCommitter) và outbox bền vững của các terminal event (FilesystemOutboxEmitter). Trang Pro Stream nêu cùng ranh giới đó từ phía nó.

Năng lực này đi kèm trong NextPDF Enterprise (nextpdf/enterprise) và kích hoạt bằng một license envelope bậc Enterprise. Một deployment không có entitlement đó sẽ không nạp các class của năng lực này. So sánh các edition và lấy license.

Terminal window
composer require nextpdf/enterprise

Các class trên trang này nằm dưới NextPDF\Enterprise\StreamNextPDF\Enterprise\Stream\Storage. Chúng tiêu thụ các contract Pro đã đóng băng trong NextPDF\Pro\Stream — các interface engine, committer, checkpoint, idempotency, retry, và dead-letter.

Công việc của bộ xử lý là ngữ nghĩa phân phối, không phải render. Nó gom luồng manifest thành các cửa sổ theo source-offset không lớn hơn batch size của engine. Mỗi cửa sổ render qua RenderEngineInterface::renderBatch(), với retry có giới hạn, xác định các timeout theo từng item. Sau đó mọi item được hoàn tất theo thứ tự source-offset đến một kết quả cuối cùng.

Ranh giới exactly-once là theo từng item, và nó được neo trong committer, không phải trong sự phối hợp. Tiến độ bền vững là một high-watermark offset dựa trên 1: mọi offset ở mức checkpoint hoặc thấp hơn đều đã đạt một kết quả cuối cùng. Thứ tự barrier là cố định: commit bytes, tiến high-watermark, lưu checkpoint, flush các idempotency mark đã đệm, rồi phát các terminal event. Một lần sập giữa commit và checkpoint sẽ re-commit một cách idempotent khi resume, vì committer so sánh các digest. Một lần sập sau checkpoint sẽ fast-forward qua offset đó, nên không có gì publish hai lần.

ObjectStorageCommitter triển khai Pro OutputCommitterInterface đối với một object store thông qua ObjectStorageClientInterface tối giản. container của target là bucket và key của nó là object key. Re-commit các bytes giống hệt là một no-op được so sánh digest. Các bytes khác biệt mà không có overwrite sẽ nêu conflict SPEC-COMMIT-409. Một object mới chỉ bao giờ được tạo bằng thao tác conditional write nguyên tử putIfAbsent(); thua cuộc đua đó sẽ kích hoạt một vòng lặp re-read-and-resolve có giới hạn. Do đó exactly-once liên writer chỉ đúng đến mức putIfAbsent() của adapter của bạn là một conditional write thực thụ — If-None-Match: * trên S3, ifGenerationMatch: 0 trên GCS. Chu kỳ này cung cấp interface cộng với NullObjectStorageClient trong bộ nhớ; adapter S3/GCS trực tiếp do host cung cấp.

Các terminal event khép kín vòng lặp cho các hệ thống downstream. Sau barrier checkpoint, bộ xử lý cố gắng phát một JobTerminalEvent cho mỗi job đã hoàn tất — các định danh, status, receipt, chi tiết lỗi, số lần thử, và không bao giờ có bất kỳ PDF bytes nào. Với một emitter callback thuần, việc phát là at-most-once: các event sau một checkpoint có thể bị bỏ qua khi crash-resume. FilesystemOutboxEmitter khiến mọi event trở nên bền vững một khi emit() chạy: mỗi event là một file JSON nguyên tử được đặt tên theo hash của eventId xác định của nó, nên re-emit sau một lần resume là idempotent, một relay phân phối at-least-once, và các consumer dedup trên eventId. Dù cách nào thì vẫn còn một ranh giới: việc phát diễn ra sau barrier checkpoint, nên một lần sập giữa checkpoint.save()emit() sẽ bỏ qua terminal event của item đó khi resume. Các hệ thống downstream cần một sổ cái event đầy đủ nên đối chiếu với các object đã commit (store là nguồn sự thật), không phải chỉ với outbox.

Việc cấp phép được nối vào output path. Factory withBrandingFromLicense() phân giải một chiến lược evaluation-branding một lần mỗi lần chạy từ license. Một license trả phí phân giải thành một identity transform. Một license evaluation hoặc thiếu sẽ đóng dấu watermark lên mọi document đã commit, và một document không thể được branding sẽ bị dead-lettered — bộ xử lý không bao giờ commit các bytes evaluation không có branding.

Quyết định chịu tải là exactly-once dựa trên object được so sánh digest, được tạo có điều kiện của committer — không phải trên các lock phân tán hay consensus. Conditional write của object store là primitive nguyên tử duy nhất mà thiết kế yêu cầu, và mọi thứ khác được phép thất bại rồi khôi phục. Đó là lý do render engine phải giữ không có side-effect, lý do keyed state và các dedup cache trong lần chạy được coi là những gia tốc có thể tính lại, và lý do một commit mập mờ hủy bỏ lần chạy thay vì đoán mò: đường resume hội tụ qua cùng phép so sánh digest. Đó cũng là lý do single-writer theo mỗi runId là một yêu cầu được nêu ra chứ không phải một lease được cưỡng chế — checkpoint store cố ý giữ đơn giản, và lớp commit đóng vai trò lưới an toàn.

Nền tảng thiết kế: Sinh document số lượng lớn.

Các host nên khởi tạo qua factory, để control license-to-branding không bao giờ bị bỏ chưa nối:

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

$clockSymfony\Component\Clock\ClockInterface (retry backoff ngủ thông qua nó). Một license null phân giải fail-closed sang evaluation branding.

Điểm vào duy nhất xử lý một lần chạy và trả về các counter của nó:

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

Ném hoặc thất bại với: NextPDF\Enterprise\Stream\Exception\StreamProcessorException khi các điều kiện tiên quyết an toàn khi sập thất bại (một lần chạy crashSafe với các collaborator không bền vững) hoặc khi một commit mập mờ; InvalidArgumentException khi windowSize vượt quá maxBatchSize() của engine.

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

Ném hoặc thất bại với: InvalidArgumentException khi windowSize hoặc checkpointIntervalJobs dưới 1. $runId là định danh lần chạy ổn định, single-writer làm khóa cho việc resume checkpoint.

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

$scheme đặt tên scheme target mà committer này phục vụ (ví dụ s3 hoặc gcs); $clock ở đây là Psr\Clock\ClockInterface.

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

Ném hoặc thất bại với: UnsupportedTargetException khi scheme không khớp; RenderManifestException khi target key không an toàn theo kiểu container-relative; CommitIntegrityException khi sha-256 được khai báo không khớp với bytes; OutputCommitConflictException (SPEC-COMMIT-409) khi bytes khác biệt mà không có overwrite; RuntimeException khi cuộc đua create không thể hội tụ sau 5 lần thử dưới sự đột biến đồng thời.

Bề mặt adapter tối giản mà một tích hợp S3/GCS trực tiếp triển khai:

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() phải là một conditional create nguyên tử thực thụ (If-None-Match: * trên S3, ifGenerationMatch: 0 trên GCS) và trả về true chỉ khi lệnh gọi này đã ghi object. put() là overwrite vô điều kiện, chỉ dùng khi manifest yêu cầu overwrite.

JobCompletionEmitterInterfaceFilesystemOutboxEmitter

Phần tiêu đề “JobCompletionEmitterInterface và FilesystemOutboxEmitter”
public function emit(JobTerminalEvent $event): void;

Emitter là tùy chọn trên bộ xử lý. Các event chỉ kích hoạt sau khi một item được hoàn tất một cách bền vững. FilesystemOutboxEmitter là triển khai bền vững được cung cấp:

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

Ném hoặc thất bại với: InvalidArgumentException khi thư mục không tồn tại; emit() ném RuntimeException nếu một event không thể được JSON-encode. hasEvent(string $eventId): bool kiểm tra outbox; count(): int báo cáo các event chưa được phân phối.

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 là xác định — runId:sourceOffset:idempotencyKey:status — đó là điều làm cho việc dedup outbox trở nên khả thi. toArray() serialize event để truyền tải; nó không mang PDF bytes. JobTerminalStatus là một string enum: Committed (committed), DeadLettered (dead_lettered), Skipped (skipped).

Các counter bất biến được trả về bởi process(): runId, sourceRead, fastForwardedByCheckpoint, skippedByIdempotency, windows, renderBatchCalls, renderRetries, commitReceipts, deadLettered, checkpointSaves, và finalCommittedOffset (high-watermark cuối cùng của kết quả cuối cùng).

Commit object-storage exactly-once ở dạng riêng lẻ. NullObjectStorageClient trong bộ nhớ thay thế cho adapter S3/GCS của bạn; ngữ nghĩa bạn quan sát được chính là những gì một adapter trực tiếp phải bảo toàn.

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
}

Output mong đợi:

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

Một lần chạy an toàn khi sập đầy đủ: các store Pro bền vững, object-storage committer, một outbox bền vững, và branding được phân giải từ license. Chạy lại cùng một runId sau một lần sập sẽ fast-forward và hội tụ.

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

Ví dụ output (các counter phụ thuộc vào job stream của bạn):

run nightly-invoices-2026-07-03: read=1200 committed=1187 dedup-skipped=13 dead-lettered=0 checkpoints=12 final-offset=1200
  • Single-writer theo mỗi runId là trách nhiệm của bạn. Checkpoint store không có lease hay compare-and-swap. Hai writer đồng thời trên một runId nằm ngoài contract; hãy cưỡng chế tính độc quyền trong scheduler của bạn.
  • crashSafe: true thất bại nhanh trên các collaborator không bền vững. Các store committer, checkpoint, idempotency, và dead-letter đều phải triển khai marker DurableCapability, nếu không process() sẽ ném StreamProcessorException nêu tên các đối tượng vi phạm. Keyed state store được miễn trừ có chủ đích: keyed state bị mất được tính lại từ checkpoint trở đi.
  • windowSize phải vừa với engine. Một cửa sổ lớn hơn maxBatchSize() sẽ ném InvalidArgumentException trước khi bất kỳ công việc nào bắt đầu.
  • Một commit mập mờ thì hủy bỏ; một conflict thì không. SPEC-COMMIT-409 là một conflict cuối cùng xác định: item bị dead-letter và lần chạy tiếp tục. Bất kỳ commit failure nào khác đều mập mờ: prefix đã hoàn tất được checkpoint và lần chạy ném để resume an toàn.
  • Render failure không bao giờ hủy bỏ lần chạy. Một kết quả Failed theo item, một ngân sách retry đã cạn, hoặc các bytes evaluation không thể branding đều dead-letter item đó và tiếp tục.
  • Các giá trị jobId trùng là an toàn; công việc trùng được khóa theo idempotencyKey. Các kết quả tương quan với item theo source offset duy nhất, không bao giờ theo jobId. Một idempotency key trùng được nhận biết ngay cả trong cùng một khoảng barrier, trước khi re-render.
  • Tính bền vững của emitter quyết định ngữ nghĩa event. Một emitter callback thuần chỉ là observer và at-most-once qua một lần sập. FilesystemOutboxEmitter khiến outbox bền vững và có khóa dedup; khi đó relay delivery là at-least-once, và exactly-once downstream yêu cầu consumer dedup trên eventId. Thư mục của nó (như của mọi filesystem store) phải tồn tại sẵn, nếu không constructor sẽ ném InvalidArgumentException.
  • Các Skipped event mặc định tắt. Đặt emitSkippedCompletions: true để cũng phát một terminal event Skipped cho các item bị short-circuit bởi dedup.
  • Output key fail closed. commit() khẳng định lại target key an toàn theo kiểu container-relative: không có traversal .., không có escape tuyệt đối, không có null byte, không có scheme stream-wrapper nhúng, và không có dấu hai chấm (thứ đóng vector NTFS alternate-data-stream). Các key không an toàn ném trước bất kỳ lệnh gọi storage nào.
  • Tính toàn vẹn được xác minh lại tại ranh giới. Committer tính lại sha-256 trên chính các bytes thực và từ chối một mismatch bằng CommitIntegrityException, nên một handoff hỏng không thể lọt xuống một cách âm thầm.
  • Các event không mang nội dung document. JobTerminalEvent và các hàng outbox chỉ giữ các định danh, digest, timestamp, và chuỗi lỗi. Các thông báo lỗi có thể lặp lại các chẩn đoán engine; hãy làm sạch chúng, và bất kỳ scheme jobId nào định danh tenant, trước khi chuyển các file outbox đến các sink của bên thứ ba.
  • Output evaluation không bao giờ được publish khi chưa branding. Khi branding là bắt buộc và không thể áp dụng, item bị dead-letter thay vì commit.
  • Exactly-once liên writer chỉ mạnh bằng adapter của bạn. Nếu putIfAbsent() không phải là một conditional write nguyên tử thực thụ, đảm bảo này suy giảm xuống ngữ nghĩa single-writer. Thông tin xác thực object-store và policy bucket là mối quan tâm của host; module không bao giờ quản lý chúng.

Không có tiêu chuẩn công bố nào định nghĩa hành vi của module này. Các đảm bảo exactly-once, checkpoint, và outbox trên trang này là các contract kỹ thuật của API NextPDF Enterprise, được nêu ở đây như hành vi có thể quan sát từ bên ngoài — chúng không phải là sự tuân thủ, hay chứng nhận đối với, bất kỳ tiêu chuẩn nào. Việc sử dụng nội bộ SHA-256 làm integrity digest cũng chỉ là hạ tầng, không phải một tuyên bố tuân thủ. Như khắp mọi nơi trong NextPDF: support không phải là conformance, và conformance không phải là certification. NextPDF không nắm giữ chứng nhận nào và không cấp chứng nhận nào; liệu một deployment xây trên module này có đáp ứng các nghĩa vụ pháp lý hay hợp đồng của bạn hay không là một quyết định cho các assessor của bạn.

Các terminal event không phải là một phần của transaction checkpoint: việc phát chạy sau checkpoint.save(), nên outbox giữ mọi event đã phát một cách bền vững nhưng không phải là một sổ cái đầy đủ qua các lần sập. Các object đã commit vẫn là nguồn sự thật.

  • Mọi source offset ở mức finalCommittedOffset hoặc thấp hơn đều đã đạt đúng một kết quả cuối cùng: Committed, Skipped, hoặc DeadLettered.
  • Các item được hoàn tất theo thứ tự source-offset; thứ tự barrier là commit, lưu checkpoint, flush idempotency-mark, rồi phát event.
  • Chạy lại một lần chạy với cùng runId không bao giờ publish trùng: các offset đã checkpoint fast-forward, và các re-commit giống hệt bytes là các no-op được so sánh digest với idempotentReuse: true.
  • Một object mới chỉ bao giờ được tạo qua conditional create nguyên tử; các bytes khác biệt tại một key đã bị chiếm mà không có overwrite là một dead-letter SPEC-COMMIT-409 xác định, không bao giờ là một clobber.
  • Một commit mập mờ checkpoint prefix đã hoàn tất và hủy bỏ với StreamProcessorException; offset thất bại không được tiến lên.
  • Một lần chạy crashSafe từ chối các collaborator committer, checkpoint, idempotency, hoặc dead-letter không bền vững trước khi đọc bất kỳ input nào.
  • Các event id là một hàm thuần túy của run, offset, idempotency key, và status, nên một outbox bền vững giữ tối đa một hàng cho mỗi event.

NextPDF Core render một document mỗi lần thông qua writer và contract render-manifest — xem Writer. Chỉ riêng Core thì không có các job stream bền vững, không có resume checkpoint, không có dedup idempotency, không có commit object-storage, và không có outbox terminal-event. NextPDF Pro thêm render engine bền vững, chạy đồng thời và các store filesystem single-host (Stream trong Pro). Nửa liên máy chủ — bộ xử lý này, object-storage committer, và outbox bền vững — yêu cầu NextPDF Enterprise.

Trang này chỉ ghi lại hành vi có thể quan sát từ bên ngoài và bề mặt public API được hỗ trợ. Các đường dẫn namespace nội bộ, các class trợ giúp, các bảng cơ chế, các tên file runbook, và các tiền tố ticket nằm ngoài phạm vi.