| 483 | } |
| 484 | |
| 485 | void TableScan::checkPreload() { |
| 486 | auto executor = connector_->executor(); |
| 487 | if (maxSplitPreloadPerDriver_ == 0 || !executor || |
| 488 | !connector_->supportsSplitPreload() || !asyncThreadCtx_->allowPreload()) { |
| 489 | return; |
| 490 | } |
| 491 | if (dataSource_->allPrefetchIssued()) { |
| 492 | maxPreloadedSplits_ = driverCtx_->task->numDrivers(driverCtx_->driver) * |
| 493 | maxSplitPreloadPerDriver_; |
| 494 | if (!splitPreloader_) { |
| 495 | splitPreloader_ = |
| 496 | [executor, this](std::shared_ptr<connector::ConnectorSplit> split) { |
| 497 | preload(split); |
| 498 | |
| 499 | int64_t preloadBytes = split->splitSizeBytes(); |
| 500 | executor->add([connectorSplit = split, |
| 501 | ctx = asyncThreadCtx_, |
| 502 | preloadBytes]() mutable { |
| 503 | connector::AsyncThreadCtx::Guard guard(ctx.get(), preloadBytes); |
| 504 | if (!guard) { |
| 505 | return; |
| 506 | } |
| 507 | connectorSplit->dataSource->prepare(); |
| 508 | connectorSplit.reset(); |
| 509 | }); |
| 510 | }; |
| 511 | } |
| 512 | } |
| 513 | } |
| 514 | |
| 515 | bool TableScan::isFinished() { |
| 516 | return noMoreSplits_; |
nothing calls this directly
no test coverage detected