| 379 | // AsyncTaskGroup |
| 380 | template <typename T> |
| 381 | bool AsyncTaskScheduler::AddAsyncGenerator(std::function<Future<T>()> generator, |
| 382 | std::function<Status(const T&)> visitor, |
| 383 | std::string_view name) { |
| 384 | struct State { |
| 385 | State(std::function<Future<T>()> generator, std::function<Status(const T&)> visitor, |
| 386 | std::unique_ptr<AsyncTaskGroup> task_group, std::string_view name) |
| 387 | : generator(std::move(generator)), |
| 388 | visitor(std::move(visitor)), |
| 389 | task_group(std::move(task_group)), |
| 390 | name(name) {} |
| 391 | std::function<Future<T>()> generator; |
| 392 | std::function<Status(const T&)> visitor; |
| 393 | std::unique_ptr<AsyncTaskGroup> task_group; |
| 394 | std::string_view name; |
| 395 | }; |
| 396 | struct SubmitTask : public Task { |
| 397 | explicit SubmitTask(std::unique_ptr<State> state_holder) |
| 398 | : state_holder(std::move(state_holder)) {} |
| 399 | |
| 400 | struct SubmitTaskCallback { |
| 401 | SubmitTaskCallback(std::unique_ptr<State> state_holder, Future<> task_completion) |
| 402 | : state_holder(std::move(state_holder)), |
| 403 | task_completion(std::move(task_completion)) {} |
| 404 | void operator()(const Result<T>& maybe_item) { |
| 405 | if (!maybe_item.ok()) { |
| 406 | task_completion.MarkFinished(maybe_item.status()); |
| 407 | return; |
| 408 | } |
| 409 | const auto& item = *maybe_item; |
| 410 | if (IsIterationEnd(item)) { |
| 411 | task_completion.MarkFinished(); |
| 412 | return; |
| 413 | } |
| 414 | Status visit_st = state_holder->visitor(item); |
| 415 | if (!visit_st.ok()) { |
| 416 | task_completion.MarkFinished(std::move(visit_st)); |
| 417 | return; |
| 418 | } |
| 419 | state_holder->task_group->AddTask( |
| 420 | std::make_unique<SubmitTask>(std::move(state_holder))); |
| 421 | task_completion.MarkFinished(); |
| 422 | } |
| 423 | std::unique_ptr<State> state_holder; |
| 424 | Future<> task_completion; |
| 425 | }; |
| 426 | |
| 427 | Result<Future<>> operator()() { |
| 428 | Future<> task = Future<>::Make(); |
| 429 | // Consume as many items as we can (those that are already finished) |
| 430 | // synchronously to avoid recursion / stack overflow. |
| 431 | while (true) { |
| 432 | Future<T> next = state_holder->generator(); |
| 433 | if (next.TryAddCallback( |
| 434 | [&] { return SubmitTaskCallback(std::move(state_holder), task); })) { |
| 435 | return task; |
| 436 | } |
| 437 | ARROW_ASSIGN_OR_RAISE(T item, next.result()); |
| 438 | if (IsIterationEnd(item)) { |