Consumes a stream of [ArrowLeafColumn] via a channel and serializes them using an [ArrowColumnWriter] Once the channel is exhausted, returns the ArrowColumnWriter.
(
mut rx: Receiver<ArrowLeafColumn>,
mut writer: ArrowColumnWriter,
reservation: MemoryReservation,
encoding_time: Time,
)
| 416 | /// Consumes a stream of [ArrowLeafColumn] via a channel and serializes them using an [ArrowColumnWriter] |
| 417 | /// Once the channel is exhausted, returns the ArrowColumnWriter. |
| 418 | async fn column_serializer_task( |
| 419 | mut rx: Receiver<ArrowLeafColumn>, |
| 420 | mut writer: ArrowColumnWriter, |
| 421 | reservation: MemoryReservation, |
| 422 | encoding_time: Time, |
| 423 | ) -> Result<(ArrowColumnWriter, MemoryReservation)> { |
| 424 | while let Some(col) = rx.recv().await { |
| 425 | let _timer = encoding_time.timer(); |
| 426 | writer.write(&col)?; |
| 427 | reservation.try_resize(writer.memory_size())?; |
| 428 | } |
| 429 | Ok((writer, reservation)) |
| 430 | } |
| 431 | |
| 432 | type ColumnWriterTask = SpawnedTask<Result<(ArrowColumnWriter, MemoryReservation)>>; |
| 433 | type ColSender = Sender<ArrowLeafColumn>; |
no test coverage detected
searching dependent graphs…