swanandx opened a new issue, #3276:
URL: https://github.com/apache/iceberg-rust/issues/3276

   ### Apache Iceberg Rust version
   
   main, 0.10.1
   
   ### Describe the bug
   
   A scan stops making progress and returns no error. We saw it stall against
   GCS.
   
   In `crates/iceberg/src/arrow/reader/pipeline.rs`:
   
   ```rust
   tasks
       .map_ok(move |task| task_reader.clone().process(task))
       // .map_err(..) elided
       .try_buffer_unordered(concurrency_limit_data_files)
       .try_flatten_unordered(concurrency_limit_data_files)
   ```
   
   There are two stages. `try_buffer_unordered` opens files, which means reading
   each file's Parquet footer. `try_flatten_unordered` then reads the rows out.
   
   The second stage handles at most `concurrency_limit_data_files` files at 
once.
   When it is full it stops taking new ones, so the opens behind it still 
finish:
   the footer arrives, and nobody ever reads it.
   
   That is what breaks the scan. All these requests share one HTTP/2 connection,
   and on such a connection the server may only send so far ahead of what the
   client has read. HTTP/2 calls that allowance the connection window. A footer
   nobody reads counts against the window and never gives it back. Once the 
unread
   footers use up the window, the server cannot send anything more, including 
for
   the files the scan is still reading. Everything stops, with no error.
   
   Filling the window takes a lot of files open at once, which is why it shows
   on big machines: `concurrency_limit_data_files` defaults to
   `available_parallelism()`.
   
   #806 introduced this, replacing a `spawn` and a bounded channel with these
   combinators. In the old shape every file opened was being read.
   
   ### To Reproduce
   
   It reproduces with no network, deterministically. Model one HTTP/2 connection
   as a `Storage`: serve frames to the open streams round-robin, and size the
   window to hold two footers. Scan twelve files, each larger than the window.
   
   With `concurrency_limit_data_files` from 3 to 10 the scan stalls. It passes 
at
   2, where the window still covers the unread footers, and at 11 and above, 
where
   fewer than two files are left over to go unread. I have this as a test and 
will
   attach it to a fix.
   
   Two conditions have to hold, and both are defaults on a large machine 
reading a
   table with more data files than cores:
   
   - more data files than `concurrency_limit_data_files`, so the second stage 
fills
   - enough unread footers to use up the window
   
   We saw it against GCS on a 64 core machine, consistently. It did not 
reproduce
   on 4 cores. We did not try the sizes in between, so I cannot give a threshold
   in cores.
   
   ### Expected behavior
   
   The scan finishes.
   
   ### Willingness to contribute
   
   I can contribute a fix for this bug independently


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