Hacktoberfest 2026:维护者为十月标记出来的 issue,仍然开放、适合新手。 浏览 Hacktoberfest issue

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

未关闭 适合新手
#490 0 条评论 0 个 reaction 已指派 0 人 在 GitHub 查看

维护者通常 1 天内回复

还没有人认领这个 Issue。

评估

难度
2/5
预计耗时
1-3 小时
新手友好度
84/100
Issue 类型
缺陷
描述清晰度
描述清楚
活跃度
活跃
技术栈
typescript
领域
backend

调研方向

从 src/io/adapters.ts 中的 fromAsyncIterable 开始,然后检查 fromIterable、fromDOMStream 和 fromNodeStream 中对应的循环。使用每个 chunk 一条 IPC 消息运行 Node 24 复现,并验证在一个 batch 的完整字节到达时会产生该 batch,包括源保持打开状态时的最后一个 batch。

由索引模型根据 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 小时
30 天内合并 PR
13

环境准备

从这里开始

  1. 先读完整个 Issue,再读项目的贡献指南。
  2. 在 Issue 下留言说明你要接手 —— 这能避免两个人做同样的事。
  3. Fork 仓库,在一个分支上完成修改。
  4. 提交 Pull Request,并在描述里引用这个 Issue 编号。

apache/arrow-js 的其他 Issue

查看 apache/arrow-js 的全部 Issue

相似的 Issue

更多 TypeScript Issue

把新 issue 发到你的邮箱

精选适合新手参与的 GitHub issue 摘要。