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]

Reply via email to