From 1a46192fe5a6bf84dd54356889b39fbb890cc109 Mon Sep 17 00:00:00 2001 From: winningsix Date: Thu, 30 Jul 2026 08:10:47 +0000 Subject: [PATCH 1/2] perf(exec): schedule preloaded splits in ready order Avoid binding scan drivers to an arbitrary in-flight preload. Wake all current split waiters after each preload completion so accumulated ready work cannot be stranded by a lost edge notification. --- velox/exec/TableScan.cpp | 13 +++-- velox/exec/Task.cpp | 30 ++++++++++- velox/exec/Task.h | 6 +++ velox/exec/TaskStructs.cpp | 16 ++++-- velox/exec/TaskStructs.h | 15 +++++- velox/exec/tests/TaskTest.cpp | 99 +++++++++++++++++++++++++++++++++++ 6 files changed, 167 insertions(+), 12 deletions(-) diff --git a/velox/exec/TableScan.cpp b/velox/exec/TableScan.cpp index a85674dd950..b17a7f977a1 100644 --- a/velox/exec/TableScan.cpp +++ b/velox/exec/TableScan.cpp @@ -508,10 +508,15 @@ void TableScan::checkPreload() { [ioExecutor, this](const std::shared_ptr& split) { preload(split); - ioExecutor->add([connectorSplit = split]() mutable { - connectorSplit->dataSource->prepare(); - connectorSplit.reset(); - }); + ioExecutor->add( + [connectorSplit = split, + task = operatorCtx_->task(), + splitGroupId = driverCtx_->splitGroupId, + planNodeId = planNodeId()]() mutable { + connectorSplit->dataSource->prepare(); + connectorSplit.reset(); + task->splitPreloadFinished(splitGroupId, planNodeId); + }); }; } } diff --git a/velox/exec/Task.cpp b/velox/exec/Task.cpp index f170b620658..b9b740b5acb 100644 --- a/velox/exec/Task.cpp +++ b/velox/exec/Task.cpp @@ -282,8 +282,11 @@ class QueueSplitsStore : public SplitsStore { Split& split, ContinueFuture& future) override { if (!splits_.empty()) { - split = getSplit(maxPreloadSplits, preload); - return true; + if (getSplit(maxPreloadSplits, preload, split)) { + return true; + } + future = makeFuture(); + return false; } if (tryGetBarrier(driverId, split)) { return true; @@ -2307,6 +2310,29 @@ BlockingReason Task::getSplitOrFuture( : BlockingReason::kWaitForSplit; } +void Task::splitPreloadFinished( + uint32_t splitGroupId, + const core::PlanNodeId& planNodeId) { + std::vector promises; + { + std::lock_guard l(mutex_); + const auto stateIt = splitsStates_.find(planNodeId); + if (stateIt == splitsStates_.end()) { + return; + } + const auto storeIt = + stateIt->second.groupSplitsStores.find(splitGroupId); + if (storeIt == stateIt->second.groupSplitsStores.end() || + storeIt->second == nullptr) { + return; + } + promises = storeIt->second->splitPreloadFinished(); + } + for (auto& promise : promises) { + promise.setValue(); + } +} + bool Task::testingHasDriverWaitForSplit() const { std::lock_guard l(mutex_); for (const auto& splitState : splitsStates_) { diff --git a/velox/exec/Task.h b/velox/exec/Task.h index b8d5ff63e14..f2e4184061d 100644 --- a/velox/exec/Task.h +++ b/velox/exec/Task.h @@ -527,6 +527,12 @@ class Task : public std::enable_shared_from_this { exec::Split& split, ContinueFuture& future); + /// Notifies scan drivers that an asynchronously preloaded split has + /// completed and ready-first split selection should be retried. + void splitPreloadFinished( + uint32_t splitGroupId, + const core::PlanNodeId& planNodeId); + /// Returns the scaled scan controller for a given table scan node if the /// query has configured. std::shared_ptr getScaledScanControllerLocked( diff --git a/velox/exec/TaskStructs.cpp b/velox/exec/TaskStructs.cpp index 2d5cbc27d42..e5a2b61363b 100644 --- a/velox/exec/TaskStructs.cpp +++ b/velox/exec/TaskStructs.cpp @@ -49,9 +49,10 @@ ContinueFuture SplitsStore::makeFuture() { return std::move(future); } -Split SplitsStore::getSplit( +bool SplitsStore::getSplit( int maxPreloadSplits, - const ConnectorSplitPreloadFunc& preload) { + const ConnectorSplitPreloadFunc& preload, + Split& split) { int readySplitIndex = -1; if (maxPreloadSplits > 0) { for (int i = 0, end = std::min(maxPreloadSplits, splits_.size()); @@ -72,12 +73,19 @@ Split SplitsStore::getSplit( preloadingSplits_->erase(connectorSplit); } } + // Do not bind a scan driver to an arbitrary in-flight preload. The + // completion path wakes a waiter, which retries and takes whichever split + // is ready first. This keeps I/O and compute pipelined without + // head-of-line blocking on the queue front. + if (readySplitIndex == -1) { + return false; + } } if (readySplitIndex == -1) { readySplitIndex = 0; } VELOX_CHECK(!splits_.empty()); - auto split = std::move(splits_[readySplitIndex]); + split = std::move(splits_[readySplitIndex]); splits_.erase(splits_.begin() + readySplitIndex); --taskStats_->numQueuedSplits; ++taskStats_->numRunningSplits; @@ -93,7 +101,7 @@ Split SplitsStore::getSplit( if (taskStats_->firstSplitStartTimeMs == 0) { taskStats_->firstSplitStartTimeMs = taskStats_->lastSplitStartTimeMs; } - return split; + return true; } bool SplitsStore::tryGetBarrier( diff --git a/velox/exec/TaskStructs.h b/velox/exec/TaskStructs.h index 59c8f0912b6..b0db8329296 100644 --- a/velox/exec/TaskStructs.h +++ b/velox/exec/TaskStructs.h @@ -136,6 +136,16 @@ class SplitsStore { return std::move(promises_); } + /// Wakes drivers waiting for a preloaded split to become ready. + /// + /// Readiness is level-triggered: by the time this completion arrives there + /// may already be multiple ready splits, and completions that happened + /// before waiters registered do not leave an edge to wake them later. + /// Wake all current waiters so each rechecks the ready queue. + std::vector splitPreloadFinished() { + return std::move(promises_); + } + void setTaskStats(TaskStats& taskStats) { taskStats_ = &taskStats; } @@ -147,9 +157,10 @@ class SplitsStore { } protected: - Split getSplit( + bool getSplit( int maxPreloadSplits, - const ConnectorSplitPreloadFunc& preload); + const ConnectorSplitPreloadFunc& preload, + Split& split); ContinueFuture makeFuture(); diff --git a/velox/exec/tests/TaskTest.cpp b/velox/exec/tests/TaskTest.cpp index 6cb6680e841..babf771b492 100644 --- a/velox/exec/tests/TaskTest.cpp +++ b/velox/exec/tests/TaskTest.cpp @@ -17,6 +17,7 @@ #include "velox/exec/Task.h" #include #include "folly/synchronization/EventCount.h" +#include #include "velox/common/base/tests/GTestUtils.h" #include "velox/common/file/tests/FaultyFileSystem.h" #include "velox/common/future/VeloxPromise.h" @@ -872,6 +873,104 @@ TEST_F(TaskTest, stateChangeFutureNotFiredWhenNoMoreSplits) { task->requestCancel().wait(); } +TEST_F(TaskTest, preloadedSplitsAreConsumedInReadyOrder) { + auto data = makeRowVector({makeFlatVector({1, 2, 3})}); + auto task = Task::create( + "task-ready-preload", + PlanBuilder().tableScan(asRowType(data->type())).planFragment(), + 0, + core::QueryCtx::create(), + Task::ExecutionMode::kSerial, + exec::Consumer{}); + + const std::string slowPath = "file:/tmp/slow"; + const std::string fastPath = "file:/tmp/fast"; + task->addSplit("0", exec::Split(makeHiveConnectorSplit(slowPath))); + task->addSplit("0", exec::Split(makeHiveConnectorSplit(fastPath))); + task->noMoreSplits("0"); + + folly::Baton<> allowSlow; + folly::Baton<> allowFast; + std::vector preloadThreads; + ConnectorSplitPreloadFunc preload = + [&](const std::shared_ptr& connectorSplit) { + const auto hiveSplit = + std::dynamic_pointer_cast( + connectorSplit); + ASSERT_NE(hiveSplit, nullptr); + auto* allow = + hiveSplit->filePath == slowPath ? &allowSlow : &allowFast; + connectorSplit->dataSource = + std::make_unique>( + [allow]() -> std::unique_ptr { + allow->wait(); + VELOX_FAIL("Test preload completion"); + }); + preloadThreads.emplace_back([task, connectorSplit]() { + connectorSplit->dataSource->prepare(); + task->splitPreloadFinished(kUngroupedGroupId, "0"); + }); + }; + + exec::Split split; + ContinueFuture splitFuture = ContinueFuture::makeEmpty(); + EXPECT_EQ( + task->getSplitOrFuture( + /*driverId=*/0, + kUngroupedGroupId, + "0", + /*maxPreloadSplits=*/2, + preload, + split, + splitFuture), + BlockingReason::kWaitForSplit); + EXPECT_FALSE(splitFuture.isReady()); + + // A second driver can start waiting after preloads have already been + // launched. One readiness transition must wake every current waiter so + // accumulated ready work cannot be stranded after the last completion. + exec::Split secondSplit; + ContinueFuture secondSplitFuture = ContinueFuture::makeEmpty(); + EXPECT_EQ( + task->getSplitOrFuture( + /*driverId=*/1, + kUngroupedGroupId, + "0", + /*maxPreloadSplits=*/2, + preload, + secondSplit, + secondSplitFuture), + BlockingReason::kWaitForSplit); + EXPECT_FALSE(secondSplitFuture.isReady()); + + allowFast.post(); + std::move(splitFuture).wait(); + EXPECT_TRUE(secondSplitFuture.isReady()); + + splitFuture = ContinueFuture::makeEmpty(); + EXPECT_EQ( + task->getSplitOrFuture( + /*driverId=*/0, + kUngroupedGroupId, + "0", + /*maxPreloadSplits=*/2, + preload, + split, + splitFuture), + BlockingReason::kNotBlocked); + const auto readyHiveSplit = + std::dynamic_pointer_cast( + split.connectorSplit); + ASSERT_NE(readyHiveSplit, nullptr); + EXPECT_EQ(readyHiveSplit->filePath, fastPath); + + allowSlow.post(); + for (auto& thread : preloadThreads) { + thread.join(); + } + task->requestCancel().wait(); +} + TEST_F(TaskTest, wrongPlanNodeForSplit) { auto connectorSplit = std::make_shared( "test", From 0151f9d8583d6c7335c7cf451a74488a935c4ef9 Mon Sep 17 00:00:00 2001 From: winningsix Date: Thu, 30 Jul 2026 08:12:38 +0000 Subject: [PATCH 2/2] fix(s3): serialize default credential refresh Wrap the default AWS credential chain in an eager synchronized cache. Concurrent S3 signers share one refresh and continue to rotate credentials five minutes before expiration. --- .../storage_adapters/s3fs/S3FileSystem.cpp | 3 +- .../hive/storage_adapters/s3fs/S3Util.cpp | 61 +++++++++++++++++++ .../hive/storage_adapters/s3fs/S3Util.h | 12 ++++ .../s3fs/tests/S3UtilTest.cpp | 36 +++++++++++ 4 files changed, 111 insertions(+), 1 deletion(-) diff --git a/velox/connectors/hive/storage_adapters/s3fs/S3FileSystem.cpp b/velox/connectors/hive/storage_adapters/s3fs/S3FileSystem.cpp index d96fd7707dc..66ca5b5a343 100644 --- a/velox/connectors/hive/storage_adapters/s3fs/S3FileSystem.cpp +++ b/velox/connectors/hive/storage_adapters/s3fs/S3FileSystem.cpp @@ -324,7 +324,8 @@ class S3FileSystem::Impl { // Return a default AWSCredentialsProvider. std::shared_ptr getDefaultCredentialsProvider() const { - return std::make_shared(); + return makeSynchronizedCachingCredentialsProvider( + std::make_shared()); } // Configure and return an AWSCredentialsProvider with S3 IAM Role. diff --git a/velox/connectors/hive/storage_adapters/s3fs/S3Util.cpp b/velox/connectors/hive/storage_adapters/s3fs/S3Util.cpp index 8f982131d4c..7c0679b4ef3 100644 --- a/velox/connectors/hive/storage_adapters/s3fs/S3Util.cpp +++ b/velox/connectors/hive/storage_adapters/s3fs/S3Util.cpp @@ -24,7 +24,68 @@ #include "velox/connectors/hive/storage_adapters/s3fs/S3Util.h" +#include + namespace facebook::velox::filesystems { +namespace { + +class SynchronizedCachingCredentialsProvider final + : public Aws::Auth::AWSCredentialsProvider { + public: + explicit SynchronizedCachingCredentialsProvider( + std::shared_ptr source) + : source_(std::move(source)) { + refresh(); + } + + Aws::Auth::AWSCredentials GetAWSCredentials() override { + { + std::lock_guard lock(mutex_); + if (isUsable(credentials_)) { + return credentials_; + } + } + refresh(); + std::lock_guard lock(mutex_); + return credentials_; + } + + private: + static bool isUsable(const Aws::Auth::AWSCredentials& credentials) { + return !credentials.IsEmpty() && + !credentials.ExpiresSoon(5 * 60 * 1000); + } + + void refresh() { + std::lock_guard refreshLock(refreshMutex_); + { + std::lock_guard lock(mutex_); + if (isUsable(credentials_)) { + return; + } + } + auto refreshed = source_->GetAWSCredentials(); + std::lock_guard lock(mutex_); + if (!refreshed.IsEmpty()) { + credentials_ = std::move(refreshed); + } + } + + const std::shared_ptr source_; + std::mutex mutex_; + std::mutex refreshMutex_; + Aws::Auth::AWSCredentials credentials_; +}; + +} // namespace + +std::shared_ptr +makeSynchronizedCachingCredentialsProvider( + std::shared_ptr source) { + VELOX_CHECK_NOT_NULL(source); + return std::make_shared( + std::move(source)); +} std::string getErrorStringFromS3Error( const Aws::Client::AWSError& error) { diff --git a/velox/connectors/hive/storage_adapters/s3fs/S3Util.h b/velox/connectors/hive/storage_adapters/s3fs/S3Util.h index 125d9cc805b..14488c7b3f9 100644 --- a/velox/connectors/hive/storage_adapters/s3fs/S3Util.h +++ b/velox/connectors/hive/storage_adapters/s3fs/S3Util.h @@ -21,6 +21,7 @@ #pragma once +#include #include #include #include @@ -32,6 +33,17 @@ namespace facebook::velox::filesystems { +/// Wraps a refreshable AWS credentials provider with synchronized caching. +/// +/// A refreshable provider chain can invoke the underlying identity provider +/// once per concurrent S3 signer when an empty or expiring credential is +/// observed. The wrapper serializes refresh and serves the same credential +/// snapshot until five minutes before expiration. The underlying provider +/// remains responsible for normal credential rotation. +std::shared_ptr +makeSynchronizedCachingCredentialsProvider( + std::shared_ptr source); + namespace { static std::string_view kSep{"/"}; // AWS S3 EMRFS, Hadoop block storage filesystem on-top of Amazon S3 buckets. diff --git a/velox/connectors/hive/storage_adapters/s3fs/tests/S3UtilTest.cpp b/velox/connectors/hive/storage_adapters/s3fs/tests/S3UtilTest.cpp index c71065ba2cd..e65c0007f74 100644 --- a/velox/connectors/hive/storage_adapters/s3fs/tests/S3UtilTest.cpp +++ b/velox/connectors/hive/storage_adapters/s3fs/tests/S3UtilTest.cpp @@ -18,8 +18,44 @@ #include "gtest/gtest.h" +#include +#include +#include +#include + namespace facebook::velox::filesystems { +TEST(S3UtilTest, synchronizedCachingCredentialsProvider) { + class CountingProvider final : public Aws::Auth::AWSCredentialsProvider { + public: + Aws::Auth::AWSCredentials GetAWSCredentials() override { + const auto call = ++calls; + if (call == 1) { + return {}; + } + std::this_thread::sleep_for(std::chrono::milliseconds(10)); + return {"access", "secret", "session"}; + } + + std::atomic calls{0}; + }; + + auto source = std::make_shared(); + auto cached = makeSynchronizedCachingCredentialsProvider(source); + std::vector threads; + for (size_t index = 0; index < 128; ++index) { + threads.emplace_back([cached] { + EXPECT_EQ(cached->GetAWSCredentials().GetAWSAccessKeyId(), "access"); + }); + } + for (auto& thread : threads) { + thread.join(); + } + // The eager construction attempt is empty. All concurrent callers share + // the single subsequent refresh. + EXPECT_EQ(source->calls.load(), 2); +} + // TODO: Each prefix should be implemented as its own filesystem. TEST(S3UtilTest, isS3File) { EXPECT_FALSE(isS3File("ss3://"));