| 676 | void TableBatchReader::set_chunksize(int64_t chunksize) { max_chunksize_ = chunksize; } |
| 677 | |
| 678 | Status TableBatchReader::ReadNext(std::shared_ptr<RecordBatch>* out) { |
| 679 | if (absolute_row_position_ == table_.num_rows()) { |
| 680 | *out = nullptr; |
| 681 | return Status::OK(); |
| 682 | } |
| 683 | |
| 684 | // Determine the minimum contiguous slice across all columns |
| 685 | int64_t chunksize = |
| 686 | std::min(table_.num_rows() - absolute_row_position_, max_chunksize_); |
| 687 | std::vector<const Array*> chunks(table_.num_columns()); |
| 688 | for (int i = 0; i < table_.num_columns(); ++i) { |
| 689 | auto chunk = column_data_[i]->chunk(chunk_numbers_[i]).get(); |
| 690 | int64_t chunk_remaining = chunk->length() - chunk_offsets_[i]; |
| 691 | |
| 692 | if (chunk_remaining < chunksize) { |
| 693 | chunksize = chunk_remaining; |
| 694 | } |
| 695 | |
| 696 | chunks[i] = chunk; |
| 697 | } |
| 698 | |
| 699 | // Slice chunks and advance chunk index as appropriate |
| 700 | std::vector<std::shared_ptr<ArrayData>> batch_data(table_.num_columns()); |
| 701 | |
| 702 | for (int i = 0; i < table_.num_columns(); ++i) { |
| 703 | // Exhausted chunk |
| 704 | const Array* chunk = chunks[i]; |
| 705 | const int64_t offset = chunk_offsets_[i]; |
| 706 | std::shared_ptr<ArrayData> slice_data; |
| 707 | if ((chunk->length() - offset) == chunksize) { |
| 708 | ++chunk_numbers_[i]; |
| 709 | chunk_offsets_[i] = 0; |
| 710 | if (offset > 0) { |
| 711 | // Need to slice |
| 712 | slice_data = chunk->Slice(offset, chunksize)->data(); |
| 713 | } else { |
| 714 | // No slice |
| 715 | slice_data = chunk->data(); |
| 716 | } |
| 717 | } else { |
| 718 | chunk_offsets_[i] += chunksize; |
| 719 | slice_data = chunk->Slice(offset, chunksize)->data(); |
| 720 | } |
| 721 | batch_data[i] = std::move(slice_data); |
| 722 | } |
| 723 | |
| 724 | absolute_row_position_ += chunksize; |
| 725 | *out = RecordBatch::Make(table_.schema(), chunksize, std::move(batch_data)); |
| 726 | |
| 727 | return Status::OK(); |
| 728 | } |
| 729 | |
| 730 | } // namespace arrow |