Abhisheklearn12 opened a new issue, #25830: URL: https://github.com/apache/datafusion/issues/25830
### 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: https://github.com/apache/datafusion/blob/cee15b7cca30faaf894d625d245397e4ec4320f9/datafusion/datasource/src/boundary_stream.rs#L361-L381 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. ```rust #[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` (cee15b7cc) 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 1. **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. 2. **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. 3. **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 - Parent issue #8723. Related merged work: #20823 (JSON) and #22962 (CSV, closed #21419). Related tooling in open PRs: #23803 (counts simulated GET requests per query) and #24093 (models remote object storage for `--simulate-latency`). - CSV and JSON share `AlignedBoundaryStream` (`datafusion/datasource/src/boundary_stream.rs`), so both have the same request pattern. The end to end numbers were measured with NDJSON. - Overflow GETs only come from finishing the record at a range's end. Start alignment, which every range except the first does, never issues extra requests: if its initial window has no newline, no record starts inside the range, so it yields nothing. The first range starts at byte 0 and can still need overflow GETs to finish a long first record. - The latency is simulated, not real S3. The request counts and the linear vs logarithmic scaling don't depend on the store. I'll open a PR with the patch and the reproducer above as a regression test. -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
