RecordBatchReader on an open-ended async stream withholds each record batch until the next chunk arrives
维护者通常 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
环境准备
- 提供 Dockerfile 或 Docker Compose 文件
- 有 Pull Request 模板
- 阅读贡献指南
从这里开始
- 先读完整个 Issue,再读项目的贡献指南。
- 在 Issue 下留言说明你要接手 —— 这能避免两个人做同样的事。
- Fork 仓库,在一个分支上完成修改。
- 提交 Pull Request,并在描述里引用这个 Issue 编号。
apache/arrow-js 的其他 Issue
-
难度 3/5 1-2 天 新手友好度 68/100
apache/arrow-js#484 · 1 条评论 · 1 个 reaction ·
维护者通常 1 天内回复
-
难度 3/5 1-2 天 新手友好度 55/100
维护者通常 1 天内回复
-
难度 3/5 1-2 天 新手友好度 35/100
维护者通常 1 天内回复
-
难度 5/5 一周以上 新手友好度 35/100
apache/arrow-js#423 · 1 个 reaction ·
维护者通常 1 天内回复
-
难度 3/5 1-2 天 新手友好度 50/100
维护者通常 1 天内回复
相似的 Issue
-
area/core status/need-triage
难度 2/5 1-3 小时 新手友好度 88/100
google-gemini/gemini-cli#29602 ·
维护者通常 1 天内回复
-
area: backend enhancement priority: low
难度 2/5 1-3 小时 新手友好度 88/100
snapotter-hq/SnapOtter#1879 ·
维护者通常 1 天内回复
-
难度 2/5 1-3 小时 新手友好度 86/100
Tencent/BrowserSkill#390 ·
维护者通常 1 天内回复
-
good first issue status: needs triaging type: bug version: 2.0
难度 2/5 1-3 小时 新手友好度 85/100
medusajs/medusa#17094 · 2 条评论 ·
维护者通常 1 天内回复
-
难度 2/5 1-3 小时 新手友好度 86/100
维护者通常 1 天内回复