コンテンツにスキップ
getnextpdf.com

Enterprise エディション

ストリーム: ドキュメントジョブ処理

NextPDF\Enterprise\Stream\DocumentJobStreamProcessor は、レンダーマニフェストのストリームを、永続的で追跡可能な結果へと変換します。iterable<RenderManifest> をジェネレーターとして消費し、Pro レンダーエンジンを通じて上限付きウィンドウをレンダリングし、すべてのジョブをソース順に確定します。各ジョブは、ちょうど 1 つの終端状態で終わります。出力のコミット済み、すでにコミット済みと認識、またはデッドレター化のいずれかです。進捗はチェックポイントされるため、クラッシュした実行は何も再発行せずに再開します。

Stream のストーリーは 2 つのエディションに分かれており、この分割は意図的なものです。Pro は、永続的で並行なレンダーエンジンと、ローカルの単一ホストのファイルシステムストアを提供します。これがインプロセス側の半分です。Enterprise は、このドキュメントジョブストリームプロセッサーに加えて、ホスト境界をまたぐ部品を提供します。オブジェクトストレージコミッター(ObjectStorageCommitter)と、終端イベントの永続的なアウトボックス(FilesystemOutboxEmitter)です。Pro Stream ページは、反対側から同じ境界を述べています。

この機能は NextPDF Enterprisenextpdf/enterprise)に含まれ、Enterprise ティアのライセンスエンベロープで有効化されます。そのエンタイトルメントを持たないデプロイでは、この機能のクラスはロードされません。エディションを比較してライセンスを取得

Terminal window
composer require nextpdf/enterprise

このページのクラスは NextPDF\Enterprise\Stream および NextPDF\Enterprise\Stream\Storage の下にあります。これらは NextPDF\Pro\Stream の凍結された Pro コントラクト、すなわちエンジン、コミッター、チェックポイント、冪等性、リトライ、デッドレターの各インターフェースを利用します。

プロセッサーの役割はデリバリーセマンティクスであり、レンダリングではありません。マニフェストストリームを、エンジンのバッチサイズを超えないソースオフセットウィンドウにグループ化します。各ウィンドウは RenderEngineInterface::renderBatch() を通じてレンダリングされ、アイテムごとのタイムアウトに対して上限付きで決定論的なリトライが行われます。その後、すべてのアイテムがソースオフセット順に終端結果へと確定されます。

exactly-once の境界はアイテムごとであり、それはコーディネーションではなくコミッターに固定されています。永続的な進捗は 1 始まりのオフセットハイウォーターマークです。チェックポイント以下のすべてのオフセットは終端結果に到達しています。バリアの順序は固定されています。バイト列をコミットし、ウォーターマークを前進させ、チェックポイントを保存し、バッファされた冪等性マークをフラッシュし、その後に終端イベントを発行します。コミットとチェックポイントの間でのクラッシュは、再開時に冪等に再コミットされます。コミッターがダイジェストを比較するためです。チェックポイント後のクラッシュはオフセットを越えて早送りするため、何も二重に発行されません。

ObjectStorageCommitter は、最小限の ObjectStorageClientInterface を通じて、オブジェクトストアに対して Pro の OutputCommitterInterface を実装します。ターゲットの container はバケットであり、その key はオブジェクトキーです。同一のバイト列の再コミットは、ダイジェスト比較による no-op です。overwrite なしで相違するバイト列は SPEC-COMMIT-409 コンフリクトを発生させます。新規オブジェクトは、アトミックな putIfAbsent() 条件付き書き込みによってのみ作成されます。そのレースに敗れると、上限付きの再読み込みと解決のループがトリガーされます。したがって、ライターをまたぐ exactly-once は、アダプターの putIfAbsent() が真の条件付き書き込みである限りにおいてのみ成立します。S3 では If-None-Match: *、GCS では ifGenerationMatch: 0 です。今回のサイクルでは、インターフェースに加えてインメモリの NullObjectStorageClient を提供します。本番の S3/GCS アダプターはホスト提供です。

