swanandx commented on code in PR #3277:
URL: https://github.com/apache/iceberg-rust/pull/3277#discussion_r4223479405


##########
crates/iceberg/src/arrow/reader/pipeline.rs:
##########
@@ -91,14 +91,23 @@ impl ArrowReader {
                     .try_flatten(),
             )
         } else {
+            // Open each file inside the stream that reads it. A file opened
+            // ahead of the flatten can sit unread, and on one HTTP/2
+            // connection an unread response stops the server sending for
+            // every other read on that connection.
             Box::pin(
                 tasks
-                    .map_ok(move |task| task_reader.clone().process(task))
+                    .map_ok(move |task| {
+                        task_reader
+                            .clone()
+                            .process(task)
+                            .try_flatten_stream()
+                            .boxed()
+                    })
                     .map_err(|err| {
                         Error::new(ErrorKind::Unexpected, "file scan task 
generate failed")
                             .with_source(err)
                     })
-                    .try_buffer_unordered(concurrency_limit_data_files)
                     .try_flatten_unordered(concurrency_limit_data_files),
             )

Review Comment:
   I measured this on GCS and the first version of the PR was clearly slower, 
so I have replaced it.
   
   Scan of 64 files of 25 KB, concurrency 8, metadata size hint 8 KiB, and for 
the deletes case one positional delete file per data file. I ran the variants 
turn by turn in the same process, 12 rounds each, and compared them round by 
round, because a plain before/after run on a home connection is all noise.
   
   |  | no deletes | with deletes |
   | --- | --- | --- |
   | main | 8489 ms | 13344 ms |
   | first version of this PR | +52.8% | +37.8% |
   | new version | +0.7% | +1.1% |
   
   The first version was slower in 12 rounds out of 12 (95% CI [+3348, +5470] 
ms without deletes). The new version is within [-241, +243] ms of main, so I 
cannot measure a difference. You were right that removing the prefetch was the 
wrong trade.
   
   The prefetch is not what causes the stall. The cause is a prefetched open 
that nobody polls any more: its response sits half read and keeps part of the 
HTTP/2 connection window, and enough of those stop the server sending on that 
connection at all. That is also why `load_deletes` never triggered it, it 
spawns, so it finishes whether or not anyone polls it.
   
   So the new version leaves the combinators exactly as they are on main and 
only spawns the open, the same way `caching_delete_file_loader.rs` already does:
   
   ```rust
   let opening = runtime.io().spawn(task_reader.clone().process(task));
   ```
   
   The stall test now passes from concurrency 2 to 24, where main stalls from 3 
to 10. It is a `#[tokio::test]`, so this holds on a current-thread runtime too.
   
   Two things to be aware of: a panic inside `process` now comes back as an 
error through `JoinError` instead of unwinding, and a scan stream that is 
dropped early leaves up to `concurrency_limit_data_files` opens running with 
their results thrown away, which is what `load_deletes` already does.



-- 
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