Fetches another batch from the task queue. Starts the task if not started yet.
| 323 | task_->start(maxDrivers_, numConcurrentSplitGroups_); |
| 324 | queue_->setNumProducers(numSplitGroups_ * task_->numOutputDrivers()); |
| 325 | } catch (const BoltException& e) { |
| 326 | // Could not find output pipeline, due to Task terminated before |
| 327 | // start. Do not override the error. |
| 328 | if (e.message().find("Output pipeline not found for task") == |
| 329 | std::string::npos) { |
| 330 | throw; |
| 331 | } |
| 332 | } |
| 333 | } |
| 334 | } |
| 335 | |
| 336 | /// Fetches another batch from the task queue. |
| 337 | /// Starts the task if not started yet. |
| 338 | bool moveNext() override { |
| 339 | start(); |
| 340 | if (error_) { |
| 341 | std::rethrow_exception(error_); |
| 342 | } |
| 343 |
no test coverage detected