終端イベントは、ダウンストリームシステムに対してループを閉じます。チェックポイントバリアの後、プロセッサーは確定した各ジョブについて JobTerminalEvent の発行を試みます。識別子、ステータス、レシート、エラー詳細、試行回数を含み、PDF バイト列は一切含みません。単純なコールバックエミッターでは、発行は at-most-once です。チェックポイント後のイベントは、クラッシュ再開時にスキップされる可能性があります。FilesystemOutboxEmitter は、emit() が実行されると各イベントを永続化します。各イベントは、決定論的な eventId のハッシュで名付けられた 1 つのアトミックな JSON ファイルであるため、再開後の再発行は冪等であり、リレーは at-least-once で配信し、コンシューマーは eventId で重複排除します。いずれにせよ、1 つの境界が残ります。発行はチェックポイントバリアの後に起こるため、checkpoint.save()emit() の間のクラッシュは、そのアイテムの終端イベントを再開時にスキップします。完全なイベント台帳を必要とするダウンストリームシステムは、アウトボックス単独ではなく、コミット済みオブジェクト(ストアが信頼できる情報源です)と突き合わせて照合すべきです。

ライセンスは出力パスに組み込まれています。withBrandingFromLicense() ファクトリは、実行ごとに 1 回、ライセンスから評価版ブランディング戦略を解決します。有償ライセンスは恒等変換に解決されます。評価版ライセンスまたはライセンス欠如の場合は、コミットされるすべてのドキュメントにウォーターマークを付与し、ブランディングできないドキュメントはデッドレター化されます。プロセッサーは、ブランディングされていない評価版バイト列を決してコミットしません。

要となる決定は、exactly-once が分散ロックやコンセンサスではなく、コミッターのダイジェスト比較された条件付き作成オブジェクトに依拠するという点です。オブジェクトストアの条件付き書き込みは、この設計が必要とする唯一のアトミックプリミティブであり、それ以外のすべては失敗して回復することが許容されます。だからこそ、レンダーエンジンは副作用のない状態を保たなければならず、キー付き状態と実行中の重複排除キャッシュは再計算可能な高速化として扱われ、曖昧なコミットは推測する代わりに実行を中止します。再開パスは同じダイジェスト比較を通じて収束するのです。また、これが 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

$clockSymfony\Component\Clock\ClockInterface です(リトライのバックオフはこれを通じてスリープします)。null ライセンスはフェイルクローズドで評価版ブランディングに解決されます。

単一のエントリポイントが 1 つの実行を処理し、そのカウンターを返します。

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

スロー、または失敗する条件: クラッシュ安全性の前提条件が満たされない場合(crashSafe 実行で非永続的なコラボレーターを使用した場合)、またはコミットが曖昧な場合は NextPDF\Enterprise\Stream\Exception\StreamProcessorExceptionwindowSize がエンジンの maxBatchSize() を超える場合は InvalidArgumentException

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

スロー、または失敗する条件: windowSize または checkpointIntervalJobs が 1 未満の場合は InvalidArgumentException$runId は、チェックポイント再開のキーとなる、安定した単一ライターの実行識別子です。

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

$scheme は、このコミッターが処理するターゲットスキームを指定します(例えば s3gcs)。ここでの $clockPsr\Clock\ClockInterface です。

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

スロー、または失敗する条件: スキームの不一致では UnsupportedTargetException。ターゲットキーがコンテナ相対で安全でない場合は RenderManifestException。宣言された sha-256 がバイト列と一致しない場合は CommitIntegrityExceptionoverwrite なしで相違するバイト列では OutputCommitConflictExceptionSPEC-COMMIT-409)。並行変更下で作成レースが 5 回の試行後も収束できない場合は RuntimeException

本番の 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() は、真のアトミックな条件付き作成(S3 では If-None-Match: *、GCS では ifGenerationMatch: 0)でなければならず、この呼び出しがオブジェクトを書き込んだ場合にのみ true を返します。put() は、マニフェストが overwrite を要求した場合にのみ使用される、無条件の上書きです。

public function emit(JobTerminalEvent $event): void;

エミッターはプロセッサーではオプションです。イベントは、アイテムが永続的に確定された後にのみ発火します。FilesystemOutboxEmitter は、提供される永続的な実装です。

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

