ReadableStreamをfor await...ofで読む
ReadableStreamが非同期イテラブルになり、for await...ofでチャンクを順に読めるようになりました。getReader()とread()のループとの違い、途中で抜けるとストリームがキャンセルされる既定の挙動とpreventCancelでの回避、fetchのレスポンスを行単位で処理する例を扱います。
- Chrome124
- Edge124
- Firefox110
- Safari27
はじめに
fetch()のレスポンスを少しずつ処理したいとき、これまではresponse.body.getReader()でリーダーを取り、read()がdone: trueを返すまでwhile (true)で回す書き方が定番でした。動きはしますが、ループの終了条件は自分で書く必要があり、途中で抜けるときはreleaseLock()やcancel()の呼び分けも引き受けることになります。
Streams APIの仕様には、ReadableStreamを非同期イテラブルとして扱う定義が2020年ごろから入っています。Firefox 110とChrome 124が実装したあとも、for awaitで書いたコードがSafariでだけTypeErrorになる状態が長く続きました。Safari 27の対応により、非同期イテラブルなストリームはBaseline 2026に加わりました。
リーダーのループを置き換える
ReadableStreamが非同期イテラブルになったので、for await...ofにそのまま渡せます。
// これまで: リーダーを取って done になるまで回す
const reader = response.body.getReader();
while (true) {
const { done, value } = await reader.read();
if (done) break;
process(value);
}
// for await...of: 終了条件とリーダーの管理が消える
for await (const chunk of response.body) {
process(chunk);
}ループの実行中に別の場所からgetReader()を呼ぶとTypeErrorになります。ループの中ではストリームが内部でリーダーにロックされているためで、ロックはループを抜けた時点で解放されます。ストリームがエラーで終わったときはawaitの位置でそのエラーが投げられるので、try...catchでループごと囲めば受け取れます。
ループを抜けるとストリームはキャンセルされる
for await...ofをbreakやreturn、例外で途中終了すると、既定ではストリームがキャンセルされます。データソースにはcancel()が伝わり、以降そのストリームからは読めません。先頭の数チャンクだけ見てあとは捨てるなら既定のままで十分ですが、続きを別のリーダーで読みたいときは行き詰まります。
続きを読みたいなら、values()にpreventCancel: trueを渡します。values()はfor awaitが内部で呼んでいる非同期イテレータを返すメソッドで、オプションを付けたいときに明示的に呼びます。
const stream = getStream();
// 先頭のチャンクだけ確認してループを抜ける
for await (const chunk of stream.values({ preventCancel: true })) {
inspect(chunk);
break;
}
// キャンセルされていないので、続きをリーダーで読める
const reader = stream.getReader();preventCancelを付けずにbreakすると、そのあとのgetReader()自体は成功しますが、read()は即座にdone: trueを返します。ストリームがキャンセル済みで、データがすでに捨てられているためです。
fetchのレスポンスを行単位で処理する
サーバーが1行に1件ずつJSONを書いたログを流してくるとき、全部届くのを待たず、届いた行から順に処理したいとします。response.bodyをそのままfor awaitで読むと、届くのはバイト列のUint8Arrayで、しかもチャンクの切れ目は行の切れ目と一致しません。1つのチャンクの末尾に行の途中までが入り、続きが次のチャンクの先頭に来ます。
そこで2つの手当てをします。TextDecoderStreamをpipeThrough()でつないでバイト列を文字列のストリームに変え、行の途中で切れた末尾は次のチャンクに持ち越します。
const response = await fetch('/api/events');
const texts = response.body.pipeThrough(new TextDecoderStream());
// 前のチャンクの末尾に残った、行の途中までの文字列
let rest = '';
for await (const text of texts) {
const parts = (rest + text).split('\n');
// 最後の要素は行の途中で切れている可能性があるので、次のチャンクに回す
rest = parts.pop() ?? '';
for (const line of parts) {
if (line !== '') handle(JSON.parse(line));
}
}
// 最後のチャンクの末尾に残った行を処理する
if (rest !== '') handle(JSON.parse(rest));restには前のチャンクで完結しなかった行の先頭部分が入っていて、次のチャンクの先頭とつないでから改行で分割します。分割結果の最後の要素は次のチャンクへ持ち越し、それ以外は完結した行として処理します。ストリームが終わったあとにrestに残っているものは、改行で終わらなかった最後の行です。
途中で読むのをやめたいときは、AbortControllerでfetch()を中断します。ストリームがエラー状態になり、for awaitのawaitでAbortErrorが投げられるので、中断を正常終了として扱うならループをtry...catchで囲んでAbortErrorだけ握りつぶします。
おわりに
ReadableStreamをfor await...ofで読めるようになり、リーダーの取得とdoneの判定を自分で書く必要がなくなりました。ループを途中で抜けるとストリームがキャンセルされるのが既定で、続きを読みたいときだけvalues({ preventCancel: true })を使います。Safariの未対応で足踏みしていた機能なので、これまでリーダーのループを書いてきた場所は、使うブラウザの範囲を確かめた上で置き換えられます。