MCPcopy Create free account
hub / github.com/apache/arrow / AddAsyncGenerator

Method AddAsyncGenerator

cpp/src/arrow/util/async_util.h:381–457  ·  view source on GitHub ↗

Source from the content-addressed store, hash-verified

379// AsyncTaskGroup
380template <typename T>
381bool 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)) {

Callers 1

TESTFunction · 0.80

Calls 4

MakeFunction · 0.70
OKFunction · 0.50
getMethod · 0.45
AddTaskMethod · 0.45

Tested by 1

TESTFunction · 0.64