| 94 | } |
| 95 | |
| 96 | arrow::Result<FlightPayload> Next() override { |
| 97 | FlightPayload payload; |
| 98 | if (records_sent_ >= total_records_) { |
| 99 | // Signal that iteration is over |
| 100 | payload.ipc_message.metadata = nullptr; |
| 101 | return payload; |
| 102 | } |
| 103 | |
| 104 | if (verify_) { |
| 105 | // mutate first array |
| 106 | auto data = |
| 107 | reinterpret_cast<int64_t*>(arrays_[0]->data()->buffers[1]->mutable_data()); |
| 108 | for (int64_t i = 0; i < batch_length_; ++i) { |
| 109 | data[i] = start_ + records_sent_ + i; |
| 110 | } |
| 111 | } |
| 112 | |
| 113 | auto batch = batch_; |
| 114 | |
| 115 | // Last partial batch |
| 116 | if (records_sent_ + batch_length_ > total_records_) { |
| 117 | batch = batch_->Slice(0, total_records_ - records_sent_); |
| 118 | records_sent_ += total_records_ - records_sent_; |
| 119 | } else { |
| 120 | records_sent_ += batch_length_; |
| 121 | } |
| 122 | RETURN_NOT_OK(ipc::GetRecordBatchPayload(*batch, ipc_options_, &payload.ipc_message)); |
| 123 | return payload; |
| 124 | } |
| 125 | |
| 126 | private: |
| 127 | const int64_t start_; |
no test coverage detected