Is your feature request related to a problem or challenge?
When a CSV or NDJSON file is split into byte ranges (datafusion.optimizer.repartition_file_scans, on by default), each range is read through AlignedBoundaryStream. The first request covers [start - 1, end + 16 KiB). For the first range it starts at byte 0, and it never goes past the end of the file. If the record crossing end doesn't finish inside that window, ScanningLastTerminator fetches the rest in 16 KiB GETs (END_SCAN_LOOKAHEAD, the last one clamped to the file size), each issued only after the previous one is exhausted:
|
Phase::ScanningLastTerminator => { |
|
match this.inner.poll_next_unpin(cx) { |
|
Poll::Pending => return Poll::Pending, |
|
Poll::Ready(None) => { |
|
// Inner exhausted. Issue the next overflow GET if |
|
// the file has not been fully read yet. |
|
let pos = this.abs_pos(); |
|
if pos < this.file_size { |
|
let fetch_end = pos |
|
.saturating_add(END_SCAN_LOOKAHEAD) |
|
.min(this.file_size); |
|
let store = Arc::clone(&this.store); |
|
let location = this.location.clone(); |
|
this.inner = get_stream(store, location, pos..fetch_end) |
|
.try_flatten_stream() |
|
.boxed(); |
|
continue; |
|
} |
|
this.phase = Phase::Done; |
|
return Poll::Ready(None); |
|
} |
Say R is the number of bytes from the range end up to and including that record's newline, and the file doesn't end before it. Then the range needs ceil((R - 16 KiB) / 16 KiB) extra sequential GETs when R > 16 KiB, and none otherwise. A range ending just after the start of a 1 MiB record makes 64 round trips (the initial GET plus 63 overflow GETs) before that partition can finish. Records of 16 KiB or less never trigger this.
This is the retry case #8723 mentions (no newline found because the overfetched range was too small). #20823 (JSON) and #22962 (CSV) delivered #8723's main goal, one GET per range instead of three, whenever the initial lookahead reaches the newline. This issue is about the remaining overflow case in the current AlignedBoundaryStream. I couldn't find an open issue or PR for it.
It hits uncompressed NDJSON and newline-delimited CSV scans that the optimizer splits into byte ranges, when records are over 16 KiB (large JSON documents, long text fields) and the store has real per-request latency. Splitting only happens when a file group's total size reaches datafusion.optimizer.repartition_file_min_size (default 1 MiB). Compressed files are never split, CSV with newlines_in_values enabled isn't split, and JSON array files don't support ranged scans. On local disk the extra requests cost almost nothing.
Reproducer. Add this to datafusion/core/tests/datasource/object_store_access.rs, which already has the RequestCountingObjectStore harness. The file has one 200,000 byte payload in the middle, split into two ranges. The empty inline snapshot makes insta fail and print the recorded requests.
#[tokio::test]
async fn query_json_file_with_long_record_across_byte_ranges() {
let test = Test::new().with_single_file_json_long_record().await;
test.query("SET datafusion.optimizer.repartition_file_min_size = 0")
.await;
test.query("SET datafusion.execution.target_partitions = 2")
.await;
assert_snapshot!(
test.query("select id, length(payload) from json_long_record_table")
.await,
@r""
);
}
// in `impl Test`
async fn with_single_file_json_long_record(self) -> Test {
let json_data = format!(
"{{\"id\":0,\"payload\":\"a\"}}\n\
{{\"id\":1,\"payload\":\"{}\"}}\n\
{{\"id\":2,\"payload\":\"c\"}}\n",
"b".repeat(200_000)
);
self.with_bytes("/json_long_record_table.json", json_data)
.await
.register_json("json_long_record_table", "/json_long_record_table.json")
.await
}
On main (cee15b7) the first range needs six overflow GETs, one after another (the last clamped at the end of the file), to finish that one record:
Total Requests: 9
- GET (opts) path=json_long_record_table.json head=true
- GET (opts) path=json_long_record_table.json range=0-116418
- GET (opts) path=json_long_record_table.json range=116418-132802
- GET (opts) path=json_long_record_table.json range=132802-149186
- GET (opts) path=json_long_record_table.json range=149186-165570
- GET (opts) path=json_long_record_table.json range=165570-181954
- GET (opts) path=json_long_record_table.json range=181954-198338
- GET (opts) path=json_long_record_table.json range=198338-200068
- GET (opts) path=json_long_record_table.json range=100033-200068
Scaling. One range ending just inside a record of length L, read through AlignedBoundaryStream over InMemory, counting every get_opts call for that range ("patch" is the doubling fix proposed below):
| Record |
GETs on main |
GETs with patch |
| 1,000 B |
1 |
1 |
| 16 KiB |
1 |
1 |
| 64 KiB |
4 |
3 |
| 256 KiB |
16 |
5 |
| 1 MiB |
64 |
7 |
| 4 MiB |
256 |
9 |
| 16 MiB |
1,024 |
11 |
At 30 ms per request the 1 MiB row on main is 64 x 30 ms = 1.9 s of waiting for a single partition.
End to end. SELECT count(*), sum(id), sum(length(payload)) FROM t over one ~64 MiB NDJSON file per record size L. Rows are {"id":N,"payload":"AAAA..."} with payload length uniform in [L/2, 3L/2) (fixed seed), so boundaries land anywhere inside records. Explicit schema, target_partitions = 8, default range splitting. The store is LocalFileSystem with no latency, a fixed 30 ms sleep before every get_opts, or the benchmarks' S3-like LatencyObjectStore (benchmarks/src/util/latency_object_store.rs, P50 about 30 ms, P99 about 200 ms). "Requests" counts every get_opts call per query. release-nonlto build, 1 warm-up then 5 measured runs, i7-11700F (8 cores, 16 threads), 32 GB RAM, Linux 6.1, rustc 1.98.1. Times are the median with the min and max of the 5 runs in brackets, all in ms.
| Latency |
L |
Requests main / patch |
main |
patch |
Median speedup |
| none |
4 KiB |
9 / 9 |
21.6 [18.2, 25.8] |
18.2 [16.5, 20.9] |
ranges overlap |
| none |
256 KiB |
56 / 26 |
18.9 [14.4, 23.1] |
17.9 [16.8, 24.6] |
ranges overlap |
| none |
1 MiB |
313 / 40 |
16.8 [14.6, 20.0] |
16.4 [15.3, 22.3] |
ranges overlap |
| 30 ms fixed |
4 KiB |
9 / 9 |
88.9 [80.3, 91.5] |
87.1 [82.4, 88.8] |
ranges overlap |
| 30 ms fixed |
256 KiB |
56 / 26 |
644.6 [642.0, 650.9] |
206.9 [202.8, 222.8] |
3.1x |
| 30 ms fixed |
1 MiB |
313 / 40 |
1,998.6 [1,996.4, 2,009.5] |
248.2 [245.2, 253.5] |
8.1x |
| S3-like |
4 KiB |
9 / 9 |
322.1 [150.0, 367.6] |
310.0 [145.6, 358.2] |
ranges overlap |
| S3-like |
256 KiB |
56 / 26 |
1,463.6 [1,361.6, 1,609.9] |
500.3 [445.9, 609.9] |
2.9x |
| S3-like |
1 MiB |
313 / 40 |
4,957.4 [4,671.0, 5,535.7] |
744.5 [673.9, 750.8] |
6.7x |
The large slowdown appears only when per-request latency is simulated. Without it, all three sizes finish in 16.8 to 21.6 ms (medians) on main. At 30 ms per request the 1 MiB case takes about 2.0 s, consistent with roughly 66 sequential requests on the slowest partition, for about the same data volume as the 4 KiB case. Query results matched between main and the patch in every configuration.
Describe the solution you'd like
Double each overflow GET (32 KiB, 64 KiB, and so on) so a record needing N bytes past the initial window takes O(log N) sequential requests instead of O(N / 16 KiB). It's one new field, overflow_len, doubled before each overflow GET. The initial request and start alignment don't change, and records up to 16 KiB issue exactly the same requests as today. I have a patch, and with it the reproducer drops from 9 requests to 5, the six overflow GETs becoming two (the second clamped to the file size):
- GET (opts) path=json_long_record_table.json range=116418-149186
- GET (opts) path=json_long_record_table.json range=149186-200068
Cost. If N bytes are needed past the window, the overflow GETs request less than 2N + 32 KiB in total. On single ranges with plenty of data after the record, requested over needed bytes went from about 1.00 on main to 1.75 (64 KiB record), rising to 2.00 (4 MiB and 16 MiB). A record just over 16 KiB gets one 32 KiB overflow GET instead of 16 KiB. Across the whole 64 MiB query, requested bytes went up 0.74% (256 KiB records) and 1.57% (1 MiB records). These are requested bytes. The stream stops reading at the newline, so a streaming backend may transfer less, but I haven't measured transferred bytes on a real remote store.
Correctness. Only the end of each overflow range changes. The existing boundary_stream tests, each run across many chunk sizes, all pass. I also ran a randomized differential test against a verbatim copy of the current implementation: records up to 300,000 bytes, with and without a trailing newline, random range splits including an end past the file size, and chunk sizes from 7 bytes up to a single chunk. All 1,359 ranges were byte-identical, and each file's ranges concatenated back to the original. The test catches a deliberately planted off-by-one.
Describe alternatives you've considered
- Raise
END_SCAN_LOOKAHEAD. Every range overfetches more, even when its boundary record is short, and the request count is still linear in record length.
- One open-ended
GetRange::Offset(pos) for the overflow. Always a single extra request, but unbounded, so stores or wrappers that buffer or prefetch a whole requested range could read to the end of the file. It also departs from the bounded requests used everywhere else here.
- Doubling with a cap, say 8 MiB per request. Bounds any single request, but growth is linear again past the cap. Easy to add if reviewers prefer a hard bound.
Additional context
I'll open a PR with the patch and the reproducer above as a regression test.
Is your feature request related to a problem or challenge?
When a CSV or NDJSON file is split into byte ranges (
datafusion.optimizer.repartition_file_scans, on by default), each range is read throughAlignedBoundaryStream. The first request covers[start - 1, end + 16 KiB). For the first range it starts at byte 0, and it never goes past the end of the file. If the record crossingenddoesn't finish inside that window,ScanningLastTerminatorfetches the rest in 16 KiB GETs (END_SCAN_LOOKAHEAD, the last one clamped to the file size), each issued only after the previous one is exhausted:datafusion/datafusion/datasource/src/boundary_stream.rs
Lines 361 to 381 in cee15b7
Say
Ris the number of bytes from the range end up to and including that record's newline, and the file doesn't end before it. Then the range needsceil((R - 16 KiB) / 16 KiB)extra sequential GETs whenR > 16 KiB, and none otherwise. A range ending just after the start of a 1 MiB record makes 64 round trips (the initial GET plus 63 overflow GETs) before that partition can finish. Records of 16 KiB or less never trigger this.This is the retry case #8723 mentions (no newline found because the overfetched range was too small). #20823 (JSON) and #22962 (CSV) delivered #8723's main goal, one GET per range instead of three, whenever the initial lookahead reaches the newline. This issue is about the remaining overflow case in the current
AlignedBoundaryStream. I couldn't find an open issue or PR for it.It hits uncompressed NDJSON and newline-delimited CSV scans that the optimizer splits into byte ranges, when records are over 16 KiB (large JSON documents, long text fields) and the store has real per-request latency. Splitting only happens when a file group's total size reaches
datafusion.optimizer.repartition_file_min_size(default 1 MiB). Compressed files are never split, CSV withnewlines_in_valuesenabled isn't split, and JSON array files don't support ranged scans. On local disk the extra requests cost almost nothing.Reproducer. Add this to
datafusion/core/tests/datasource/object_store_access.rs, which already has theRequestCountingObjectStoreharness. The file has one 200,000 byte payload in the middle, split into two ranges. The empty inline snapshot makes insta fail and print the recorded requests.On
main(cee15b7) the first range needs six overflow GETs, one after another (the last clamped at the end of the file), to finish that one record:Scaling. One range ending just inside a record of length L, read through
AlignedBoundaryStreamoverInMemory, counting everyget_optscall for that range ("patch" is the doubling fix proposed below):mainAt 30 ms per request the 1 MiB row on
mainis 64 x 30 ms = 1.9 s of waiting for a single partition.End to end.
SELECT count(*), sum(id), sum(length(payload)) FROM tover one ~64 MiB NDJSON file per record size L. Rows are{"id":N,"payload":"AAAA..."}with payload length uniform in[L/2, 3L/2)(fixed seed), so boundaries land anywhere inside records. Explicit schema,target_partitions = 8, default range splitting. The store isLocalFileSystemwith no latency, a fixed 30 ms sleep before everyget_opts, or the benchmarks' S3-likeLatencyObjectStore(benchmarks/src/util/latency_object_store.rs, P50 about 30 ms, P99 about 200 ms). "Requests" counts everyget_optscall per query.release-nonltobuild, 1 warm-up then 5 measured runs, i7-11700F (8 cores, 16 threads), 32 GB RAM, Linux 6.1, rustc 1.98.1. Times are the median with the min and max of the 5 runs in brackets, all in ms.main/ patchmainThe large slowdown appears only when per-request latency is simulated. Without it, all three sizes finish in 16.8 to 21.6 ms (medians) on
main. At 30 ms per request the 1 MiB case takes about 2.0 s, consistent with roughly 66 sequential requests on the slowest partition, for about the same data volume as the 4 KiB case. Query results matched betweenmainand the patch in every configuration.Describe the solution you'd like
Double each overflow GET (32 KiB, 64 KiB, and so on) so a record needing N bytes past the initial window takes
O(log N)sequential requests instead ofO(N / 16 KiB). It's one new field,overflow_len, doubled before each overflow GET. The initial request and start alignment don't change, and records up to 16 KiB issue exactly the same requests as today. I have a patch, and with it the reproducer drops from 9 requests to 5, the six overflow GETs becoming two (the second clamped to the file size):Cost. If N bytes are needed past the window, the overflow GETs request less than
2N + 32 KiBin total. On single ranges with plenty of data after the record, requested over needed bytes went from about 1.00 onmainto 1.75 (64 KiB record), rising to 2.00 (4 MiB and 16 MiB). A record just over 16 KiB gets one 32 KiB overflow GET instead of 16 KiB. Across the whole 64 MiB query, requested bytes went up 0.74% (256 KiB records) and 1.57% (1 MiB records). These are requested bytes. The stream stops reading at the newline, so a streaming backend may transfer less, but I haven't measured transferred bytes on a real remote store.Correctness. Only the end of each overflow range changes. The existing
boundary_streamtests, each run across many chunk sizes, all pass. I also ran a randomized differential test against a verbatim copy of the current implementation: records up to 300,000 bytes, with and without a trailing newline, random range splits including an end past the file size, and chunk sizes from 7 bytes up to a single chunk. All 1,359 ranges were byte-identical, and each file's ranges concatenated back to the original. The test catches a deliberately planted off-by-one.Describe alternatives you've considered
END_SCAN_LOOKAHEAD. Every range overfetches more, even when its boundary record is short, and the request count is still linear in record length.GetRange::Offset(pos)for the overflow. Always a single extra request, but unbounded, so stores or wrappers that buffer or prefetch a whole requested range could read to the end of the file. It also departs from the bounded requests used everywhere else here.Additional context
--simulate-latency).AlignedBoundaryStream(datafusion/datasource/src/boundary_stream.rs), so both have the same request pattern. The end to end numbers were measured with NDJSON.I'll open a PR with the patch and the reproducer above as a regression test.