diff --git a/axiom/cli/SqlQueryRunner.cpp b/axiom/cli/SqlQueryRunner.cpp index 357e0020d..c13f00bc9 100644 --- a/axiom/cli/SqlQueryRunner.cpp +++ b/axiom/cli/SqlQueryRunner.cpp @@ -1412,7 +1412,12 @@ std::string SqlQueryRunner::runExplainIo( velox::exec::SimpleExpressionEvaluator evaluator( queryCtx.get(), optimizerPool_.get()); optimizer::v2::Optimizer optimizer( - *logicalPlan, *resolver, *session, evaluator, queryCtx); + *logicalPlan, + *resolver, + *session, + evaluator, + queryCtx, + /*runtimeStats=*/nullptr); return optimizer.explainIo(std::move(outputTable)); } std::string text; @@ -1890,7 +1895,8 @@ optimizer::PlanAndStats SqlQueryRunner::optimize( *schemaResolver, *optimizerSession, evaluator, - queryCtx) + queryCtx, + runtimeStats.get()) .optimize(opts); } @@ -2101,10 +2107,14 @@ std::vector SqlQueryRunner::runShowStatsForQuery( queryCtx.get(), optimizerPool_.get()); auto resolver = orDefaultSchemaResolver(nullptr); - const auto stats = - optimizer::v2::Optimizer( - *logicalPlan, *resolver, *session, evaluator, queryCtx) - .estimateQueryStats(); + const auto stats = optimizer::v2::Optimizer( + *logicalPlan, + *resolver, + *session, + evaluator, + queryCtx, + /*runtimeStats=*/nullptr) + .estimateQueryStats(); presto::ShowStatsBuilder builder(roundCardinality(stats.cardinality)); for (const auto& column : stats.columns) { diff --git a/axiom/common/QueryRuntimeStats.h b/axiom/common/QueryRuntimeStats.h index 130166379..99e73fa18 100644 --- a/axiom/common/QueryRuntimeStats.h +++ b/axiom/common/QueryRuntimeStats.h @@ -55,6 +55,12 @@ class QueryRuntimeStats : public velox::ConcurrentRuntimeStatWriter { "axiom-optimizeToVeloxCpuNanos"}; // Split manager. + static constexpr std::string_view kSelectPartitionsWallNanos{ + "axiom-selectPartitionsWallNanos"}; + static constexpr std::string_view kSelectPartitionsCpuNanos{ + "axiom-selectPartitionsCpuNanos"}; + static constexpr std::string_view kSelectPartitionsCount{ + "axiom-selectPartitionsCount"}; static constexpr std::string_view kListPartitionsWallNanos{ "axiom-listPartitionsWallNanos"}; static constexpr std::string_view kListPartitionsCpuNanos{ diff --git a/axiom/connectors/ConnectorMetadata.h b/axiom/connectors/ConnectorMetadata.h index fe4a86290..f39803918 100644 --- a/axiom/connectors/ConnectorMetadata.h +++ b/axiom/connectors/ConnectorMetadata.h @@ -677,6 +677,9 @@ class TableLayout { /// @param session Connector session for the current query. /// @param tableHandle Table handle for the table; the connector reads its /// accepted filters from it in whatever representation it stored them. + /// @param partitionSelection Exact partitions selected for this scan. Null is + /// reserved for callers that do not perform partition selection during + /// planning. /// @param columns Names of table columns the optimizer is interested in. /// Column names correspond to actual table columns (not synthetic subfield /// projections). If the connector provides per-column statistics, it must @@ -692,6 +695,7 @@ class TableLayout { virtual folly::coro::Task> co_estimateStats( ConnectorSessionPtr /*session*/, velox::connector::ConnectorTableHandlePtr /*tableHandle*/, + PartitionSelectionPtr /*partitionSelection*/, std::vector /*columns*/, const FilterSelectivityEstimator& /*estimator*/) const { co_return std::nullopt; @@ -714,6 +718,9 @@ class TableLayout { /// @param session Connector session for the current query. /// @param tableHandle Table handle carrying the filters pushed via /// createTableHandle. + /// @param partitionSelection Exact partitions selected for this scan. Null is + /// reserved for callers that do not perform partition selection during + /// planning. /// @param groupingColumns Names of columns to group the counts by, in /// output-key order. Empty requests a single global count. The connector /// returns std::nullopt if it cannot group by these columns from metadata @@ -728,6 +735,7 @@ class TableLayout { co_metadataCounts( ConnectorSessionPtr /*session*/, velox::connector::ConnectorTableHandlePtr /*tableHandle*/, + PartitionSelectionPtr /*partitionSelection*/, std::vector /*groupingColumns*/, std::vector /*columns*/) const { co_return std::nullopt; diff --git a/axiom/connectors/ConnectorSplitManager.h b/axiom/connectors/ConnectorSplitManager.h index 0802b5664..2a18005fc 100644 --- a/axiom/connectors/ConnectorSplitManager.h +++ b/axiom/connectors/ConnectorSplitManager.h @@ -18,6 +18,7 @@ #include #include #include +#include #include "axiom/common/QueryRuntimeStats.h" #include "axiom/connectors/ConnectorSession.h" @@ -106,10 +107,33 @@ class PartitionHandle { using PartitionHandlePtr = std::shared_ptr; +/// The exact partitions selected for one scan and the storage partitioning the +/// connector guarantees across those partitions. A null storagePartitionType +/// means the scan must be planned as unbucketed. +struct PartitionSelection { + std::vector partitions; + std::shared_ptr storagePartitionType; +}; + +using PartitionSelectionPtr = std::shared_ptr; + class ConnectorSplitManager { public: virtual ~ConnectorSplitManager() = default; + /// Selects the partitions read by 'tableHandle' and validates the declared + /// storage partitioning for exactly that set. The default preserves the + /// connector's declared partitioning. + virtual folly::coro::Task co_selectPartitions( + const ConnectorSessionPtr& session, + const velox::connector::ConnectorTableHandlePtr& tableHandle, + std::shared_ptr declaredStoragePartitionType) { + auto partitions = co_await co_listPartitions(session, tableHandle); + co_return PartitionSelection{ + .partitions = std::move(partitions), + .storagePartitionType = std::move(declaredStoragePartitionType)}; + } + /// Returns a list of all partitions that match the filters in /// 'tableHandle'. A non-partitioned table returns one partition. virtual folly::coro::Task> co_listPartitions( diff --git a/axiom/connectors/hive/LocalHiveConnectorMetadata.cpp b/axiom/connectors/hive/LocalHiveConnectorMetadata.cpp index eb9918ddb..95a1bbb2f 100644 --- a/axiom/connectors/hive/LocalHiveConnectorMetadata.cpp +++ b/axiom/connectors/hive/LocalHiveConnectorMetadata.cpp @@ -576,6 +576,7 @@ folly::coro::Task> LocalHiveTableLayout::co_estimateStats( ConnectorSessionPtr /*session*/, velox::connector::ConnectorTableHandlePtr tableHandle, + PartitionSelectionPtr /*partitionSelection*/, std::vector columns, const FilterSelectivityEstimator& estimator) const { auto hiveHandle = @@ -632,6 +633,7 @@ folly::coro::Task>> LocalHiveTableLayout::co_metadataCounts( ConnectorSessionPtr /*session*/, velox::connector::ConnectorTableHandlePtr tableHandle, + PartitionSelectionPtr /*partitionSelection*/, std::vector groupingColumns, std::vector columns) const { auto hiveHandle = diff --git a/axiom/connectors/hive/LocalHiveConnectorMetadata.h b/axiom/connectors/hive/LocalHiveConnectorMetadata.h index 54c4f5ef5..3004e02ad 100644 --- a/axiom/connectors/hive/LocalHiveConnectorMetadata.h +++ b/axiom/connectors/hive/LocalHiveConnectorMetadata.h @@ -168,6 +168,7 @@ class LocalHiveTableLayout : public HiveTableLayout { folly::coro::Task> co_estimateStats( ConnectorSessionPtr session, velox::connector::ConnectorTableHandlePtr tableHandle, + PartitionSelectionPtr partitionSelection, std::vector columns, const FilterSelectivityEstimator& estimator) const override; @@ -181,6 +182,7 @@ class LocalHiveTableLayout : public HiveTableLayout { co_metadataCounts( ConnectorSessionPtr session, velox::connector::ConnectorTableHandlePtr tableHandle, + PartitionSelectionPtr partitionSelection, std::vector groupingColumns, std::vector columns) const override; diff --git a/axiom/connectors/tests/TestConnector.cpp b/axiom/connectors/tests/TestConnector.cpp index fe5683bd9..4e60d09bc 100644 --- a/axiom/connectors/tests/TestConnector.cpp +++ b/axiom/connectors/tests/TestConnector.cpp @@ -448,6 +448,23 @@ const TestTable& findTestTableForHandle( } // namespace +folly::coro::Task TestSplitManager::co_selectPartitions( + const ConnectorSessionPtr& session, + const velox::connector::ConnectorTableHandlePtr& tableHandle, + std::shared_ptr declaredStoragePartitionType) { + const auto& table = findTestTableForHandle(tableHandle); + VELOX_CHECK( + !table.layouts().empty(), "Test table must have at least one layout"); + const auto* layout = table.layouts().front()->as(); + VELOX_CHECK_NOT_NULL(layout); + auto partitions = co_await co_listPartitions(session, tableHandle); + co_return PartitionSelection{ + .partitions = std::move(partitions), + .storagePartitionType = layout->partitionsConsistent() + ? std::move(declaredStoragePartitionType) + : nullptr}; +} + folly::coro::Task> TestSplitManager::co_listPartitions( const ConnectorSessionPtr& session, @@ -1009,6 +1026,15 @@ void TestConnectorMetadata::setStats( it->second->setStats(numRows, columnStats); } +void TestConnectorMetadata::setPartitionsConsistent( + const SchemaTableName& tableName, + bool consistent) { + auto it = tables_.find(tableName); + VELOX_CHECK( + it != tables_.end(), "Table doesn't exist: {}", tableName.toString()); + it->second->mutableLayout()->setPartitionsConsistent(consistent); +} + TestDataSource::TestDataSource( const velox::RowTypePtr& outputType, const velox::connector::ColumnHandleMap& handles, @@ -1334,6 +1360,12 @@ void TestConnector::setStats( metadata_->setStats(tableName, numRows, columnStats); } +void TestConnector::setPartitionsConsistent( + const SchemaTableName& tableName, + bool consistent) { + metadata_->setPartitionsConsistent(tableName, consistent); +} + std::shared_ptr TestConnectorFactory::newConnector( const std::string& id, std::shared_ptr config, diff --git a/axiom/connectors/tests/TestConnector.h b/axiom/connectors/tests/TestConnector.h index 1f59666a5..fdd5aad5f 100644 --- a/axiom/connectors/tests/TestConnector.h +++ b/axiom/connectors/tests/TestConnector.h @@ -127,6 +127,14 @@ class TestTableLayout : public TableLayout { return partitionType_; } + void setPartitionsConsistent(bool consistent) { + partitionsConsistent_ = consistent; + } + + bool partitionsConsistent() const { + return partitionsConsistent_; + } + /// Records discrete values to use in 'discretePredicateColumns' and /// 'discretePredicates' APIs. If called repeatedly, overwrites previous /// values. @@ -161,6 +169,7 @@ class TestTableLayout : public TableLayout { std::vector discreteValueColumns_; std::vector discreteValues_; std::shared_ptr partitionType_; + bool partitionsConsistent_{true}; }; /// RowVectors are appended using the addData() interface and the vector @@ -295,6 +304,12 @@ class TestConnectorMetadata; /// only for the requested buckets' entries. class TestSplitManager : public ConnectorSplitManager { public: + folly::coro::Task co_selectPartitions( + const ConnectorSessionPtr& session, + const velox::connector::ConnectorTableHandlePtr& tableHandle, + std::shared_ptr declaredStoragePartitionType) + override; + folly::coro::Task> co_listPartitions( const ConnectorSessionPtr& session, const velox::connector::ConnectorTableHandlePtr& tableHandle) override; @@ -582,6 +597,11 @@ class TestConnectorMetadata : public ConnectorMetadata { uint64_t numRows, const std::unordered_map& columnStats); + /// Sets whether the table's partitions share its declared bucket layout. + void setPartitionsConsistent( + const SchemaTableName& tableName, + bool consistent); + TablePtr createTable( const ConnectorSessionPtr& session, const SchemaTableName& tableName, @@ -913,6 +933,18 @@ class TestConnector : public velox::connector::Connector { setStats({std::string(kDefaultSchema), tableName}, numRows, columnStats); } + /// Sets whether the table's selected partitions preserve its declared + /// bucketing. + void setPartitionsConsistent( + const SchemaTableName& tableName, + bool consistent); + + /// Convenience overload that uses kDefaultSchema as the schema. + void setPartitionsConsistent(const std::string& tableName, bool consistent) { + setPartitionsConsistent( + {std::string(kDefaultSchema), tableName}, consistent); + } + bool dropTableIfExists(const SchemaTableName& name); bool dropTableIfExists(std::string_view tableName) { diff --git a/axiom/optimizer/MultiFragmentPlan.h b/axiom/optimizer/MultiFragmentPlan.h index deff7ce3f..5544677a3 100644 --- a/axiom/optimizer/MultiFragmentPlan.h +++ b/axiom/optimizer/MultiFragmentPlan.h @@ -19,6 +19,7 @@ #include #include "axiom/common/Enums.h" #include "axiom/connectors/ConnectorMetadata.h" +#include "axiom/connectors/ConnectorSplitManager.h" #include "velox/core/PlanFragment.h" #include "velox/vector/ComplexVector.h" #include "velox/vector/SimpleVector.h" @@ -76,6 +77,10 @@ struct InputStage { std::string producerTaskPrefix; }; +/// Selected partitions keyed by the emitted TableScanNode ID. +using ScanPartitionSelectionMap = folly:: + F14FastMap; + /// Callbacks to finalize writing to a connector after the query completes. /// On success, the runner calls 'commit'. On failure, 'abort'. Only one of the /// 'commit' or 'abort' is called. @@ -222,8 +227,13 @@ class MultiFragmentPlan { } }; - MultiFragmentPlan(std::vector fragments, Options options) - : fragments_{std::move(fragments)}, options_{std::move(options)} {} + MultiFragmentPlan( + std::vector fragments, + Options options, + ScanPartitionSelectionMap scanPartitionSelections = {}) + : fragments_{std::move(fragments)}, + options_{std::move(options)}, + scanPartitionSelections_{std::move(scanPartitionSelections)} {} const std::vector& fragments() const { return fragments_; @@ -233,6 +243,10 @@ class MultiFragmentPlan { return options_; } + const ScanPartitionSelectionMap& scanPartitionSelections() const { + return scanPartitionSelections_; + } + /// @param detailed If true, includes details of each plan node. Otherwise, /// only node types are included. /// @param addContext Optional lambda to add context to plan nodes. Receives @@ -259,6 +273,7 @@ class MultiFragmentPlan { private: const std::vector fragments_; const Options options_; + const ScanPartitionSelectionMap scanPartitionSelections_; }; using MultiFragmentPlanPtr = std::shared_ptr; diff --git a/axiom/optimizer/Optimization.cpp b/axiom/optimizer/Optimization.cpp index 89d5ff81c..8f78032cc 100644 --- a/axiom/optimizer/Optimization.cpp +++ b/axiom/optimizer/Optimization.cpp @@ -197,6 +197,7 @@ void Optimization::estimateAllBaseTableSelectivity(DerivedTable& dt) { tasks.push_back(layout->co_estimateStats( std::move(connectorSession), data->handle, + /*partitionSelection=*/nullptr, std::move(columnNames), estimator)); } diff --git a/axiom/optimizer/tests/BucketedExecutionPlanTest.cpp b/axiom/optimizer/tests/BucketedExecutionPlanTest.cpp index 459917e59..e7161b27e 100644 --- a/axiom/optimizer/tests/BucketedExecutionPlanTest.cpp +++ b/axiom/optimizer/tests/BucketedExecutionPlanTest.cpp @@ -1217,6 +1217,66 @@ TEST_P(BucketedExecutionTest, bucketColumnRenamedByProjection) { } } +TEST_P(BucketedExecutionTest, partitionSelectionControlsStorageBucketing) { + if (!useV2_) { + GTEST_SKIP(); + } + + addBucketedTable("selected_orders", {"customer_id"}, 16); + testConnector_->setPartitionsConsistent("selected_orders", false); + + auto plan = planDistributed(parseSelect( + "SELECT customer_id, count(*) FROM selected_orders " + "GROUP BY customer_id", + kTestConnectorId)); + + ASSERT_EQ(plan.plan->scanPartitionSelections().size(), 1); + const auto& selection = plan.plan->scanPartitionSelections().begin()->second; + ASSERT_NE(selection, nullptr); + EXPECT_EQ(selection->partitions.size(), 16); + EXPECT_EQ(selection->storagePartitionType, nullptr); + for (const auto& fragment : plan.plan->fragments()) { + EXPECT_TRUE(fragment.groupedNodes.empty()); + } +} + +TEST_P(BucketedExecutionTest, partitionSelectionReachesPlanOutput) { + if (!useV2_) { + GTEST_SKIP(); + } + + addBucketedTable("selected_bucketed_orders", {"customer_id"}, 16); + auto plan = planDistributed(parseSelect( + "SELECT customer_id, count(*) FROM selected_bucketed_orders " + "GROUP BY customer_id", + kTestConnectorId)); + + ASSERT_EQ(plan.plan->scanPartitionSelections().size(), 1); + const auto& selection = plan.plan->scanPartitionSelections().begin()->second; + ASSERT_NE(selection, nullptr); + EXPECT_EQ(selection->partitions.size(), 16); + EXPECT_NE(selection->storagePartitionType, nullptr); + + const auto metrics = runtimeStats().runtimeStats(); + const auto selectWall = + metrics.find(std::string(QueryRuntimeStats::kSelectPartitionsWallNanos)); + ASSERT_NE(selectWall, metrics.end()); + EXPECT_GT(selectWall->second.sum, 0); + const auto selectCpu = + metrics.find(std::string(QueryRuntimeStats::kSelectPartitionsCpuNanos)); + ASSERT_NE(selectCpu, metrics.end()); + EXPECT_GT(selectCpu->second.sum, 0); + const auto selectCount = + metrics.find(std::string(QueryRuntimeStats::kSelectPartitionsCount)); + ASSERT_NE(selectCount, metrics.end()); + EXPECT_EQ(selectCount->second.sum, 16); + EXPECT_EQ( + metrics.count(std::string(QueryRuntimeStats::kListPartitionsWallNanos)), + 0); + + expectBucketedFragmentWithWidth(*plan.plan, 4); +} + AXIOM_INSTANTIATE_V1_V2(BucketedExecutionTest); } // namespace diff --git a/axiom/optimizer/tests/QueryTestBase.cpp b/axiom/optimizer/tests/QueryTestBase.cpp index 4e8d2280d..26bdc3200 100644 --- a/axiom/optimizer/tests/QueryTestBase.cpp +++ b/axiom/optimizer/tests/QueryTestBase.cpp @@ -331,9 +331,14 @@ optimizer::PlanAndStats QueryTestBase::planVelox( optimizer::PlanAndStats planAndStats; if (useV2_) { - planAndStats = - v2::Optimizer(*plan, schemaResolver, *session, evaluator, queryCtx) - .optimize(options); + planAndStats = v2::Optimizer( + *plan, + schemaResolver, + *session, + evaluator, + queryCtx, + &runtimeStats_) + .optimize(options); } else { optimizer::Optimization opt( session, diff --git a/axiom/optimizer/tests/QueryTestBase.h b/axiom/optimizer/tests/QueryTestBase.h index c1c94900f..5a42decdd 100644 --- a/axiom/optimizer/tests/QueryTestBase.h +++ b/axiom/optimizer/tests/QueryTestBase.h @@ -301,6 +301,10 @@ class QueryTestBase : public velox::exec::test::HiveConnectorTestBase { std::shared_ptr testConnector_; + QueryRuntimeStats& runtimeStats() { + return runtimeStats_; + } + private: std::shared_ptr optimizerPool_; diff --git a/axiom/optimizer/tests/SqlTestBase.cpp b/axiom/optimizer/tests/SqlTestBase.cpp index 19c4d6e00..1b6ebad57 100644 --- a/axiom/optimizer/tests/SqlTestBase.cpp +++ b/axiom/optimizer/tests/SqlTestBase.cpp @@ -226,7 +226,8 @@ std::shared_ptr SqlTestBase::makeLocalRunnerV2( schemaResolver, *optimizerSession, evaluator, - queryCtx) + queryCtx, + /*runtimeStats=*/nullptr) .optimize(options); }); } diff --git a/axiom/optimizer/v2/EmitPass.cpp b/axiom/optimizer/v2/EmitPass.cpp index 81f74fed0..50c5bc13b 100644 --- a/axiom/optimizer/v2/EmitPass.cpp +++ b/axiom/optimizer/v2/EmitPass.cpp @@ -582,6 +582,9 @@ class Emitter { // takePrediction. NodePredictionMap prediction_; + // Selected partitions keyed by emitted TableScanNode ID. + ScanPartitionSelectionMap scanPartitionSelections_; + public: // Lowers 'root' (plus the output rename to 'outputNames') into fragments, // returning them with the root fragment last. @@ -600,6 +603,10 @@ class Emitter { NodePredictionMap takePrediction() { return std::move(prediction_); } + + ScanPartitionSelectionMap takeScanPartitionSelections() { + return std::move(scanPartitionSelections_); + } }; velox::core::PlanNodePtr Emitter::projectToColumns( @@ -657,6 +664,10 @@ velox::core::PlanNodePtr Emitter::emitScan(const Scan& scan) { tableHandle, std::move(assignments)); + if (handle.partitionSelection != nullptr) { + scanPartitionSelections_.emplace(scanNode->id(), handle.partitionSelection); + } + // A grouped scan makes its fragment bucketed. if (const auto* partitionType = scan.groupedPartitionType()) { groupedLeaves_.push_back({scanNode->id(), partitionType}); @@ -2337,7 +2348,8 @@ EmitPass::Result EmitPass::run( return Result{ std::move(fragments), emitter.takeFinishWrite(), - emitter.takePrediction()}; + emitter.takePrediction(), + emitter.takeScanPartitionSelections()}; } } // namespace facebook::axiom::optimizer::v2 diff --git a/axiom/optimizer/v2/EmitPass.h b/axiom/optimizer/v2/EmitPass.h index 645254dc4..61c4f301b 100644 --- a/axiom/optimizer/v2/EmitPass.h +++ b/axiom/optimizer/v2/EmitPass.h @@ -37,6 +37,8 @@ class EmitPass { /// Per-node estimates, keyed by emitted `PlanNodeId`, for EXPLAIN. Empty /// when estimates are unavailable. NodePredictionMap prediction; + /// Query-local partition selections keyed by emitted TableScanNode ID. + ScanPartitionSelectionMap scanPartitionSelections; }; /// Lowers the tree-IR rooted at 'root' into fragments, projecting to the diff --git a/axiom/optimizer/v2/EstimateLeafStatsPass.cpp b/axiom/optimizer/v2/EstimateLeafStatsPass.cpp index dd71db9ec..3f2f52875 100644 --- a/axiom/optimizer/v2/EstimateLeafStatsPass.cpp +++ b/axiom/optimizer/v2/EstimateLeafStatsPass.cpp @@ -126,6 +126,7 @@ void EstimateLeafStatsPass::run(NodeCP root, const OptimizerSession& session) { requests.push_back(layout->co_estimateStats( std::move(connectorSession), handle->tableHandle, + handle->partitionSelection, std::move(columnNames), estimator)); } diff --git a/axiom/optimizer/v2/FoldMetadataAggregatePass.cpp b/axiom/optimizer/v2/FoldMetadataAggregatePass.cpp index 2a2b4a737..442d21466 100644 --- a/axiom/optimizer/v2/FoldMetadataAggregatePass.cpp +++ b/axiom/optimizer/v2/FoldMetadataAggregatePass.cpp @@ -136,6 +136,7 @@ class Folder : public NodeRewriter { auto result = folly::coro::blockingWait(layout->co_metadataCounts( std::move(connectorSession), handle.tableHandle, + handle.partitionSelection, std::move(groupingColumns), std::move(nullCountColumns))); if (!result.has_value()) { diff --git a/axiom/optimizer/v2/Node.cpp b/axiom/optimizer/v2/Node.cpp index 2369e3010..8ca130384 100644 --- a/axiom/optimizer/v2/Node.cpp +++ b/axiom/optimizer/v2/Node.cpp @@ -16,6 +16,8 @@ #include "axiom/optimizer/v2/Node.h" +#include "axiom/optimizer/v2/ScanHandle.h" + #include "axiom/connectors/ConnectorMetadata.h" #include "axiom/optimizer/Schema.h" #include "axiom/optimizer/v2/KeyHash.h" @@ -611,8 +613,14 @@ size_t Scan::KeyHash::operator()(const Scan* node) const { Partitioning Scan::storageBucketing() const { const connector::TableLayout* layout = baseTable_->layout(); - return bucketPartition( - layout, outputColumns(), layout->partitionType().get()); + // A missing selection means this connector did not opt into per-scan + // certification, so preserve its declared layout. A present selection with a + // null storagePartitionType is an explicit unbucketed result. + const auto* storagePartitionType = + scanHandle_ != nullptr && scanHandle_->partitionSelection != nullptr + ? scanHandle_->partitionSelection->storagePartitionType.get() + : layout->partitionType().get(); + return bucketPartition(layout, outputColumns(), storagePartitionType); } size_t Scan::KeyHash::operator()(const Key& key) const { diff --git a/axiom/optimizer/v2/Optimize.cpp b/axiom/optimizer/v2/Optimize.cpp index 754b8e81f..e178066eb 100644 --- a/axiom/optimizer/v2/Optimize.cpp +++ b/axiom/optimizer/v2/Optimize.cpp @@ -49,7 +49,8 @@ FrontendResult translateAndPushdown( Builder& builder, const OptimizerSession& session, const std::shared_ptr& queryCtx, - PushdownAndPrunePass::ConnectorPushdown connectorPushdown) { + PushdownAndPrunePass::ConnectorPushdown connectorPushdown, + QueryRuntimeStats* runtimeStats) { ConstantPlanRunner constantPlanRunner{queryCtx}; auto translated = TranslatePass::run( plan, schema, evaluator, builder, session, constantPlanRunner); @@ -61,7 +62,8 @@ FrontendResult translateAndPushdown( builder, evaluator, session, - connectorPushdown); + connectorPushdown, + runtimeStats); return {std::move(translated), pushed}; } @@ -151,7 +153,8 @@ PlanAndStats Optimizer::optimize(const MultiFragmentPlan::Options& options) { builder, session_, queryCtx_, - PushdownAndPrunePass::ConnectorPushdown::kOffer); + PushdownAndPrunePass::ConnectorPushdown::kOffer, + runtimeStats_); NodeCP folded = FoldMetadataAggregatePass::run(frontend.pushed, builder, session_); if (session_.options().useFilteredTableStats) { @@ -175,7 +178,9 @@ PlanAndStats Optimizer::optimize(const MultiFragmentPlan::Options& options) { PlanAndStats result; result.plan = std::make_shared( - std::move(emitted.fragments), planOptions); + std::move(emitted.fragments), + planOptions, + std::move(emitted.scanPartitionSelections)); result.plan->checkConsistency( /*mayBeEmpty=*/plan_.is(logical_plan::NodeKind::kTableWrite)); result.finishWrite = std::move(emitted.finishWrite); @@ -215,7 +220,8 @@ std::string Optimizer::explainIo( builder, session_, queryCtx_, - PushdownAndPrunePass::ConnectorPushdown::kSkip); + PushdownAndPrunePass::ConnectorPushdown::kSkip, + /*runtimeStats=*/nullptr); std::vector> tableFilters; collectScans(frontend.pushed, tableFilters); @@ -235,7 +241,8 @@ QueryStats Optimizer::estimateQueryStats() { builder, session_, queryCtx_, - PushdownAndPrunePass::ConnectorPushdown::kOffer); + PushdownAndPrunePass::ConnectorPushdown::kOffer, + runtimeStats_); if (session_.options().useFilteredTableStats) { EstimateLeafStatsPass::run(frontend.pushed, session_); } diff --git a/axiom/optimizer/v2/Optimize.h b/axiom/optimizer/v2/Optimize.h index c92ed1ebb..63c01aab8 100644 --- a/axiom/optimizer/v2/Optimize.h +++ b/axiom/optimizer/v2/Optimize.h @@ -21,6 +21,7 @@ #include #include "axiom/common/CatalogSchemaTableName.h" +#include "axiom/common/QueryRuntimeStats.h" #include "axiom/connectors/SchemaResolver.h" #include "axiom/logical_plan/LogicalPlanNode.h" #include "axiom/optimizer/MultiFragmentPlan.h" @@ -66,12 +67,14 @@ class Optimizer { const connector::SchemaResolver& schemaResolver, const OptimizerSession& session, velox::core::ExpressionEvaluator& evaluator, - std::shared_ptr queryCtx) + std::shared_ptr queryCtx, + QueryRuntimeStats* runtimeStats) : plan_{plan}, schemaResolver_{schemaResolver}, session_{session}, evaluator_{evaluator}, - queryCtx_{std::move(queryCtx)} {} + queryCtx_{std::move(queryCtx)}, + runtimeStats_{runtimeStats} {} /// Lowers the plan to a distributed Velox execution plan (a /// `MultiFragmentPlan` of one or more fragments). @@ -123,6 +126,7 @@ class Optimizer { const OptimizerSession& session_; velox::core::ExpressionEvaluator& evaluator_; const std::shared_ptr queryCtx_; + QueryRuntimeStats* const runtimeStats_; }; } // namespace facebook::axiom::optimizer::v2 diff --git a/axiom/optimizer/v2/PushdownAndPrunePass.cpp b/axiom/optimizer/v2/PushdownAndPrunePass.cpp index ee597c3d3..7bb25b31b 100644 --- a/axiom/optimizer/v2/PushdownAndPrunePass.cpp +++ b/axiom/optimizer/v2/PushdownAndPrunePass.cpp @@ -641,12 +641,14 @@ class Pushdown : public NodeRewriter { Builder& builder, velox::core::ExpressionEvaluator& evaluator, const OptimizerSession& session, - PushdownAndPrunePass::ConnectorPushdown connectorPushdown) + PushdownAndPrunePass::ConnectorPushdown connectorPushdown, + QueryRuntimeStats* runtimeStats) : NodeRewriter(builder), exprs_(builder), evaluator_(evaluator), session_(session), connectorPushdown_(connectorPushdown), + runtimeStats_(runtimeStats), simplifier_(builder, evaluator) {} protected: @@ -1185,7 +1187,9 @@ class Pushdown : public NodeRewriter { filters, session_, evaluator_, - rejectedHere)); + rejectedHere, + /*resolvePartitionSelection=*/true, + runtimeStats_)); appendAll(rejected, rejectedHere); negotiatedByBaseTableId_.emplace( baseTable.id(), Negotiated{handle, filters, std::move(rejectedHere)}); @@ -2021,6 +2025,7 @@ class Pushdown : public NodeRewriter { velox::core::ExpressionEvaluator& evaluator_; const OptimizerSession& session_; const PushdownAndPrunePass::ConnectorPushdown connectorPushdown_; + QueryRuntimeStats* const runtimeStats_; ExprSimplifier simplifier_; // Outcome of one negotiation with the connector. @@ -2044,8 +2049,9 @@ NodeCP PushdownAndPrunePass::run( Builder& builder, velox::core::ExpressionEvaluator& evaluator, const OptimizerSession& session, - ConnectorPushdown connectorPushdown) { - Pushdown pass{builder, evaluator, session, connectorPushdown}; + ConnectorPushdown connectorPushdown, + QueryRuntimeStats* runtimeStats) { + Pushdown pass{builder, evaluator, session, connectorPushdown, runtimeStats}; PushdownContext context; context.required = PlanObjectSet::fromObjects(outputColumns); context.requiredAbove = context.required; diff --git a/axiom/optimizer/v2/PushdownAndPrunePass.h b/axiom/optimizer/v2/PushdownAndPrunePass.h index 73e927614..f8a66ea57 100644 --- a/axiom/optimizer/v2/PushdownAndPrunePass.h +++ b/axiom/optimizer/v2/PushdownAndPrunePass.h @@ -16,6 +16,7 @@ #pragma once +#include "axiom/common/QueryRuntimeStats.h" #include "axiom/optimizer/OptimizerSession.h" #include "axiom/optimizer/v2/Builder.h" #include "axiom/optimizer/v2/Node.h" @@ -89,7 +90,8 @@ class PushdownAndPrunePass { Builder& builder, velox::core::ExpressionEvaluator& evaluator, const OptimizerSession& session, - ConnectorPushdown connectorPushdown); + ConnectorPushdown connectorPushdown, + QueryRuntimeStats* runtimeStats); }; } // namespace facebook::axiom::optimizer::v2 diff --git a/axiom/optimizer/v2/ScanHandle.cpp b/axiom/optimizer/v2/ScanHandle.cpp index 0e366c32b..f3101d0e4 100644 --- a/axiom/optimizer/v2/ScanHandle.cpp +++ b/axiom/optimizer/v2/ScanHandle.cpp @@ -16,9 +16,14 @@ #include "axiom/optimizer/v2/ScanHandle.h" +#include +#include + +#include "axiom/connectors/ConnectorMetadataRegistry.h" #include "axiom/optimizer/QueryGraph.h" #include "axiom/optimizer/Schema.h" #include "axiom/optimizer/v2/ExprEmitter.h" +#include "folly/coro/BlockingWait.h" namespace facebook::axiom::optimizer::v2 { @@ -28,7 +33,9 @@ ScanHandle ScanHandle::build( const ExprVector& filters, const OptimizerSession& session, velox::core::ExpressionEvaluator& evaluator, - ExprVector& rejected) { + ExprVector& rejected, + bool resolvePartitionSelection, + QueryRuntimeStats* runtimeStats) { const auto* layout = baseTable.schemaTable->columnGroups[0]->layout; auto connectorSession = session.toConnectorSession(layout->connector()->connectorId()); @@ -105,6 +112,37 @@ ScanHandle ScanHandle::build( rejected.push_back(filters[index]); } + connector::PartitionSelectionPtr partitionSelection; + // Some connectors provide table layouts without a registered + // ConnectorMetadata. Registered connectors resolve the exact selection + // through their split manager when the caller needs it. + if (resolvePartitionSelection) { + if (auto metadata = connector::ConnectorMetadataRegistry::tryGet( + layout->connectorId())) { + if (auto* splitManager = metadata->splitManager()) { + const auto cpuStart = velox::process::threadCpuNanos(); + const auto startThreadId = std::this_thread::get_id(); + const auto start = std::chrono::steady_clock::now(); + partitionSelection = std::make_shared( + folly::coro::blockingWait(splitManager->co_selectPartitions( + connectorSession, tableHandle, layout->partitionType()))); + if (runtimeStats != nullptr) { + runtimeStats->addTiming( + QueryRuntimeStats::kSelectPartitionsWallNanos, + std::chrono::steady_clock::now() - start); + recordCpuIfSameThread( + *runtimeStats, + QueryRuntimeStats::kSelectPartitionsCpuNanos, + cpuStart, + startThreadId); + runtimeStats->addCount( + QueryRuntimeStats::kSelectPartitionsCount, + partitionSelection->partitions.size()); + } + } + } + } + PlanObjectSet rejectedColumns; rejectedColumns.unionColumns(rejected); for (auto& [column, handle] : filterOnlyHandles) { @@ -115,6 +153,7 @@ ScanHandle ScanHandle::build( return ScanHandle{ .tableHandle = std::move(tableHandle), + .partitionSelection = std::move(partitionSelection), .columnHandles = std::move(columnHandles), }; } diff --git a/axiom/optimizer/v2/ScanHandle.h b/axiom/optimizer/v2/ScanHandle.h index a03352703..e9da76f63 100644 --- a/axiom/optimizer/v2/ScanHandle.h +++ b/axiom/optimizer/v2/ScanHandle.h @@ -18,6 +18,8 @@ #include +#include "axiom/common/QueryRuntimeStats.h" +#include "axiom/connectors/ConnectorSplitManager.h" #include "axiom/optimizer/OptimizerSession.h" #include "axiom/optimizer/QueryGraph.h" #include "velox/connectors/Connector.h" @@ -39,11 +41,17 @@ struct ScanHandle { const ExprVector& filters, const OptimizerSession& session, velox::core::ExpressionEvaluator& evaluator, - ExprVector& rejected); + ExprVector& rejected, + bool resolvePartitionSelection, + QueryRuntimeStats* runtimeStats); /// Table handle with the accepted filters pushed into the connector. velox::connector::ConnectorTableHandlePtr tableHandle; + /// Exact versioned partitions selected for this scan and the storage + /// partitioning certified across them. + connector::PartitionSelectionPtr partitionSelection; + /// Column handle per column of the connector read schema: every column the /// `Scan` outputs, plus the ones only the connector's own filters read, /// which it reads without projecting. diff --git a/axiom/optimizer/v2/TranslatePass.cpp b/axiom/optimizer/v2/TranslatePass.cpp index 6f6f73b0d..55b0f7c29 100644 --- a/axiom/optimizer/v2/TranslatePass.cpp +++ b/axiom/optimizer/v2/TranslatePass.cpp @@ -3352,7 +3352,9 @@ ExprCP Translator::tryFoldConstantScalar(NodeCP body) { filters, session_, evaluator_, - rejected); + rejected, + /*resolvePartitionSelection=*/false, + /*runtimeStats=*/nullptr); auto connectorSession = session_.toConnectorSession(discreteLayout.layout->connectorId()); diff --git a/axiom/runner/LocalRunner.cpp b/axiom/runner/LocalRunner.cpp index 9d9e42ea8..83db361af 100644 --- a/axiom/runner/LocalRunner.cpp +++ b/axiom/runner/LocalRunner.cpp @@ -65,6 +65,7 @@ std::shared_ptr SimpleSplitSourceFactory::splitSourceForScan( const connector::ConnectorSessionPtr& /* session */, const velox::core::TableScanNode& scan, + connector::PartitionSelectionPtr /*partitionSelection*/, const std::shared_ptr& /*partitionType*/, std::optional samplePercentage) { VELOX_USER_CHECK( @@ -81,6 +82,7 @@ std::shared_ptr ConnectorSplitSourceFactory::splitSourceForScan( const connector::ConnectorSessionPtr& session, const velox::core::TableScanNode& scan, + connector::PartitionSelectionPtr partitionSelection, const std::shared_ptr& partitionType, std::optional samplePercentage) { const auto& handle = scan.tableHandle(); @@ -88,26 +90,33 @@ ConnectorSplitSourceFactory::splitSourceForScan( connector::ConnectorMetadataRegistry::get(handle->connectorId()); auto splitManager = metadata->splitManager(); - auto listCpuStart = velox::process::threadCpuNanos(); - auto listThreadId = std::this_thread::get_id(); - auto listStart = std::chrono::steady_clock::now(); - auto partitions = folly::coro::blockingWait( - splitManager->co_listPartitions(session, handle)); - recordCpuIfSameThread( - runtimeStats_, - QueryRuntimeStats::kListPartitionsCpuNanos, - listCpuStart, - listThreadId); - runtimeStats_.addTiming( - QueryRuntimeStats::kListPartitionsWallNanos, - std::chrono::steady_clock::now() - listStart); - runtimeStats_.addCount( - QueryRuntimeStats::kListPartitionsCount, partitions.size()); + std::vector listedPartitions; + const std::vector* partitions; + if (partitionSelection != nullptr) { + partitions = &partitionSelection->partitions; + } else { + auto listCpuStart = velox::process::threadCpuNanos(); + auto listThreadId = std::this_thread::get_id(); + auto listStart = std::chrono::steady_clock::now(); + listedPartitions = folly::coro::blockingWait( + splitManager->co_listPartitions(session, handle)); + recordCpuIfSameThread( + runtimeStats_, + QueryRuntimeStats::kListPartitionsCpuNanos, + listCpuStart, + listThreadId); + runtimeStats_.addTiming( + QueryRuntimeStats::kListPartitionsWallNanos, + std::chrono::steady_clock::now() - listStart); + runtimeStats_.addCount( + QueryRuntimeStats::kListPartitionsCount, listedPartitions.size()); + partitions = &listedPartitions; + } return splitManager->getSplitSource( session, handle, - partitions, + *partitions, partitionType, samplePercentage, runtimeStats_); @@ -528,8 +537,17 @@ std::shared_ptr LocalRunner::splitSourceForScan( const velox::core::TableScanNode& scan, const std::shared_ptr& partitionType, std::optional samplePercentage) { + connector::PartitionSelectionPtr partitionSelection; + if (auto it = plan_->scanPartitionSelections().find(scan.id()); + it != plan_->scanPartitionSelections().end()) { + partitionSelection = it->second; + } return splitSourceFactory_->splitSourceForScan( - session, scan, partitionType, samplePercentage); + session, + scan, + std::move(partitionSelection), + partitionType, + samplePercentage); } void LocalRunner::cancelTasks() { diff --git a/axiom/runner/LocalRunner.h b/axiom/runner/LocalRunner.h index 7c9c48e27..41c431328 100644 --- a/axiom/runner/LocalRunner.h +++ b/axiom/runner/LocalRunner.h @@ -46,6 +46,7 @@ class SplitSourceFactory { virtual std::shared_ptr splitSourceForScan( const connector::ConnectorSessionPtr& session, const velox::core::TableScanNode& scan, + connector::PartitionSelectionPtr partitionSelection, const std::shared_ptr& partitionType, std::optional samplePercentage) = 0; }; @@ -62,6 +63,7 @@ class SimpleSplitSourceFactory : public SplitSourceFactory { std::shared_ptr splitSourceForScan( const connector::ConnectorSessionPtr& session, const velox::core::TableScanNode& scan, + connector::PartitionSelectionPtr partitionSelection, const std::shared_ptr& partitionType, std::optional samplePercentage) override; @@ -81,6 +83,7 @@ class ConnectorSplitSourceFactory : public SplitSourceFactory { std::shared_ptr splitSourceForScan( const connector::ConnectorSessionPtr& session, const velox::core::TableScanNode& scan, + connector::PartitionSelectionPtr partitionSelection, const std::shared_ptr& partitionType, std::optional samplePercentage) override; diff --git a/axiom/runner/tests/ProgressReporterTest.cpp b/axiom/runner/tests/ProgressReporterTest.cpp index b252737b3..b3cf7f0fa 100644 --- a/axiom/runner/tests/ProgressReporterTest.cpp +++ b/axiom/runner/tests/ProgressReporterTest.cpp @@ -72,11 +72,16 @@ class GatedSplitSourceFactory : public SplitSourceFactory { std::shared_ptr splitSourceForScan( const connector::ConnectorSessionPtr& session, const velox::core::TableScanNode& scan, + connector::PartitionSelectionPtr partitionSelection, const std::shared_ptr& partitionType, std::optional samplePercentage) override { return std::make_shared( inner_.splitSourceForScan( - session, scan, partitionType, samplePercentage), + session, + scan, + std::move(partitionSelection), + partitionType, + samplePercentage), gate_); }