Hacktoberfest 2026:メンテナが10月に向けて印を付けた、オープンで初心者向けの issue。 Hacktoberfest の issue を見る

RecordBatchReader on an open-ended async stream withholds each record batch until the next chunk arrives

オープン 初心者向け
#490 コメント 0 件 リアクション 0 件 担当者 0 名 GitHub で見る

メンテナーはふだん 1 日以内に返信

まだ誰も着手していません。

評価

難易度
2/5
見積もり時間
1〜3時間
初心者へのやさしさ
84/100
issue の種類
バグ
明瞭さ
明確に書かれている
活発さ
活発
技術スタック
typescript
領域
backend

調査の方向性

src/io/adapters.ts の fromAsyncIterable から始め、次に fromIterable、fromDOMStream、fromNodeStream の対応するループを調べます。チャンクごとに 1 つの IPC メッセージを使って Node 24 の再現を実行し、完全なバイト列が到着した時点でバッチが返されることを、ソースが開いたままの状態での最後のバッチも含めて確認します。

索引モデルが issue の本文から書いたものです。

説明

Describe the bug, including details regarding any error messages, version, and platform.

When RecordBatchReader.from() reads an async byte source that stays open (a WebSocket, a network stream) and a chunk ends exactly at the end of a record batch message, that batch is not yielded until the next chunk arrives. On a live stream the latest batch is always one step behind, and the last batch never appears until the stream ends.

The cause looks like an off-by-one in src/io/adapters.ts, fromAsyncIterable (quoted, lightly simplified, from the published 21.2.0 build):

do {
    ({ cmd, size } = yield byteRange());   // serve the request; receive the next one
} while (size < bufferLength);             // BUG: stops when the next request is exactly
                                           // the bytes still buffered, then waits for
                                           // the next chunk before serving them

The reader asks for a message body with size === bufferLength (exactly the bytes left in the chunk). The inner loop exits because size < bufferLength is false, and the outer loop awaits the next chunk before serving bytes that are already buffered.

A likely fix:

-    } while (size < bufferLength);
+    } while (size <= bufferLength);

The same while (size < bufferLength) appears in fromIterable, fromDOMStream and fromNodeStream, so browser ReadableStream sources are probably affected too (not tested).

To reproduce

apache-arrow 21.2.0, Node 24. Three one-batch IPC messages, one per chunk, one chunk per second, source never ends:

// make the inputs with pyarrow:
//   s = pa.schema([("x", pa.int64())]); open("schema.bin","wb").write(s.serialize().to_pybytes())
//   for i in range(3): open(f"b{i}.bin","wb").write(
//       pa.record_batch([pa.array([i*10, i*10+1])], schema=s).serialize().to_pybytes())
const { RecordBatchReader } = require("apache-arrow");
const fs = require("fs");
const frames = ["schema", "b0", "b1", "b2"].map(f => new Uint8Array(fs.readFileSync(f + ".bin")));
const sleep = ms => new Promise(r => setTimeout(r, ms));
const t0 = Date.now();

async function* source() {          // one message per chunk, like a WebSocket; never ends
  for (const f of frames) { yield f; await sleep(1000); }
  await sleep(1e9);
}

(async () => {
  const reader = await RecordBatchReader.from(source());
  await reader.open();
  for await (const b of reader)
    console.log(`batch ${b.getChildAt(0).get(0)} at ${((Date.now() - t0) / 1000).toFixed(1)} s`);
})();
setTimeout(() => process.exit(0), 5500);

Output:

batch 0 at 2.0 s     <- sent at 1.0 s
batch 10 at 3.0 s    <- sent at 2.0 s
                     <- batch 20, sent at 3.0 s, never yielded

Expected: each batch at the time its chunk arrives (1.0 s, 2.0 s, 3.0 s).

Splitting each chunk so that its last byte arrives as a separate chunk makes every batch arrive on time, which is consistent with the off-by-one above.


This report was prepared with the help of an AI coding assistant (Claude), which traced the cause in the library source and wrote and ran the reproduction.

主要言語
TypeScript
スター
112
フォーク
23
平均マージ
1日 10時間
マージ済み PR(30日)
13

環境構築

はじめの一歩

  1. issue を最後まで読み、次にプロジェクトのコントリビューションガイドを読みます。
  2. 着手することを issue にコメントします — 二人が同じ作業をするのを防げます。
  3. リポジトリをフォークし、ブランチを切って変更します。
  4. issue 番号を参照したプルリクエストを送ります。

apache/arrow-js のほかの issue

apache/arrow-js の issue をすべて見る

似ている issue

TypeScript の issue をもっと見る

新しい issue をメールで受け取る

初心者向けの GitHub issue を短くまとめたダイジェスト。