Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
22 changes: 16 additions & 6 deletions axiom/cli/SqlQueryRunner.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -1890,7 +1895,8 @@ optimizer::PlanAndStats SqlQueryRunner::optimize(
*schemaResolver,
*optimizerSession,
evaluator,
queryCtx)
queryCtx,
runtimeStats.get())
.optimize(opts);
}

Expand Down Expand Up @@ -2101,10 +2107,14 @@ std::vector<velox::RowVectorPtr> 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) {
Expand Down
6 changes: 6 additions & 0 deletions axiom/common/QueryRuntimeStats.h
Original file line number Diff line number Diff line change
Expand Up @@ -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{
Expand Down
8 changes: 8 additions & 0 deletions axiom/connectors/ConnectorMetadata.h
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -692,6 +695,7 @@ class TableLayout {
virtual folly::coro::Task<std::optional<FilteredTableStats>> co_estimateStats(
ConnectorSessionPtr /*session*/,
velox::connector::ConnectorTableHandlePtr /*tableHandle*/,
PartitionSelectionPtr /*partitionSelection*/,
std::vector<std::string> /*columns*/,
const FilterSelectivityEstimator& /*estimator*/) const {
co_return std::nullopt;
Expand All @@ -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
Expand All @@ -728,6 +735,7 @@ class TableLayout {
co_metadataCounts(
ConnectorSessionPtr /*session*/,
velox::connector::ConnectorTableHandlePtr /*tableHandle*/,
PartitionSelectionPtr /*partitionSelection*/,
std::vector<std::string> /*groupingColumns*/,
std::vector<std::string> /*columns*/) const {
co_return std::nullopt;
Expand Down
24 changes: 24 additions & 0 deletions axiom/connectors/ConnectorSplitManager.h
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@
#include <folly/coro/Task.h>
#include <velox/connectors/Connector.h>
#include <optional>
#include <utility>
#include "axiom/common/QueryRuntimeStats.h"
#include "axiom/connectors/ConnectorSession.h"

Expand Down Expand Up @@ -106,10 +107,33 @@ class PartitionHandle {

using PartitionHandlePtr = std::shared_ptr<const PartitionHandle>;

/// 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<PartitionHandlePtr> partitions;
std::shared_ptr<const PartitionType> storagePartitionType;
};

using PartitionSelectionPtr = std::shared_ptr<const PartitionSelection>;

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<PartitionSelection> co_selectPartitions(
const ConnectorSessionPtr& session,
const velox::connector::ConnectorTableHandlePtr& tableHandle,
std::shared_ptr<const PartitionType> 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<std::vector<PartitionHandlePtr>> co_listPartitions(
Expand Down
2 changes: 2 additions & 0 deletions axiom/connectors/hive/LocalHiveConnectorMetadata.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -576,6 +576,7 @@ folly::coro::Task<std::optional<FilteredTableStats>>
LocalHiveTableLayout::co_estimateStats(
ConnectorSessionPtr /*session*/,
velox::connector::ConnectorTableHandlePtr tableHandle,
PartitionSelectionPtr /*partitionSelection*/,
std::vector<std::string> columns,
const FilterSelectivityEstimator& estimator) const {
auto hiveHandle =
Expand Down Expand Up @@ -632,6 +633,7 @@ folly::coro::Task<std::optional<std::vector<MetadataCountGroup>>>
LocalHiveTableLayout::co_metadataCounts(
ConnectorSessionPtr /*session*/,
velox::connector::ConnectorTableHandlePtr tableHandle,
PartitionSelectionPtr /*partitionSelection*/,
std::vector<std::string> groupingColumns,
std::vector<std::string> columns) const {
auto hiveHandle =
Expand Down
2 changes: 2 additions & 0 deletions axiom/connectors/hive/LocalHiveConnectorMetadata.h
Original file line number Diff line number Diff line change
Expand Up @@ -168,6 +168,7 @@ class LocalHiveTableLayout : public HiveTableLayout {
folly::coro::Task<std::optional<FilteredTableStats>> co_estimateStats(
ConnectorSessionPtr session,
velox::connector::ConnectorTableHandlePtr tableHandle,
PartitionSelectionPtr partitionSelection,
std::vector<std::string> columns,
const FilterSelectivityEstimator& estimator) const override;

Expand All @@ -181,6 +182,7 @@ class LocalHiveTableLayout : public HiveTableLayout {
co_metadataCounts(
ConnectorSessionPtr session,
velox::connector::ConnectorTableHandlePtr tableHandle,
PartitionSelectionPtr partitionSelection,
std::vector<std::string> groupingColumns,
std::vector<std::string> columns) const override;

Expand Down
32 changes: 32 additions & 0 deletions axiom/connectors/tests/TestConnector.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -448,6 +448,23 @@ const TestTable& findTestTableForHandle(

} // namespace

folly::coro::Task<PartitionSelection> TestSplitManager::co_selectPartitions(
const ConnectorSessionPtr& session,
const velox::connector::ConnectorTableHandlePtr& tableHandle,
std::shared_ptr<const PartitionType> 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<TestTableLayout>();
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<std::vector<PartitionHandlePtr>>
TestSplitManager::co_listPartitions(
const ConnectorSessionPtr& session,
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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<velox::connector::Connector> TestConnectorFactory::newConnector(
const std::string& id,
std::shared_ptr<const velox::config::ConfigBase> config,
Expand Down
32 changes: 32 additions & 0 deletions axiom/connectors/tests/TestConnector.h
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -161,6 +169,7 @@ class TestTableLayout : public TableLayout {
std::vector<const Column*> discreteValueColumns_;
std::vector<velox::Variant> discreteValues_;
std::shared_ptr<const PartitionType> partitionType_;
bool partitionsConsistent_{true};
};

/// RowVectors are appended using the addData() interface and the vector
Expand Down Expand Up @@ -295,6 +304,12 @@ class TestConnectorMetadata;
/// only for the requested buckets' entries.
class TestSplitManager : public ConnectorSplitManager {
public:
folly::coro::Task<PartitionSelection> co_selectPartitions(
const ConnectorSessionPtr& session,
const velox::connector::ConnectorTableHandlePtr& tableHandle,
std::shared_ptr<const PartitionType> declaredStoragePartitionType)
override;

folly::coro::Task<std::vector<PartitionHandlePtr>> co_listPartitions(
const ConnectorSessionPtr& session,
const velox::connector::ConnectorTableHandlePtr& tableHandle) override;
Expand Down Expand Up @@ -582,6 +597,11 @@ class TestConnectorMetadata : public ConnectorMetadata {
uint64_t numRows,
const std::unordered_map<std::string, ColumnStatistics>& 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,
Expand Down Expand Up @@ -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) {
Expand Down
19 changes: 17 additions & 2 deletions axiom/optimizer/MultiFragmentPlan.h
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@
#include <folly/container/F14Map.h>
#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"
Expand Down Expand Up @@ -76,6 +77,10 @@ struct InputStage {
std::string producerTaskPrefix;
};

/// Selected partitions keyed by the emitted TableScanNode ID.
using ScanPartitionSelectionMap = folly::
F14FastMap<velox::core::PlanNodeId, connector::PartitionSelectionPtr>;

/// 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.
Expand Down Expand Up @@ -222,8 +227,13 @@ class MultiFragmentPlan {
}
};

MultiFragmentPlan(std::vector<ExecutableFragment> fragments, Options options)
: fragments_{std::move(fragments)}, options_{std::move(options)} {}
MultiFragmentPlan(
std::vector<ExecutableFragment> fragments,
Options options,
ScanPartitionSelectionMap scanPartitionSelections = {})
: fragments_{std::move(fragments)},
options_{std::move(options)},
scanPartitionSelections_{std::move(scanPartitionSelections)} {}

const std::vector<ExecutableFragment>& fragments() const {
return fragments_;
Expand All @@ -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
Expand All @@ -259,6 +273,7 @@ class MultiFragmentPlan {
private:
const std::vector<ExecutableFragment> fragments_;
const Options options_;
const ScanPartitionSelectionMap scanPartitionSelections_;
};

using MultiFragmentPlanPtr = std::shared_ptr<const MultiFragmentPlan>;
Expand Down
1 change: 1 addition & 0 deletions axiom/optimizer/Optimization.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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));
}
Expand Down
Loading
Loading