RecordBatchReader on an open-ended async stream withholds each record batch until the next chunk arrives
Maintainers usually reply within 1 day
Nobody has claimed this yet.
Assessment
- Difficulty
- 2/5
- Estimated time
- 1-3 hours
- Newbie friendliness
- 84/100
- Issue type
- Bug
- Clarity
- Clearly specified
- Activity status
- Active
- Tech stack
- typescript
- Domain
- backend
Research direction
Start in src/io/adapters.ts at fromAsyncIterable, then inspect the matching loops in fromIterable, fromDOMStream, and fromNodeStream. Run the Node 24 reproduction with one IPC message per chunk and verify that batches are yielded when their complete bytes arrive, including the final batch while the source remains open.
Written by the indexing model from the issue text.
Description
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.
- Dominant language
- TypeScript
- Stars
- 112
- Forks
- 23
- Avg merge
- 1d 10h
- Merged PRs (30d)
- 13
Getting set up
First steps
- Read the whole issue, then the project's contributing guide.
- Comment on the issue to say you are picking it up — it saves two people doing the same work.
- Fork the repository and make your change on a branch.
- Open a pull request that references the issue number.
More from apache/arrow-js
-
Difficulty 3/5 1-2 days Newbie friendliness 68/100
apache/arrow-js#484 · 1 comment · 1 reaction ·
Maintainers usually reply within 1 day
-
Difficulty 3/5 1-2 days Newbie friendliness 55/100
apache/arrow-js#468 · 2 comments ·
Maintainers usually reply within 1 day
-
Difficulty 3/5 1-2 days Newbie friendliness 35/100
Maintainers usually reply within 1 day
-
Difficulty 5/5 Over a week Newbie friendliness 35/100
apache/arrow-js#423 · 1 reaction ·
Maintainers usually reply within 1 day
-
Difficulty 3/5 1-2 days Newbie friendliness 50/100
Maintainers usually reply within 1 day
Similar issues
-
Mend: dependency security vulnerability untriaged
Difficulty 2/5 1-3 hours Newbie friendliness 68/100
opensearch-project/security-dashboards-plugin#2545 ·
Maintainers usually reply within 1 day
-
Add: Dream TR SDOpencheck:passed streams:add
Difficulty 2/5 1-3 hours Newbie friendliness 70/100
Maintainers usually reply within 1 day
-
doctor integrity sample scans soft-deleted pages on Postgres (batch path has no deleted_at filter)Open
Difficulty 2/5 1-3 hours Newbie friendliness 88/100
Maintainers usually reply within 1 day
-
Difficulty 1/5 Under an hour Newbie friendliness 72/100
SocialGouv/egapro#4672 · 1 comment ·
Maintainers usually reply within 2 days
-
area:agents area:tui bug
Difficulty 2/5 1-3 hours Newbie friendliness 74/100
anthropics/claude-code#98358 ·
Maintainers usually reply within 1 day