| 1543 | } |
| 1544 | |
| 1545 | void Task::noMoreSplitsForGroup( |
| 1546 | const core::PlanNodeId& planNodeId, |
| 1547 | int32_t splitGroupId) { |
| 1548 | std::vector<ContinuePromise> promises; |
| 1549 | EventCompletionNotifier stateChangeNotifier; |
| 1550 | { |
| 1551 | std::lock_guard<std::timed_mutex> l(mutex_); |
| 1552 | |
| 1553 | auto& splitsState = getPlanNodeSplitsStateLocked(planNodeId); |
| 1554 | auto& splitsStore = splitsState.groupSplitsStores[splitGroupId]; |
| 1555 | splitsStore.noMoreSplits = true; |
| 1556 | promises = std::move(splitsStore.splitPromises); |
| 1557 | |
| 1558 | // There were no splits in this group, hence, no active drivers. Mark the |
| 1559 | // group complete. |
| 1560 | if (seenSplitGroups_.count(splitGroupId) == 0) { |
| 1561 | taskStats_.completedSplitGroups.insert(splitGroupId); |
| 1562 | stateChangeNotifier.activate(std::move(stateChangePromises_)); |
| 1563 | } |
| 1564 | } |
| 1565 | stateChangeNotifier.notify(); |
| 1566 | for (auto& promise : promises) { |
| 1567 | promise.setValue(); |
| 1568 | } |
| 1569 | } |
| 1570 | |
| 1571 | void Task::noMoreSplits(const core::PlanNodeId& planNodeId) { |
| 1572 | std::vector<ContinuePromise> splitPromises; |