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

Method decode

arrow-ipc/src/reader/stream.rs:159–286  ·  view source on GitHub ↗

Try to read the next [`RecordBatch`] from the provided [`Buffer`] [`Buffer::advance`] will be called on `buffer` for any consumed bytes. The push-based interface facilitates integration with sources that yield arbitrarily delimited bytes ranges, such as a chunked byte stream received from object storage ``` # use arrow_array::RecordBatch; # use arrow_buffer::Buffer; # use arrow_ipc::reader::Str

(&mut self, buffer: &mut Buffer)

Source from the content-addressed store, hash-verified

157 /// }
158 /// ```
159 pub fn decode(&mut self, buffer: &mut Buffer) -> Result<Option<RecordBatch>, ArrowError> {
160 while !buffer.is_empty() {
161 match &mut self.state {
162 DecoderState::Header {
163 buf,
164 read,
165 continuation,
166 } => {
167 let offset_buf = &mut buf[*read as usize..];
168 let to_read = buffer.len().min(offset_buf.len());
169 offset_buf[..to_read].copy_from_slice(&buffer[..to_read]);
170 *read += to_read as u8;
171 buffer.advance(to_read);
172 if *read == 4 {
173 if !*continuation && buf == &CONTINUATION_MARKER {
174 *continuation = true;
175 *read = 0;
176 continue;
177 }
178 let size = u32::from_le_bytes(*buf);
179
180 if size == 0 {
181 self.state = DecoderState::Finished;
182 continue;
183 }
184 self.state = DecoderState::Message { size };
185 }
186 }
187 DecoderState::Message { size } => {
188 let len = *size as usize;
189 if self.buf.is_empty() && buffer.len() > len {
190 let message = MessageBuffer::try_new(buffer.slice_with_length(0, len))?;
191 self.state = DecoderState::Body { message };
192 buffer.advance(len);
193 continue;
194 }
195
196 let to_read = buffer.len().min(len - self.buf.len());
197 self.buf.extend_from_slice(&buffer[..to_read]);
198 buffer.advance(to_read);
199 if self.buf.len() == len {
200 let message = MessageBuffer::try_new(std::mem::take(&mut self.buf).into())?;
201 self.state = DecoderState::Body { message };
202 }
203 }
204 DecoderState::Body { message } => {
205 let message = message.as_ref();
206 let body_length = message.bodyLength() as usize;
207
208 let body = if self.buf.is_empty() && buffer.len() >= body_length {
209 let body = buffer.slice_with_length(0, body_length);
210 buffer.advance(body_length);
211 body
212 } else {
213 let to_read = buffer.len().min(body_length - self.buf.len());
214 self.buf.extend_from_slice(&buffer[..to_read]);
215 buffer.advance(to_read);
216

Callers 4

verify_arrow_streamFunction · 0.45
test_eosFunction · 0.45
test_schemaFunction · 0.45

Calls 15

try_newFunction · 0.85
takeFunction · 0.85
fb_to_schemaFunction · 0.85
defaultFunction · 0.85
read_dictionary_implFunction · 0.85
slice_with_lengthMethod · 0.80
extend_from_sliceMethod · 0.80
header_typeMethod · 0.80
header_as_schemaMethod · 0.80
read_record_batchMethod · 0.80

Tested by 4

verify_arrow_streamFunction · 0.36
test_eosFunction · 0.36
test_schemaFunction · 0.36