ストリーミング
大容量データのストリーム保存と読み取り方法を説明します。
型の選択
ストリームを使うキーは StreamValue<T> か Value<T> で定義します。PlainValue<T> では stream() を使えません。
const kvs = UniKvs.config<{
logs: StreamValue<Uint8Array>;
blob: Value<Uint8Array>;
}>()
.appendStorage(new Memory())
.create();
| 型 | set |
get |
stream |
|---|---|---|---|
PlainValue<T> |
対応します。 | 対応します。 | 対応しません。 |
StreamValue<T> |
値またはストリームに対応します。 | 対応しません。 | 対応します。 |
Value<T> |
値またはストリームに対応します。 | 対応します。 | 対応します。 |
schema を指定した場合はチャンクごとに検証します。書き込みチャンクの不良は InvalidInputError を含む PluginOperationAggregateError として報告し、読み取りチャンクの不良は InvalidOutputError として報告します。
書き込み
値そのものか ReadableStream を渡せます。
単一値で保存する
小さなデータは値そのものを渡します。
await kvs.set("logs", new Uint8Array([1, 2, 3]));ストリームで保存する
大きなデータは ReadableStream を渡して逐次書き込みします。
const src = new ReadableStream({
start(controller) {
controller.enqueue(new Uint8Array([1]));
controller.enqueue(new Uint8Array([2]));
controller.close();
},
});
await kvs.set("logs", src);読み取り
stream() は ValueStream を返します。用途に合わせて 3 つから選べます。
チャンクごとに細かく制御する場合に使います。
const reader = (await kvs.stream("logs")).getReader();
while (true) {
const { done, value } = await reader.read();
if (done) {
break;
}
console.log(value);
}短く書く場合に使います。
for await (const chunk of await kvs.stream("logs")) {
console.log(chunk);
}ロック解放を自動化する場合に使います。
await using (const valueStream = await kvs.stream("logs")) {
const reader = valueStream.getReader();
const { value } = await reader.read();
console.log(value);
}ストレージ別の注意
- ファイル・S3・OPFS はバイト列のストリーム保存に適しています。詳細は各パッケージリファレンスを参照してください。
- IndexedDB は任意データを保存できますが、ストリームの制約を確認してください。
- トランスフォーマーによってはストリーム変換に未対応の場合があります。対応状況は各パッケージリファレンスを参照してください。
キャンセル
長時間の転送は AbortSignal で中断できます。
const ac = new AbortController();
setTimeout(() => ac.abort(), 1000);
await kvs.set("logs", src, { signal: ac.signal });