| 293 | }; |
| 294 | |
| 295 | Result<EnumeratedRecordBatchGenerator> FragmentToBatches( |
| 296 | const Enumerated<std::shared_ptr<Fragment>>& fragment, |
| 297 | const std::shared_ptr<ScanOptions>& options) { |
| 298 | #ifdef ARROW_WITH_OPENTELEMETRY |
| 299 | util::tracing::Span span; |
| 300 | START_SPAN(span, "Scanner::FragmentToBatches", |
| 301 | { |
| 302 | {"arrow.dataset.fragment", fragment.value->ToString()}, |
| 303 | {"arrow.dataset.fragment.index", fragment.index}, |
| 304 | {"arrow.dataset.fragment.last", fragment.last}, |
| 305 | {"arrow.dataset.fragment.type_name", fragment.value->type_name()}, |
| 306 | }); |
| 307 | #endif |
| 308 | ARROW_ASSIGN_OR_RAISE(auto batch_gen, fragment.value->ScanBatchesAsync(options)); |
| 309 | ArrayVector columns; |
| 310 | for (const auto& field : options->dataset_schema->fields()) { |
| 311 | // TODO(ARROW-7051): use helper to make empty batch |
| 312 | ARROW_ASSIGN_OR_RAISE(auto array, |
| 313 | MakeArrayOfNull(field->type(), /*length=*/0, options->pool)); |
| 314 | columns.push_back(std::move(array)); |
| 315 | } |
| 316 | WRAP_ASYNC_GENERATOR(batch_gen); |
| 317 | batch_gen = MakeDefaultIfEmptyGenerator( |
| 318 | std::move(batch_gen), |
| 319 | RecordBatch::Make(options->dataset_schema, /*num_rows=*/0, std::move(columns))); |
| 320 | auto enumerated_batch_gen = MakeEnumeratedGenerator(std::move(batch_gen)); |
| 321 | |
| 322 | auto combine_fn = [fragment, cache_metadata = options->cache_metadata]( |
| 323 | const Enumerated<std::shared_ptr<RecordBatch>>& record_batch) { |
| 324 | if (!cache_metadata && record_batch.last) { |
| 325 | ARROW_WARN_NOT_OK(fragment.value->ClearCachedMetadata(), |
| 326 | "Could not clear cached metadata on fragment"); |
| 327 | } |
| 328 | return EnumeratedRecordBatch{record_batch, fragment}; |
| 329 | }; |
| 330 | |
| 331 | return MakeMappedGenerator(enumerated_batch_gen, std::move(combine_fn)); |
| 332 | } |
| 333 | |
| 334 | Result<AsyncGenerator<EnumeratedRecordBatchGenerator>> FragmentsToBatches( |
| 335 | FragmentGenerator fragment_gen, const std::shared_ptr<ScanOptions>& options) { |
no test coverage detected