Hacktoberfest 2026: những issue maintainer đã đánh dấu cho tháng Mười, đang mở và phù hợp người mới. Xem issue Hacktoberfest

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

Đang mở Phù hợp với người mới
#490 0 bình luận 0 reaction 0 người được giao Xem trên GitHub

Maintainer thường phản hồi trong vòng 1 ngày

Chưa có ai nhận issue này.

Đánh giá

Độ khó
2/5
Thời gian dự kiến
1-3 giờ
Mức phù hợp với người mới
84/100
Loại issue
Lỗi
Độ rõ ràng
Đặc tả rõ ràng
Mức độ hoạt động
Sôi nổi
Công nghệ
typescript
Lĩnh vực
backend

Hướng nghiên cứu

Bắt đầu trong src/io/adapters.ts tại fromAsyncIterable, sau đó kiểm tra các vòng lặp tương ứng trong fromIterable, fromDOMStream và fromNodeStream. Chạy bản tái hiện trên Node 24 với một thông báo IPC cho mỗi chunk và xác minh rằng các batch được trả về khi toàn bộ byte của chúng đã đến, bao gồm cả batch cuối cùng trong khi nguồn vẫn mở.

Do mô hình lập chỉ mục viết ra từ nội dung của issue.

Mô tả

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.

Ngôn ngữ chính
TypeScript
Star
112
Fork
23
Merge trung bình
1 ngày 10 giờ
Pull request đã merge (30 ngày)
13

Chuẩn bị môi trường

Bắt đầu từ đâu

  1. Đọc hết issue, rồi đọc hướng dẫn đóng góp của dự án.
  2. Bình luận trên issue rằng bạn sẽ nhận — tránh hai người làm cùng một việc.
  3. Fork repository và làm thay đổi trên một nhánh.
  4. Mở pull request có tham chiếu số hiệu của issue.

Issue khác của apache/arrow-js

Tất cả issue của apache/arrow-js

Issue tương tự

Thêm issue về TypeScript

Nhận issue mới trong hộp thư của bạn

Bản tóm tắt ngắn những issue GitHub phù hợp với người mới.