| 525 | |
| 526 | impl<R: BufRead> BufReader<R> { |
| 527 | fn read(&mut self) -> Result<Option<RecordBatch>, ArrowError> { |
| 528 | loop { |
| 529 | let buf = self.reader.fill_buf()?; |
| 530 | let decoded = self.decoder.decode(buf)?; |
| 531 | self.reader.consume(decoded); |
| 532 | // Yield if decoded no bytes or the decoder is full |
| 533 | // |
| 534 | // The capacity check avoids looping around and potentially |
| 535 | // blocking reading data in fill_buf that isn't needed |
| 536 | // to flush the next batch |
| 537 | if decoded == 0 || self.decoder.capacity() == 0 { |
| 538 | break; |
| 539 | } |
| 540 | } |
| 541 | |
| 542 | self.decoder.flush() |
| 543 | } |
| 544 | } |
| 545 | |
| 546 | impl<R: BufRead> Iterator for BufReader<R> { |