スロー、または失敗する条件: ディレクトリが存在しない場合は InvalidArgumentException。イベントを JSON エンコードできない場合、emit()RuntimeException をスローします。hasEvent(string $eventId): bool はアウトボックスを確認します。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)、これがアウトボックスの重複排除を可能にします。toArray() はイベントを転送用にシリアライズします。PDF バイト列は含みません。JobTerminalStatus は文字列 enum です。Committedcommitted)、DeadLettereddead_lettered)、Skippedskipped)。

process() が返す不変のカウンター。runIdsourceReadfastForwardedByCheckpointskippedByIdempotencywindowsrenderBatchCallsrenderRetriescommitReceiptsdeadLetteredcheckpointSaves、そして finalCommittedOffset(最終的な終端ハイウォーターマーク)。

exactly-once のオブジェクトストレージコミットを単独で示します。インメモリの 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

完全なクラッシュセーフ実行。永続的な Pro ストア、オブジェクトストレージコミッター、永続的なアウトボックス、そしてライセンスで解決されたブランディングを使用します。クラッシュ後に同じ 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 もありません。1 つの runId に対する 2 つの並行ライターはコントラクトの範囲外です。排他性はスケジューラーで強制してください。
  • crashSafe: true は非永続的なコラボレーターに対して即座に失敗します。 コミッター、チェックポイント、冪等性、デッドレターの各ストアはすべて DurableCapability マーカーを実装しなければならず、さもなければ process() は違反者を名指しする StreamProcessorException をスローします。キー付き状態ストアは意図的に免除されています。失われたキー付き状態は、チェックポイント以降から再計算されます。
  • windowSize はエンジンに収まらなければなりません。 maxBatchSize() より大きいウィンドウは、いかなる作業も始まる前に InvalidArgumentException をスローします。
  • 曖昧なコミットは中止しますが、コンフリクトは中止しません。 SPEC-COMMIT-409 は決定論的な終端コンフリクトです。アイテムはデッドレター化され、実行は継続します。それ以外のコミット失敗は曖昧です。確定済みの接頭部分はチェックポイントされ、実行は安全な再開のためにスローします。
  • レンダリングの失敗が実行を中止することはありません。 アイテムごとの Failed 結果、使い果たしたリトライ予算、またはブランディング不能な評価版バイト列は、いずれもそのアイテムをデッドレター化して継続します。
  • 重複した jobId の値は安全です。重複作業は idempotencyKey をキーとします。 結果は一意のソースオフセットによってアイテムに対応付けられ、jobId によることは決してありません。重複した冪等性キーは、同じバリア間隔内であっても、再レンダリングの前に認識されます。
  • エミッターの永続性がイベントセマンティクスを決定します。 単純なコールバックエミッターはオブザーバー専用であり、クラッシュをまたぐと at-most-once です。FilesystemOutboxEmitter はアウトボックスを永続的かつ重複排除キー付きにします。その場合、リレー配信は at-least-once となり、ダウンストリームの exactly-once にはコンシューマー側での eventId による重複排除が必要です。そのディレクトリは(すべてのファイルシステムストアと同様に)事前に存在しなければならず、さもなければコンストラクターは InvalidArgumentException をスローします。
  • スキップイベントはデフォルトでオフです。 emitSkippedCompletions: true を設定すると、重複排除でショートサーキットされたアイテムについても Skipped 終端イベントを発行します。
  • 出力キーはフェイルクローズドです。 commit() は、ターゲットキーがコンテナ相対で安全であることを再確認します。.. によるトラバーサルなし、絶対パスによるエスケープなし、ヌルバイトなし、埋め込みストリームラッパースキームなし、そしてコロンなし(これにより NTFS の代替データストリームのベクトルを封じます)。安全でないキーは、いかなるストレージ呼び出しの前にもスローします。
  • 整合性は境界で再検証されます。 コミッターは実際のバイト列に対して sha-256 を再計算し、不一致を CommitIntegrityException で拒否します。したがって、破損した受け渡しが黙って着地することはありません。
  • イベントはドキュメントの内容を含みません。 JobTerminalEvent とアウトボックスの行は、識別子、ダイジェスト、タイムスタンプ、エラー文字列のみを保持します。エラーメッセージはエンジンの診断情報をそのまま反映することがあります。アウトボックスファイルをサードパーティのシンクに送る前に、それらと、テナントを識別しうる jobId のスキームをスクラブしてください。
  • 評価版の出力がブランディングなしで公開されることは決してありません。 ブランディングが必要でありながら適用できない場合、そのアイテムはコミットではなくデッドレター化されます。
  • ライターをまたぐ exactly-once は、アダプターの強さ以上にはなりません。 putIfAbsent() が真のアトミックな条件付き書き込みでない場合、保証は単一ライターセマンティクスへと低下します。オブジェクトストアの認証情報とバケットポリシーはホストの関心事であり、モジュールがそれらを管理することはありません。

