MCPcopy Create free account
hub / github.com/apache/arrow-rs / read_next

Method read_next

arrow-avro/src/reader/async_reader/mod.rs:254–534  ·  view source on GitHub ↗
(&mut self, cx: &mut Context<'_>)

Source from the content-addressed store, hash-verified

252
253 #[forbid(clippy::question_mark_used)]
254 fn read_next(&mut self, cx: &mut Context<'_>) -> Poll<Option<Result<RecordBatch, AvroError>>> {
255 loop {
256 match mem::replace(&mut self.reader_state, ReaderState::InvalidState) {
257 ReaderState::Idle { reader } => {
258 let range = self.range.clone();
259 if range.start >= range.end {
260 return self.finish_with_error(AvroError::InvalidArgument(format!(
261 "Invalid range specified for Avro file: start {} >= end {}, file_size: {}",
262 range.start, range.end, self.file_size
263 )));
264 }
265
266 let future = Self::fetch_bytes(reader, range).boxed();
267 self.reader_state = ReaderState::FetchingData {
268 future,
269 next_behaviour: FetchNextBehaviour::ReadSyncMarker,
270 };
271 }
272 ReaderState::FetchingData {
273 mut future,
274 next_behaviour,
275 } => {
276 let (reader, data_chunk) = match future.poll_unpin(cx) {
277 Poll::Ready(Ok(data)) => data,
278 Poll::Ready(Err(e)) => return self.finish_with_error(e),
279 Poll::Pending => {
280 self.reader_state = ReaderState::FetchingData {
281 future,
282 next_behaviour,
283 };
284 return Poll::Pending;
285 }
286 };
287
288 match next_behaviour {
289 FetchNextBehaviour::ReadSyncMarker => {
290 let sync_marker_pos = data_chunk
291 .windows(16)
292 .position(|slice| slice == self.sync_marker);
293 let block_start = match sync_marker_pos {
294 Some(pos) => pos + 16, // Move past the sync marker
295 None => {
296 // Sync marker not found, valid if we arbitrarily split the file at its end.
297 self.reader_state = ReaderState::Finished;
298 return Poll::Ready(None);
299 }
300 };
301
302 self.reader_state = ReaderState::DecodingBlock {
303 reader,
304 data: data_chunk.slice(block_start..),
305 };
306 }
307 FetchNextBehaviour::DecodeVLQHeader => {
308 let mut data = data_chunk;
309
310 // Feed bytes one at a time until we reach Data state (VLQ header complete)
311 while !matches!(self.block_decoder.state(), BlockDecoderState::Data) {

Callers 1

poll_nextMethod · 0.45

Calls 15

finish_with_errorMethod · 0.80
positionMethod · 0.80
bytes_remainingMethod · 0.80
remaining_block_rangeMethod · 0.80
start_flushingMethod · 0.80
decode_blockMethod · 0.80
batch_is_fullMethod · 0.80
flush_blockMethod · 0.80
poll_flushMethod · 0.80
cloneMethod · 0.45
sliceMethod · 0.45
is_emptyMethod · 0.45

Tested by

no test coverage detected