| 210 | return batch_reader_->schema(); |
| 211 | } |
| 212 | arrow::Result<FlightStreamChunk> Next() override { |
| 213 | FlightStreamChunk out; |
| 214 | internal::FlightData* data; |
| 215 | peekable_reader_->Peek(&data); |
| 216 | if (!data) { |
| 217 | out.app_metadata = nullptr; |
| 218 | out.data = nullptr; |
| 219 | RETURN_NOT_OK(stream_->Finish(Status::OK())); |
| 220 | return out; |
| 221 | } |
| 222 | |
| 223 | if (!data->metadata) { |
| 224 | // Metadata-only (data->metadata is the IPC header) |
| 225 | out.app_metadata = data->app_metadata; |
| 226 | out.data = nullptr; |
| 227 | peekable_reader_->Next(&data); |
| 228 | return out; |
| 229 | } |
| 230 | |
| 231 | if (!batch_reader_) { |
| 232 | RETURN_NOT_OK(EnsureDataStarted()); |
| 233 | // Re-peek here since EnsureDataStarted() advances the stream |
| 234 | return Next(); |
| 235 | } |
| 236 | auto status = batch_reader_->ReadNext(&out.data); |
| 237 | if (ARROW_PREDICT_FALSE(!status.ok())) { |
| 238 | return stream_->Finish(std::move(status)); |
| 239 | } |
| 240 | out.app_metadata = std::move(app_metadata_); |
| 241 | return out; |
| 242 | } |
| 243 | arrow::Result<std::vector<std::shared_ptr<RecordBatch>>> ToRecordBatches() override { |
| 244 | return ToRecordBatches(stop_token_); |
| 245 | } |