このモジュールの動作を定義する公開標準は存在しません。このページの exactly-once、チェックポイント、アウトボックスの保証は、NextPDF Enterprise API のエンジニアリング上のコントラクトであり、ここでは外部から観測可能な動作として述べられています。これらは、いかなる標準への準拠でも、いかなる標準に対する認証でもありません。整合性ダイジェストとしての SHA-256 の内部利用も同様に配管であり、コンプライアンスの主張ではありません。NextPDF のあらゆる箇所と同様に、サポートは準拠ではなく、準拠は認証ではありません。NextPDF はいかなる認証も保持せず、いかなる認証も付与しません。このモジュール上に構築されたデプロイが規制上または契約上の義務を満たすかどうかは、あなたの評価者が判断すべき事項です。

終端イベントはチェックポイントトランザクションの一部ではありません。発行は checkpoint.save() の後に実行されるため、アウトボックスは発行されたすべてのイベントを永続的に保持しますが、クラッシュをまたぐ完全な台帳ではありません。コミット済みオブジェクトが信頼できる情報源であり続けます。

  • finalCommittedOffset 以下のすべてのソースオフセットは、ちょうど 1 つの終端結果に到達しています。CommittedSkipped、または DeadLettered です。
  • アイテムはソースオフセット順に確定されます。バリアの順序は、コミット、チェックポイント保存、冪等性マークのフラッシュ、その後にイベント発行です。
  • 同じ runId で実行を再実行しても、二重発行は決して起こりません。チェックポイントされたオフセットは早送りされ、バイト単位で同一の再コミットは idempotentReuse: true を伴うダイジェスト比較の no-op です。
  • 新規オブジェクトは、アトミックな条件付き作成によってのみ作成されます。overwrite なしで占有済みキーに相違するバイト列は、決定論的な SPEC-COMMIT-409 のデッドレターであり、決して上書きではありません。
  • 曖昧なコミットは、確定済みの接頭部分をチェックポイントし、StreamProcessorException で中止します。失敗したオフセットは前進しません。
  • crashSafe 実行は、いかなる入力を読み込む前にも、非永続的なコミッター、チェックポイント、冪等性、またはデッドレターのコラボレーターを拒否します。
  • イベント ID は、実行、オフセット、冪等性キー、ステータスの純粋関数であるため、永続的なアウトボックスはイベントごとに最大 1 行を保持します。

NextPDF Core は、ライターとレンダーマニフェストコントラクトを通じて、一度に 1 つのドキュメントをレンダリングします。Writer を参照してください。Core 単独では、永続的なジョブストリーム、チェックポイント再開、冪等性による重複排除、オブジェクトストレージコミット、終端イベントアウトボックスのいずれもありません。NextPDF Pro は、永続的で並行なレンダーエンジンと単一ホストのファイルシステムストアを追加します(Pro の Stream)。ホストをまたぐ半分、すなわちこのプロセッサー、オブジェクトストレージコミッター、永続的なアウトボックスには、NextPDF Enterprise が必要です。

このページは、外部から観測可能な動作と、サポートされる公開 API サーフェスのみを文書化します。内部のネームスペースパス、ヘルパークラス、メカニズムのテーブル、ランブックのファイル名、チケットの接頭辞は対象外です。